Compare commits

..

1 Commits

Author SHA1 Message Date
Kit Langton f2d133edf7 refactor(core): separate v2 message identity 2026-06-04 12:08:00 -04:00
35 changed files with 791 additions and 357 deletions
+4 -4
View File
@@ -1,8 +1,8 @@
{
"nodeModules": {
"x86_64-linux": "sha256-dvFu5Cbs8MFoSBQXwv4HN2vyh5p20dh6QC5zZiFr0qs=",
"aarch64-linux": "sha256-l0xO7Nocl6enQxLQlLB71mG+NuT6I1eQQ1FgLtYGQOg=",
"aarch64-darwin": "sha256-WY6Lstxt4n4n63kYZUX09birHx7sNvl0Pegc6L13mGE=",
"x86_64-darwin": "sha256-sZdG40TSE9KhrmLQyQMPRugGo6R7AS3wgHiEGYtcXtc="
"x86_64-linux": "sha256-VJyM9J6fjpFjlUGDl4AKPZB2IpV+nN8Tduj4n1dPSFw=",
"aarch64-linux": "sha256-L8X/W7XbAUwwz+4luj4qVZbrnnjM2+DGuEL1aSNO2wM=",
"aarch64-darwin": "sha256-oA2xSHcs/axmGH1kcXZpELrv5DlB6DRDW3ybD3s4z1E=",
"x86_64-darwin": "sha256-mG7z3K9IOdX/XdiHLvCyLfVvXwsgMw1LIc2XbJUdRRI="
}
}
Executable → Regular
View File
Executable → Regular
View File
Executable → Regular
View File
+23
View File
@@ -0,0 +1,23 @@
# Core database migrations
## V2 beta reset boundary
During the pre-launch V2 beta period, these unreleased tables may be truncated
by an ordinary compatibility migration:
- `session_message`: disposable V2 timeline projection.
- `session_input`: disposable V2 prompt-admission inbox. Truncating it may drop
accepted but unpromoted beta prompts, so call that out explicitly.
- `event`: unreleased workspace synchronization history.
- `event_sequence`: unreleased workspace synchronization cursor and owner state.
Resetting `event` and `event_sequence` intentionally makes existing Sessions
non-warpable until new replayable history is recorded. Call that out explicitly.
Do not truncate these tables as part of a V2 compatibility migration:
- `session`, `message`, `part`: canonical V1 Session history.
If a proposed V2 schema change appears to require resetting anything outside
the wipeable beta tables, stop and design an explicit compatibility or
fresh-database cutover plan instead.
@@ -0,0 +1,4 @@
DELETE FROM `session_input`;--> statement-breakpoint
DELETE FROM `session_message`;--> statement-breakpoint
DELETE FROM `event`;--> statement-breakpoint
DELETE FROM `event_sequence`;
+1
View File
@@ -31,5 +31,6 @@ export const migrations = (
import("./migration/20260603040000_session_message_projection_order"),
import("./migration/20260603141458_session_input_inbox"),
import("./migration/20260603160727_jittery_ezekiel_stane"),
import("./migration/20260604153000_session_message_identity"),
])
).map((module) => module.default) satisfies DatabaseMigration.Migration[]
@@ -0,0 +1,16 @@
import { Effect } from "effect"
import type { DatabaseMigration } from "../migration"
export default {
id: "20260604153000_session_message_identity",
up(tx) {
return Effect.gen(function* () {
// These tables remain disposable until workspace sync and V2 Sessions launch.
// Preserve canonical V1 session, message, and part rows.
yield* tx.run(`DELETE FROM \`session_input\`;`)
yield* tx.run(`DELETE FROM \`session_message\`;`)
yield* tx.run(`DELETE FROM \`event\`;`)
yield* tx.run(`DELETE FROM \`event_sequence\`;`)
})
},
} satisfies DatabaseMigration.Migration
+1 -1
View File
@@ -8,7 +8,7 @@ import { Location } from "./location"
import { externalID, type ExternalID, NonNegativeInt, withStatics } from "./schema"
import { Identifier } from "./util/identifier"
export const ID = Schema.String.pipe(
export const ID = Schema.String.check(Schema.isStartsWith("evt_")).pipe(
Schema.brand("Event.ID"),
withStatics((schema) => ({
create: () => schema.make("evt_" + Identifier.ascending()),
-4
View File
@@ -1,4 +0,0 @@
declare module "*.md" {
const content: string
export default content
}
+4 -6
View File
@@ -1,20 +1,18 @@
/// <reference path="../markdown.d.ts" />
export * as SkillPlugin from "./skill"
import { Effect } from "effect"
import { PluginV2 } from "../plugin"
import { AbsolutePath } from "../schema"
import { SkillV2 } from "../skill"
import customizeOpencodeContent from "./skill/customize-opencode.md" with { type: "text" }
export const CustomizeOpencodeContent = customizeOpencodeContent
export const Plugin = PluginV2.define({
id: PluginV2.ID.make("skill"),
effect: Effect.gen(function* () {
const skill = yield* SkillV2.Service
const transform = yield* skill.transform()
const content = yield* Effect.promise(() =>
Bun.file(new URL("./skill/customize-opencode.md", import.meta.url)).text(),
)
yield* transform((editor) => {
editor.source(
@@ -25,7 +23,7 @@ export const Plugin = PluginV2.define({
description:
"Use ONLY when the user is editing or creating opencode's own configuration: opencode.json, opencode.jsonc, files under .opencode/, or files under ~/.config/opencode/. Also use when creating or fixing opencode agents, subagents, skills, plugins, MCP servers, or permission rules. Do not use for the user's own application code, or for any project that is not configuring opencode itself.",
location: AbsolutePath.make("/builtin/customize-opencode.md"),
content: CustomizeOpencodeContent,
content,
}),
}),
)
+4 -4
View File
@@ -3,7 +3,7 @@ export * as SessionInput from "./input"
import { and, asc, eq, inArray, isNull } from "drizzle-orm"
import { DateTime, Effect, Schema } from "effect"
import type { Database } from "../database/database"
import type { EventV2 } from "../event"
import { EventV2 } from "../event"
import { EventTable } from "../event/sql"
import { NonNegativeInt, PositiveInt } from "../schema"
import { V2Schema } from "../v2-schema"
@@ -66,7 +66,7 @@ export const admit = Effect.fn("SessionInput.admit")(function* (
const event = yield* db
.select({ id: EventTable.id })
.from(EventTable)
.where(eq(EventTable.id, input.id))
.where(eq(EventTable.id, SessionMessage.ID.toEvent(input.id)))
.get()
.pipe(Effect.orDie)
const message = yield* db
@@ -133,7 +133,7 @@ export const guardReservedID = Effect.fn("SessionInput.guardReservedID")(functio
db: DatabaseService,
event: EventV2.Payload,
) {
const admitted = yield* find(db, event.id)
const admitted = yield* find(db, SessionMessage.ID.fromEvent(event.id))
if (admitted === undefined) return
if (!Schema.is(SessionEvent.Prompted)(event))
return yield* Effect.die("Durable event conflicts with admitted prompt input")
@@ -225,7 +225,7 @@ const publish = Effect.fn("SessionInput.publish")(function* (
prompt: decodePrompt(row.prompt),
delivery: row.delivery,
},
{ id: SessionMessage.ID.make(row.id) },
{ id: SessionMessage.ID.toEvent(SessionMessage.ID.make(row.id)) },
)
}
return rows.length
+16
View File
@@ -0,0 +1,16 @@
export * as SessionMessageID from "./message-id"
import { Schema } from "effect"
import { EventV2 } from "../event"
import { withStatics } from "../schema"
import { Identifier } from "../util/identifier"
export const ID = Schema.String.check(Schema.isStartsWith("msg_")).pipe(
Schema.brand("Session.Message.ID"),
withStatics((schema) => ({
create: () => schema.make("msg_" + Identifier.ascending()),
fromEvent: (id: EventV2.ID) => schema.make("msg" + id.slice(3)),
toEvent: (id: ID) => EventV2.ID.make("evt" + id.slice(3)),
})),
)
export type ID = typeof ID.Type
+10 -9
View File
@@ -1,5 +1,6 @@
import { castDraft, produce, type WritableDraft } from "immer"
import { Effect } from "effect"
import type { EventV2 } from "../event"
import { SessionEvent } from "./event"
import { SessionMessage } from "./message"
@@ -112,9 +113,9 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
const latestReasoning = (assistant: DraftAssistant | undefined, reasoningID: string) =>
assistant?.content.findLast((item): item is DraftReasoning => item.type === "reasoning" && item.id === reasoningID)
const updateOwnedAssistant = (messageID: SessionMessage.ID, recipe: (draft: DraftAssistant) => void) =>
const updateOwnedAssistant = (messageID: EventV2.ID, recipe: (draft: DraftAssistant) => void) =>
Effect.gen(function* () {
const assistant = yield* adapter.getAssistant(messageID)
const assistant = yield* adapter.getAssistant(SessionMessage.ID.fromEvent(messageID))
if (assistant) yield* adapter.updateAssistant(produce(assistant, recipe))
})
@@ -123,7 +124,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.agent.switched": (event) => {
return adapter.appendMessage(
new SessionMessage.AgentSwitched({
id: event.id,
id: SessionMessage.ID.fromEvent(event.id),
type: "agent-switched",
metadata: event.metadata,
agent: event.data.agent,
@@ -134,7 +135,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.model.switched": (event) => {
return adapter.appendMessage(
new SessionMessage.ModelSwitched({
id: event.id,
id: SessionMessage.ID.fromEvent(event.id),
type: "model-switched",
metadata: event.metadata,
model: event.data.model,
@@ -145,7 +146,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.prompted": (event) => {
return adapter.appendMessage(
new SessionMessage.User({
id: event.id,
id: SessionMessage.ID.fromEvent(event.id),
type: "user",
metadata: event.metadata,
text: event.data.prompt.text,
@@ -161,7 +162,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
new SessionMessage.Synthetic({
sessionID: event.data.sessionID,
text: event.data.text,
id: event.id,
id: SessionMessage.ID.fromEvent(event.id),
type: "synthetic",
time: { created: event.data.timestamp },
}),
@@ -170,7 +171,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.shell.started": (event) => {
return adapter.appendMessage(
new SessionMessage.Shell({
id: event.id,
id: SessionMessage.ID.fromEvent(event.id),
type: "shell",
metadata: event.metadata,
callID: event.data.callID,
@@ -205,7 +206,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
}
yield* adapter.appendMessage(
new SessionMessage.Assistant({
id: event.id,
id: SessionMessage.ID.fromEvent(event.id),
type: "assistant",
agent: event.data.agent,
model: event.data.model,
@@ -419,7 +420,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.compaction.started": (event) => {
return adapter.appendMessage(
new SessionMessage.Compaction({
id: event.id,
id: SessionMessage.ID.fromEvent(event.id),
type: "compaction",
metadata: event.metadata,
reason: event.data.reason,
+2 -2
View File
@@ -2,14 +2,14 @@ export * as SessionMessage from "./message"
import { Schema } from "effect"
import { ProviderMetadata } from "@opencode-ai/llm"
import { EventV2 } from "../event"
import { ModelV2 } from "../model"
import { ToolOutput } from "../tool-output"
import { V2Schema } from "../v2-schema"
import { SessionEvent } from "./event"
import { SessionMessageID } from "./message-id"
import { Prompt } from "./prompt"
export const ID = EventV2.ID
export const ID = SessionMessageID.ID
export type ID = Schema.Schema.Type<typeof ID>
const Base = {
+4 -3
View File
@@ -335,10 +335,11 @@ export const layer = Layer.effectDiscard(
)
yield* events.project(SessionEvent.Prompted, (event) =>
Effect.gen(function* () {
const messageID = SessionMessage.ID.fromEvent(event.id)
const existing = yield* db
.select({ id: SessionMessageTable.id })
.from(SessionMessageTable)
.where(eq(SessionMessageTable.id, event.id))
.where(eq(SessionMessageTable.id, messageID))
.get()
.pipe(Effect.orDie)
if (existing) return yield* Effect.die(new PromptAlreadyProjected())
@@ -346,7 +347,7 @@ export const layer = Layer.effectDiscard(
const row = yield* db
.select()
.from(SessionMessageTable)
.where(eq(SessionMessageTable.id, event.id))
.where(eq(SessionMessageTable.id, messageID))
.get()
.pipe(Effect.orDie)
if (!row) return yield* Effect.die("Prompt projection was not stored")
@@ -355,7 +356,7 @@ export const layer = Layer.effectDiscard(
if (event.seq === undefined)
return yield* Effect.die("Synchronized Session event is missing aggregate sequence")
yield* SessionInput.project(db, {
id: SessionMessage.ID.make(event.id),
id: messageID,
sessionID: event.data.sessionID,
prompt: event.data.prompt,
delivery: event.data.delivery,
+2 -1
View File
@@ -13,6 +13,7 @@ import { ToolRegistry } from "../../tool-registry"
import { SessionRunnerModel } from "./model"
import { Database } from "../../database/database"
import { SessionInput } from "../input"
import { SessionMessage } from "../message"
import { QuestionV2 } from "../../question"
/**
@@ -106,7 +107,7 @@ export const layer = Layer.effect(
yield* events.publish(SessionEvent.Tool.Failed, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: message.id,
assistantMessageID: SessionMessage.ID.toEvent(message.id),
callID: tool.id,
error: { type: "unknown", message: "Tool execution interrupted" },
provider: {
+97 -1
View File
@@ -7,9 +7,11 @@ import { EffectDrizzleSqlite } from "@opencode-ai/effect-drizzle-sqlite"
import { Effect, Layer } from "effect"
import { eq, inArray, sql } from "drizzle-orm"
import { DatabaseMigration } from "@opencode-ai/core/database/migration"
import { migrations } from "@opencode-ai/core/database/migration.gen"
import sessionUsageMigration from "@opencode-ai/core/database/migration/20260510033149_session_usage"
import normalizeStoragePathsMigration from "@opencode-ai/core/database/migration/20260601010001_normalize_storage_paths"
import sessionMessageProjectionOrderMigration from "@opencode-ai/core/database/migration/20260603040000_session_message_projection_order"
import sessionMessageIdentityMigration from "@opencode-ai/core/database/migration/20260604153000_session_message_identity"
import { ProjectV2 } from "@opencode-ai/core/project"
import { ProjectTable } from "@opencode-ai/core/project/sql"
import { AbsolutePath } from "@opencode-ai/core/schema"
@@ -62,7 +64,7 @@ describe("DatabaseMigration", () => {
expect(
yield* db.get(sql`SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'session_input'`),
).toEqual({ name: "session_input" })
expect(yield* db.get(sql`SELECT count(*) as count FROM migration`)).toEqual({ count: 29 })
expect(yield* db.get(sql`SELECT count(*) as count FROM migration`)).toEqual({ count: 30 })
expect(
yield* db.all(
sql`SELECT name FROM sqlite_master WHERE type = 'index' AND name IN ('event_aggregate_seq_idx', 'event_aggregate_type_seq_idx', 'session_input_session_pending_seq_idx', 'session_input_session_pending_delivery_seq_idx', 'session_message_session_idx', 'session_message_session_type_idx', 'session_message_session_seq_idx', 'session_message_session_type_seq_idx', 'session_message_session_time_created_id_idx') ORDER BY name`,
@@ -134,6 +136,100 @@ describe("DatabaseMigration", () => {
)
})
test("resets unreleased EventV2 state without deleting canonical Session history", async () => {
await run(
Effect.gen(function* () {
const db = yield* makeDb
yield* db.run(sql`CREATE TABLE session (id text PRIMARY KEY)`)
yield* db.run(sql`CREATE TABLE message (id text PRIMARY KEY, session_id text NOT NULL)`)
yield* db.run(sql`CREATE TABLE part (id text PRIMARY KEY, message_id text NOT NULL, session_id text NOT NULL)`)
yield* db.run(sql`CREATE TABLE session_message (id text PRIMARY KEY)`)
yield* db.run(sql`CREATE TABLE session_input (id text PRIMARY KEY)`)
yield* db.run(
sql`CREATE TABLE event_sequence (aggregate_id text PRIMARY KEY, seq integer NOT NULL, owner_id text)`,
)
yield* db.run(sql`CREATE TABLE event (id text PRIMARY KEY, aggregate_id text NOT NULL, seq integer NOT NULL)`)
yield* db.run(sql`INSERT INTO session (id) VALUES ('session')`)
yield* db.run(sql`INSERT INTO message (id, session_id) VALUES ('legacy_message', 'session')`)
yield* db.run(
sql`INSERT INTO part (id, message_id, session_id) VALUES ('legacy_part', 'legacy_message', 'session')`,
)
yield* db.run(sql`INSERT INTO session_message (id) VALUES ('experimental_message')`)
yield* db.run(sql`INSERT INTO session_input (id) VALUES ('experimental_input')`)
yield* db.run(sql`INSERT INTO event_sequence (aggregate_id, seq, owner_id) VALUES ('session', 1, 'workspace')`)
yield* db.run(sql`INSERT INTO event (id, aggregate_id, seq) VALUES ('experimental_event', 'session', 1)`)
yield* DatabaseMigration.applyOnly(db, [sessionMessageIdentityMigration])
expect(yield* db.all(sql`SELECT id FROM session`)).toEqual([{ id: "session" }])
expect(yield* db.all(sql`SELECT id FROM message`)).toEqual([{ id: "legacy_message" }])
expect(yield* db.all(sql`SELECT id FROM part`)).toEqual([{ id: "legacy_part" }])
expect(yield* db.all(sql`SELECT id FROM session_message`)).toEqual([])
expect(yield* db.all(sql`SELECT id FROM session_input`)).toEqual([])
expect(yield* db.all(sql`SELECT id FROM event`)).toEqual([])
expect(yield* db.all(sql`SELECT aggregate_id FROM event_sequence`)).toEqual([])
}),
)
})
test("applies the Session-message identity reset through normal database startup", async () => {
await using tmp = await tmpdir()
const filename = path.join(tmp.path, "identity-reset.sqlite")
const before = migrations.filter((migration) => migration.id !== sessionMessageIdentityMigration.id)
await Effect.runPromise(
Effect.gen(function* () {
const db = yield* makeDb
yield* db.run(sql`PRAGMA foreign_keys = ON`)
yield* DatabaseMigration.applyOnly(db, before)
yield* db.run(
sql`INSERT INTO project (id, worktree, sandboxes, time_created, time_updated) VALUES ('project', '/project', '[]', 1, 1)`,
)
yield* db.run(
sql`INSERT INTO session (id, project_id, slug, directory, title, version, time_created, time_updated) VALUES ('session', 'project', 'session', '/project', 'Session', 'test', 1, 1)`,
)
yield* db.run(
sql`INSERT INTO message (id, session_id, time_created, time_updated, data) VALUES ('legacy_message', 'session', 1, 1, '{}')`,
)
yield* db.run(
sql`INSERT INTO part (id, message_id, session_id, time_created, time_updated, data) VALUES ('legacy_part', 'legacy_message', 'session', 1, 1, '{}')`,
)
yield* db.run(
sql`INSERT INTO todo (session_id, content, status, priority, position, time_created, time_updated) VALUES ('session', 'keep', 'pending', 'low', 0, 1, 1)`,
)
yield* db.run(
sql`INSERT INTO session_message (id, session_id, type, seq, time_created, time_updated, data) VALUES ('evt_message', 'session', 'user', 1, 1, 1, '{}')`,
)
yield* db.run(
sql`INSERT INTO session_input (id, session_id, prompt, delivery, time_created) VALUES ('evt_input', 'session', '{}', 'queue', 1)`,
)
yield* db.run(sql`INSERT INTO event_sequence (aggregate_id, seq, owner_id) VALUES ('session', 1, 'workspace')`)
yield* db.run(
sql`INSERT INTO event (id, aggregate_id, seq, type, data) VALUES ('evt_event', 'session', 1, 'session.created.1', '{}')`,
)
}).pipe(Effect.provide(SqliteClient.layer({ filename, disableWAL: true })), Effect.scoped),
)
await Effect.runPromise(Effect.scoped(Layer.build(Database.layerFromPath(filename))))
await Effect.runPromise(
Effect.gen(function* () {
const db = yield* makeDb
expect(yield* db.all(sql`SELECT id FROM session`)).toEqual([{ id: "session" }])
expect(yield* db.all(sql`SELECT id FROM message`)).toEqual([{ id: "legacy_message" }])
expect(yield* db.all(sql`SELECT id FROM part`)).toEqual([{ id: "legacy_part" }])
expect(yield* db.all(sql`SELECT content FROM todo`)).toEqual([{ content: "keep" }])
expect(yield* db.all(sql`SELECT id FROM session_message`)).toEqual([])
expect(yield* db.all(sql`SELECT id FROM session_input`)).toEqual([])
expect(yield* db.all(sql`SELECT id FROM event`)).toEqual([])
expect(yield* db.all(sql`SELECT aggregate_id FROM event_sequence`)).toEqual([])
expect(yield* db.get(sql`SELECT id FROM migration WHERE id = ${sessionMessageIdentityMigration.id}`)).toEqual({
id: sessionMessageIdentityMigration.id,
})
}).pipe(Effect.provide(SqliteClient.layer({ filename, disableWAL: true })), Effect.scoped),
)
})
test("runs session usage backfill in order with schema changes", async () => {
await run(
Effect.gen(function* () {
+16
View File
@@ -7,6 +7,7 @@ import { Location } from "@opencode-ai/core/location"
import { AbsolutePath } from "@opencode-ai/core/schema"
import { WorkspaceV2 } from "@opencode-ai/core/workspace"
import { V2Schema } from "@opencode-ai/core/v2-schema"
import { SessionMessage } from "@opencode-ai/core/session/message"
import { eq } from "drizzle-orm"
import { location } from "./fixture/location"
import { testEffect } from "./lib/effect"
@@ -84,6 +85,21 @@ const SyncTimestamp = EventV2.define({
})
describe("EventV2", () => {
it.effect("keeps event IDs in the evt namespace", () =>
Effect.sync(() => {
expect(EventV2.ID.create()).toMatch(/^evt_/)
expect(() => EventV2.ID.make("msg_wrong_namespace")).toThrow()
expect(() => EventV2.ID.make("evtx")).toThrow()
}),
)
it.effect("round-trips Session message IDs through creator event IDs", () =>
Effect.sync(() => {
expect(String(SessionMessage.ID.fromEvent(EventV2.ID.make("evt_custom")))).toBe("msg_custom")
expect(String(SessionMessage.ID.toEvent(SessionMessage.ID.make("msg_custom")))).toBe("evt_custom")
}),
)
it.effect("derives stable namespaced external IDs", () =>
Effect.sync(() => {
const input = { namespace: "opencord.agent-input", key: "input-1" }
+100 -23
View File
@@ -67,13 +67,23 @@ describe("SessionProjector", () => {
yield* events.publish(
SessionEvent.Prompted,
{ sessionID, timestamp: created, prompt: new Prompt({ text: "first" }), delivery: "steer" },
{ id: SessionMessage.ID.make("evt_z") },
{
sessionID,
timestamp: created,
prompt: new Prompt({ text: "first" }),
delivery: "steer",
},
{ id: EventV2.ID.make("evt_z") },
)
yield* events.publish(
SessionEvent.Prompted,
{ sessionID, timestamp: created, prompt: new Prompt({ text: "second" }), delivery: "steer" },
{ id: SessionMessage.ID.make("evt_a") },
{
sessionID,
timestamp: created,
prompt: new Prompt({ text: "second" }),
delivery: "steer",
},
{ id: EventV2.ID.make("evt_a") },
)
const sessions = yield* SessionV2.Service
@@ -131,13 +141,13 @@ describe("SessionProjector", () => {
.run()
.pipe(Effect.orDie)
const events = yield* EventV2.Service
const id = SessionMessage.ID.make("evt_admitted")
const id = SessionMessage.ID.make("msg_admitted")
yield* SessionInput.admit(db, { id, sessionID, prompt: new Prompt({ text: "promote me" }), delivery: "steer" })
const event = yield* events.publish(
SessionEvent.Prompted,
{ sessionID, timestamp: created, prompt: new Prompt({ text: "promote me" }), delivery: "steer" },
{ id },
{ id: SessionMessage.ID.toEvent(id) },
)
expect(
@@ -168,9 +178,21 @@ describe("SessionProjector", () => {
.pipe(Effect.orDie)
const events = yield* EventV2.Service
yield* events.publish(SessionEvent.AgentSwitched, { sessionID, timestamp: created, agent: "build" })
yield* events.publish(SessionEvent.ModelSwitched, { sessionID, timestamp: created, model })
yield* events.publish(SessionEvent.Synthetic, { sessionID, timestamp: created, text: "synthetic context" })
yield* events.publish(SessionEvent.AgentSwitched, {
sessionID,
timestamp: created,
agent: "build",
})
yield* events.publish(SessionEvent.ModelSwitched, {
sessionID,
timestamp: created,
model,
})
yield* events.publish(SessionEvent.Synthetic, {
sessionID,
timestamp: created,
text: "synthetic context",
})
yield* events.publish(SessionEvent.Shell.Started, {
sessionID,
timestamp: created,
@@ -183,7 +205,11 @@ describe("SessionProjector", () => {
callID: "shell-1",
output: "/project",
})
yield* events.publish(SessionEvent.Compaction.Started, { sessionID, timestamp: created, reason: "manual" })
yield* events.publish(SessionEvent.Compaction.Started, {
sessionID,
timestamp: created,
reason: "manual",
})
yield* events.publish(SessionEvent.Compaction.Delta, { sessionID, timestamp: created, text: "partial" })
yield* events.publish(SessionEvent.Compaction.Ended, {
sessionID,
@@ -228,6 +254,53 @@ describe("SessionProjector", () => {
}),
)
it.effect("rejects a creator event that reuses an existing projected message ID", () =>
Effect.gen(function* () {
const { db } = yield* Database.Service
yield* db
.insert(ProjectTable)
.values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
.run()
.pipe(Effect.orDie)
yield* db
.insert(SessionTable)
.values({
id: sessionID,
project_id: Project.ID.global,
slug: "test",
directory: "/project",
title: "test",
version: "test",
})
.run()
.pipe(Effect.orDie)
const events = yield* EventV2.Service
yield* events.publish(
SessionEvent.Synthetic,
{ sessionID, timestamp: created, text: "first" },
{ id: EventV2.ID.make("evt_same") },
)
const duplicate = yield* events
.publish(
SessionEvent.Synthetic,
{ sessionID, timestamp: created, text: "second" },
{ id: EventV2.ID.make("evt_same") },
)
.pipe(Effect.exit)
expect(duplicate._tag).toBe("Failure")
expect(
yield* db
.select()
.from(SessionMessageTable)
.where(eq(SessionMessageTable.id, SessionMessage.ID.make("msg_same")))
.get()
.pipe(Effect.orDie),
).toMatchObject({ data: { text: "first" } })
}),
)
it.effect("rejects a Prompted event that conflicts with an admitted inbox row", () =>
Effect.gen(function* () {
const { db } = yield* Database.Service
@@ -249,14 +322,14 @@ describe("SessionProjector", () => {
.run()
.pipe(Effect.orDie)
const events = yield* EventV2.Service
const id = SessionMessage.ID.make("evt_conflict")
const id = SessionMessage.ID.make("msg_conflict")
yield* SessionInput.admit(db, { id, sessionID, prompt: new Prompt({ text: "admitted" }), delivery: "steer" })
const exit = yield* events
.publish(
SessionEvent.Prompted,
{ sessionID, timestamp: created, prompt: new Prompt({ text: "different" }), delivery: "steer" },
{ id },
{ id: SessionMessage.ID.toEvent(id) },
)
.pipe(Effect.exit)
@@ -288,12 +361,16 @@ describe("SessionProjector", () => {
.run()
.pipe(Effect.orDie)
const events = yield* EventV2.Service
const id = SessionMessage.ID.make("evt_delivery_conflict")
const id = SessionMessage.ID.make("msg_delivery_conflict")
const prompt = new Prompt({ text: "admitted" })
yield* SessionInput.admit(db, { id, sessionID, prompt, delivery: "queue" })
const exit = yield* events
.publish(SessionEvent.Prompted, { sessionID, timestamp: created, prompt, delivery: "steer" }, { id })
.publish(
SessionEvent.Prompted,
{ sessionID, timestamp: created, prompt, delivery: "steer" },
{ id: SessionMessage.ID.toEvent(id) },
)
.pipe(Effect.exit)
expect(String(exit)).toContain("Prompt projection conflicts with admitted input")
@@ -306,7 +383,7 @@ describe("SessionProjector", () => {
it.effect("does not revive a stale incomplete in-memory assistant projection", () =>
Effect.gen(function* () {
const stale = new SessionMessage.Assistant({
id: SessionMessage.ID.make("evt_assistant_stale"),
id: SessionMessage.ID.make("msg_assistant_stale"),
type: "assistant",
agent: "build",
model,
@@ -314,7 +391,7 @@ describe("SessionProjector", () => {
time: { created },
})
const completed = new SessionMessage.Assistant({
id: SessionMessage.ID.make("evt_assistant_completed"),
id: SessionMessage.ID.make("msg_assistant_completed"),
type: "assistant",
agent: "build",
model,
@@ -351,8 +428,8 @@ describe("SessionProjector", () => {
yield* db
.insert(SessionMessageTable)
.values([
assistantRow(SessionMessage.ID.make("evt_assistant_1"), 0),
assistantRow(SessionMessage.ID.make("evt_assistant_2"), 1),
assistantRow(SessionMessage.ID.make("msg_assistant_1"), 0),
assistantRow(SessionMessage.ID.make("msg_assistant_2"), 1),
])
.run()
.pipe(Effect.orDie)
@@ -361,7 +438,7 @@ describe("SessionProjector", () => {
yield* service.publish(SessionEvent.Step.Ended, {
sessionID,
timestamp: DateTime.makeUnsafe(1),
assistantMessageID: SessionMessage.ID.make("evt_assistant_2"),
assistantMessageID: EventV2.ID.make("evt_assistant_2"),
finish: "stop",
cost: 0,
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
@@ -409,8 +486,8 @@ describe("SessionProjector", () => {
yield* db
.insert(SessionMessageTable)
.values([
assistantRow(SessionMessage.ID.make("evt_assistant_stale"), 0),
assistantRow(SessionMessage.ID.make("evt_assistant_completed"), 1, {
assistantRow(SessionMessage.ID.make("msg_assistant_stale"), 0),
assistantRow(SessionMessage.ID.make("msg_assistant_completed"), 1, {
created: DateTime.makeUnsafe(1),
completed: DateTime.makeUnsafe(2),
}),
@@ -437,7 +514,7 @@ describe("SessionProjector", () => {
)
expect(messages).toEqual([
new SessionMessage.Assistant({
id: SessionMessage.ID.make("evt_assistant_completed"),
id: SessionMessage.ID.make("msg_assistant_completed"),
type: "assistant",
agent: "build",
model,
@@ -445,7 +522,7 @@ describe("SessionProjector", () => {
time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
}),
new SessionMessage.Assistant({
id: SessionMessage.ID.make("evt_assistant_stale"),
id: SessionMessage.ID.make("msg_assistant_stale"),
type: "assistant",
agent: "build",
model,
+40 -13
View File
@@ -310,7 +310,7 @@ describe("SessionV2.prompt", () => {
yield* events.publish(
SessionEvent.Prompted,
{ sessionID, timestamp: yield* DateTime.now, prompt, delivery: "steer" },
{ id: messageID },
{ id: SessionMessage.ID.toEvent(messageID) },
)
const retried = yield* session.prompt({ id: messageID, sessionID, prompt, resume: false })
@@ -329,7 +329,7 @@ describe("SessionV2.prompt", () => {
yield* events.publish(
SessionEvent.Prompted,
{ sessionID, timestamp: yield* DateTime.now, prompt, delivery: "queue" },
{ id: messageID },
{ id: SessionMessage.ID.toEvent(messageID) },
)
const retried = yield* session.prompt({ id: messageID, sessionID, prompt, delivery: "queue", resume: false })
@@ -339,7 +339,7 @@ describe("SessionV2.prompt", () => {
}),
)
it.effect("rejects an input ID already used by a durable non-prompt event", () =>
it.effect("rejects an input ID already used by a projected non-prompt message", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
@@ -347,7 +347,7 @@ describe("SessionV2.prompt", () => {
yield* events.publish(
SessionEvent.Synthetic,
{ sessionID, timestamp: yield* DateTime.now, text: "Collision" },
{ id: messageID },
{ id: SessionMessage.ID.toEvent(messageID) },
)
const failure = yield* session
@@ -359,7 +359,7 @@ describe("SessionV2.prompt", () => {
}),
)
it.effect("rejects a durable event ID reserved by an admitted prompt without poisoning promotion", () =>
it.effect("keeps event envelope IDs separate from admitted message IDs", () =>
Effect.gen(function* () {
yield* setup
const { db } = yield* Database.Service
@@ -368,24 +368,51 @@ describe("SessionV2.prompt", () => {
const prompt = new Prompt({ text: "Reserved prompt" })
yield* session.prompt({ id: messageID, sessionID, prompt, resume: false })
yield* events.publish(
SessionEvent.Synthetic,
{ sessionID, timestamp: yield* DateTime.now, text: "Synthetic" },
{ id: EventV2.ID.make("evt_reserved_prompt") },
)
expect(yield* admitted(messageID)).not.toHaveProperty("promotedSeq")
yield* SessionInput.promoteSteers(db, events, sessionID)
expect(yield* admitted(messageID)).toMatchObject({ promotedSeq: 1 })
expect(yield* session.messages({ sessionID })).toMatchObject([
{ id: messageID, type: "user", text: "Reserved prompt" },
{ type: "synthetic", text: "Synthetic" },
])
}),
)
it.effect("keeps Session message IDs in the msg namespace", () =>
Effect.sync(() => {
expect(SessionMessage.ID.create()).toMatch(/^msg_/)
expect(() => SessionMessage.ID.make("evt_wrong_namespace")).toThrow()
expect(() => SessionMessage.ID.make("msgx")).toThrow()
}),
)
it.effect("rejects a non-prompt event that reuses an admitted message ID", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
const events = yield* EventV2.Service
const prompt = new Prompt({ text: "Reserved prompt" })
yield* session.prompt({ id: messageID, sessionID, prompt, resume: false })
const failure = yield* events
.publish(
SessionEvent.Synthetic,
{ sessionID, timestamp: yield* DateTime.now, text: "Conflicting synthetic" },
{ id: messageID },
{ id: SessionMessage.ID.toEvent(messageID) },
)
.pipe(Effect.catchDefect(Effect.succeed))
expect(failure).toBe("Durable event conflicts with admitted prompt input")
expect(yield* admitted(messageID)).not.toHaveProperty("promotedSeq")
expect(yield* session.messages({ sessionID })).toEqual([])
yield* SessionInput.promoteSteers(db, events, sessionID)
expect(yield* admitted(messageID)).toMatchObject({ promotedSeq: 0 })
expect(yield* session.messages({ sessionID })).toMatchObject([
{ id: messageID, type: "user", text: "Reserved prompt" },
])
}),
)
@@ -1,7 +1,6 @@
import { describe, expect, test } from "bun:test"
import { Message, Model } from "@opencode-ai/llm"
import * as OpenAIChat from "@opencode-ai/llm/protocols/openai-chat"
import { EventV2 } from "@opencode-ai/core/event"
import { ModelV2 } from "@opencode-ai/core/model"
import { ProviderV2 } from "@opencode-ai/core/provider"
import { SessionMessage } from "@opencode-ai/core/session/message"
@@ -12,7 +11,7 @@ import { ToolOutput } from "@opencode-ai/core/tool-output"
import { DateTime } from "effect"
const created = DateTime.makeUnsafe(0)
const id = (value: string) => EventV2.ID.make(`evt_${value}`)
const id = (value: string) => SessionMessage.ID.make(`msg_${value}`)
const model = Model.make({ id: "model", provider: "provider", route: OpenAIChat.route })
describe("toLLMMessages", () => {
+13 -13
View File
@@ -1218,30 +1218,30 @@ describe("SessionRunnerLLM", () => {
const events = yield* EventV2.Service
yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Recover interrupted tool" }), resume: false })
yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID)
const assistant = yield* events.publish(SessionEvent.Step.Started, {
const assistantMessageID = (yield* events.publish(SessionEvent.Step.Started, {
sessionID,
timestamp: yield* DateTime.now,
agent: "build",
model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
})
})).id
yield* events.publish(SessionEvent.Tool.Input.Started, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: assistant.id,
assistantMessageID,
callID: "call-interrupted",
name: "echo",
})
yield* events.publish(SessionEvent.Tool.Input.Ended, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: assistant.id,
assistantMessageID,
callID: "call-interrupted",
text: '{"text":"stale"}',
})
yield* events.publish(SessionEvent.Tool.Called, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: assistant.id,
assistantMessageID,
callID: "call-interrupted",
tool: "echo",
input: { text: "stale" },
@@ -1280,30 +1280,30 @@ describe("SessionRunnerLLM", () => {
resume: false,
})
yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID)
const assistant = yield* events.publish(SessionEvent.Step.Started, {
const assistantMessageID = (yield* events.publish(SessionEvent.Step.Started, {
sessionID,
timestamp: yield* DateTime.now,
agent: "build",
model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
})
})).id
yield* events.publish(SessionEvent.Tool.Input.Started, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: assistant.id,
assistantMessageID,
callID: "call-hosted-interrupted",
name: "web_search",
})
yield* events.publish(SessionEvent.Tool.Input.Ended, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: assistant.id,
assistantMessageID,
callID: "call-hosted-interrupted",
text: '{"query":"stale"}',
})
yield* events.publish(SessionEvent.Tool.Called, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: assistant.id,
assistantMessageID,
callID: "call-hosted-interrupted",
tool: "web_search",
input: { query: "stale" },
@@ -1338,16 +1338,16 @@ describe("SessionRunnerLLM", () => {
resume: false,
})
yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID)
const assistant = yield* events.publish(SessionEvent.Step.Started, {
const assistantMessageID = (yield* events.publish(SessionEvent.Step.Started, {
sessionID,
timestamp: yield* DateTime.now,
agent: "build",
model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
})
})).id
yield* events.publish(SessionEvent.Tool.Input.Started, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: assistant.id,
assistantMessageID,
callID: "call-pending-interrupted",
name: "echo",
})
@@ -60,7 +60,7 @@ describe("Tool.Progress", () => {
const row = yield* db
.select()
.from(SessionMessageTable)
.where(eq(SessionMessageTable.id, assistantMessageID))
.where(eq(SessionMessageTable.id, SessionMessage.ID.fromEvent(assistantMessageID)))
.get()
.pipe(Effect.orDie)
if (!row) return yield* Effect.die("Missing projected assistant")
@@ -1,4 +1,5 @@
import { useEvent } from "@tui/context/event"
import { EventV2 } from "@opencode-ai/core/event"
import type {
SessionMessage,
SessionMessageAssistant,
@@ -9,6 +10,9 @@ import type {
import { createStore, produce, reconcile } from "solid-js/store"
import { createSimpleContext } from "./helper"
import { useSDK } from "./sdk"
import { SessionMessageID } from "@opencode-ai/core/session/message-id"
const messageID = (eventID: string) => SessionMessageID.ID.fromEvent(EventV2.ID.make(eventID))
function activeAssistant(messages: SessionMessage[]) {
const index = messages.findIndex((message) => message.type === "assistant" && !message.time.completed)
@@ -82,7 +86,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.agent.switched":
update(event.properties.sessionID, (draft) => {
draft.unshift({
id: event.id,
id: messageID(event.id),
type: "agent-switched",
agent: event.properties.agent,
time: { created: event.properties.timestamp },
@@ -92,7 +96,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.model.switched":
update(event.properties.sessionID, (draft) => {
draft.unshift({
id: event.id,
id: messageID(event.id),
type: "model-switched",
model: event.properties.model,
time: { created: event.properties.timestamp },
@@ -102,7 +106,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.prompted": {
update(event.properties.sessionID, (draft) => {
draft.unshift({
id: event.id,
id: messageID(event.id),
type: "user",
text: event.properties.prompt.text,
files: event.properties.prompt.files,
@@ -116,7 +120,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.synthetic":
update(event.properties.sessionID, (draft) => {
draft.unshift({
id: event.id,
id: messageID(event.id),
type: "synthetic",
sessionID: event.properties.sessionID,
text: event.properties.text,
@@ -127,7 +131,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.shell.started":
update(event.properties.sessionID, (draft) => {
draft.unshift({
id: event.id,
id: messageID(event.id),
type: "shell",
callID: event.properties.callID,
command: event.properties.command,
@@ -149,7 +153,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
const currentAssistant = activeAssistant(draft)
if (currentAssistant) currentAssistant.time.completed = event.properties.timestamp
draft.unshift({
id: event.id,
id: messageID(event.id),
type: "assistant",
agent: event.properties.agent,
model: event.properties.model,
@@ -161,7 +165,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
break
case "session.next.step.ended":
update(event.properties.sessionID, (draft) => {
const currentAssistant = ownedAssistant(draft, event.properties.assistantMessageID)
const currentAssistant = ownedAssistant(draft, messageID(event.properties.assistantMessageID))
if (!currentAssistant) return
currentAssistant.time.completed = event.properties.timestamp
currentAssistant.finish = event.properties.finish
@@ -173,7 +177,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
break
case "session.next.step.failed":
update(event.properties.sessionID, (draft) => {
const currentAssistant = ownedAssistant(draft, event.properties.assistantMessageID)
const currentAssistant = ownedAssistant(draft, messageID(event.properties.assistantMessageID))
if (!currentAssistant) return
currentAssistant.time.completed = event.properties.timestamp
currentAssistant.finish = "error"
@@ -199,7 +203,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
break
case "session.next.tool.input.started":
update(event.properties.sessionID, (draft) => {
ownedAssistant(draft, event.properties.assistantMessageID)?.content.push({
ownedAssistant(draft, messageID(event.properties.assistantMessageID))?.content.push({
type: "tool",
id: event.properties.callID,
name: event.properties.name,
@@ -211,7 +215,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.input.delta":
update(event.properties.sessionID, (draft) => {
const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID),
ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID,
)
if (match?.state.status === "pending") match.state.input += event.properties.delta
@@ -220,7 +224,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.input.ended":
update(event.properties.sessionID, (draft) => {
const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID),
ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID,
)
if (match?.state.status === "pending") match.state.input = event.properties.text
@@ -229,7 +233,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.called":
update(event.properties.sessionID, (draft) => {
const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID),
ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID,
)
if (!match) return
@@ -241,7 +245,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.progress":
update(event.properties.sessionID, (draft) => {
const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID),
ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID,
)
if (match?.state.status !== "running") return
@@ -252,7 +256,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.success":
update(event.properties.sessionID, (draft) => {
const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID),
ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID,
)
if (match?.state.status !== "running") return
@@ -270,7 +274,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.failed":
update(event.properties.sessionID, (draft) => {
const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID),
ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID,
)
if (!match || (match.state.status !== "pending" && match.state.status !== "running")) return
@@ -317,7 +321,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.compaction.started":
update(event.properties.sessionID, (draft) => {
draft.unshift({
id: event.id,
id: messageID(event.id),
type: "compaction",
reason: event.properties.reason,
summary: "",
@@ -1,4 +1,3 @@
import { Config } from "@/config/config"
import { ConfigV1 } from "@opencode-ai/core/v1/config/config"
import { EventV2 } from "@opencode-ai/core/event"
import { InstanceDisposed } from "@/server/event"
@@ -20,10 +19,10 @@ const SyncEventSchemas = EventV2.registry
return [
Schema.Struct({
type: Schema.Literal("sync"),
id: Schema.String,
id: EventV2.ID,
syncEvent: Schema.Struct({
type: Schema.Literal(EventV2.versionedType(definition.type, definition.sync.version)),
id: Schema.String,
id: EventV2.ID,
seq: Schema.Finite,
aggregateID: Schema.String,
data: definition.data,
@@ -41,7 +40,7 @@ const GlobalEventSchema = Schema.Struct({
...EventV2.registry
.values()
.map((definition) =>
Schema.Struct({ id: Schema.String, type: Schema.Literal(definition.type), properties: definition.data }),
Schema.Struct({ id: EventV2.ID, type: Schema.Literal(definition.type), properties: definition.data }),
)
.toArray(),
InstanceDisposed,
@@ -1,4 +1,5 @@
import { NonNegativeInt } from "@opencode-ai/core/schema"
import { EventV2 } from "@opencode-ai/core/event"
import { SessionID } from "@/session/schema"
import { Schema } from "effect"
import { HttpApi, HttpApiEndpoint, HttpApiError, HttpApiGroup, OpenApi } from "effect/unstable/httpapi"
@@ -9,7 +10,7 @@ import { described } from "./metadata"
const root = "/sync"
export const ReplayEvent = Schema.Struct({
id: Schema.String,
id: EventV2.ID,
aggregateID: Schema.String,
seq: NonNegativeInt,
type: Schema.String,
@@ -27,7 +28,7 @@ export const SessionPayload = Schema.Struct({
})
export const HistoryPayload = Schema.Record(Schema.String, NonNegativeInt)
export const HistoryEvent = Schema.Struct({
id: Schema.String,
id: EventV2.ID,
aggregate_id: Schema.String,
seq: NonNegativeInt,
type: Schema.String,
+3 -2
View File
@@ -6,7 +6,6 @@ import type { Agent } from "@/agent/agent"
import { EventV2Bridge } from "@/event-v2-bridge"
import { InstanceState } from "@/effect/instance-state"
import { Global } from "@opencode-ai/core/global"
import { SkillPlugin } from "@opencode-ai/core/plugin/skill"
import { Permission } from "@/permission"
import { FSUtil } from "@opencode-ai/core/fs-util"
import { Config } from "@/config/config"
@@ -33,7 +32,9 @@ const SKILL_PATTERN = "**/SKILL.md"
const CUSTOMIZE_OPENCODE_SKILL_NAME = "customize-opencode"
const CUSTOMIZE_OPENCODE_SKILL_DESCRIPTION =
"Use ONLY when the user is editing or creating opencode's own configuration: opencode.json, opencode.jsonc, files under .opencode/, or files under ~/.config/opencode/. Also use when creating or fixing opencode agents, subagents, skills, plugins, MCP servers, or permission rules. Do not use for the user's own application code, or for any project that is not configuring opencode itself."
const CUSTOMIZE_OPENCODE_SKILL_BODY = SkillPlugin.CustomizeOpencodeContent
const CUSTOMIZE_OPENCODE_SKILL_BODY = await Bun.file(
new URL("../../../core/src/plugin/skill/customize-opencode.md", import.meta.url),
).text()
export const Info = Schema.Struct({
name: Schema.String,
+8 -21
View File
@@ -6,14 +6,8 @@ export function listen(rpc: Definition) {
onmessage = async (evt) => {
const parsed = JSON.parse(evt.data)
if (parsed.type === "rpc.request") {
try {
const result = await rpc[parsed.method](parsed.input)
postMessage(JSON.stringify({ type: "rpc.result", result, id: parsed.id }))
} catch (error) {
postMessage(
JSON.stringify({ type: "rpc.error", error: error instanceof Error ? error.message : String(error), id: parsed.id }),
)
}
const result = await rpc[parsed.method](parsed.input)
postMessage(JSON.stringify({ type: "rpc.result", result, id: parsed.id }))
}
}
}
@@ -26,22 +20,15 @@ export function client<T extends Definition>(target: {
postMessage: (data: string) => void | null
onmessage: ((this: Worker, ev: MessageEvent<any>) => any) | null
}) {
const pending = new Map<number, { resolve: (result: any) => void; reject: (error: any) => void }>()
const pending = new Map<number, (result: any) => void>()
const listeners = new Map<string, Set<(data: any) => void>>()
let id = 0
target.onmessage = async (evt) => {
const parsed = JSON.parse(evt.data)
if (parsed.type === "rpc.result") {
const request = pending.get(parsed.id)
if (request) {
request.resolve(parsed.result)
pending.delete(parsed.id)
}
}
if (parsed.type === "rpc.error") {
const request = pending.get(parsed.id)
if (request) {
request.reject(new Error(parsed.error))
const resolve = pending.get(parsed.id)
if (resolve) {
resolve(parsed.result)
pending.delete(parsed.id)
}
}
@@ -57,8 +44,8 @@ export function client<T extends Definition>(target: {
return {
call<Method extends keyof T>(method: Method, input: Parameters<T[Method]>[0]): Promise<ReturnType<T[Method]>> {
const requestId = id++
return new Promise((resolve, reject) => {
pending.set(requestId, { resolve, reject })
return new Promise((resolve) => {
pending.set(requestId, resolve)
target.postMessage(JSON.stringify({ type: "rpc.request", method, input, id: requestId }))
})
},
@@ -49,14 +49,14 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
await mounted
events.emit(
global({
id: "agent-1",
id: "evt_agent_1",
type: "session.next.agent.switched",
properties: { sessionID: "session-1", timestamp: 0, agent: "build" },
}),
)
events.emit(
global({
id: "model-1",
id: "evt_model_1",
type: "session.next.model.switched",
properties: {
sessionID: "session-1",
@@ -67,7 +67,7 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
)
events.emit(
global({
id: "assistant-1",
id: "evt_assistant_1",
type: "session.next.step.started",
properties: {
sessionID: "session-1",
@@ -79,12 +79,12 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
)
events.emit(
global({
id: "input-1",
id: "evt_input_1",
type: "session.next.tool.input.started",
properties: {
sessionID: "session-1",
timestamp: 2,
assistantMessageID: "assistant-1",
assistantMessageID: "evt_assistant_1",
callID: "call-1",
name: "bash",
},
@@ -92,12 +92,12 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
)
events.emit(
global({
id: "failed-1",
id: "evt_failed_1",
type: "session.next.tool.failed",
properties: {
sessionID: "session-1",
timestamp: 3,
assistantMessageID: "assistant-1",
assistantMessageID: "evt_assistant_1",
callID: "call-1",
error: { type: "unknown", message: "aborted" },
provider: { executed: false },
@@ -117,6 +117,7 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
const assistant = sync.session.message.fromSession("session-1")[0]
expect(assistant?.type).toBe("assistant")
if (assistant?.type !== "assistant") return
expect(assistant.id).toBe("msg_assistant_1")
const tool = assistant.content[0]
expect(tool?.type).toBe("tool")
if (tool?.type !== "tool") return
@@ -22,7 +22,6 @@ import * as HttpSessionError from "../../src/server/routes/instance/httpapi/hand
import { SessionPaths } from "../../src/server/routes/instance/httpapi/groups/session"
import { Session } from "@/session/session"
import { MessageID, PartID, SessionID, type SessionID as SessionIDType } from "../../src/session/schema"
import { MessageV2 } from "../../src/session/message-v2"
import { Database } from "@opencode-ai/core/database/database"
import { SessionInputTable, SessionMessageTable, SessionTable } from "@opencode-ai/core/session/sql"
import { SessionMessage } from "@opencode-ai/core/session/message"
@@ -584,7 +583,7 @@ describe("session HttpApi", () => {
request(`/api/session/${session.id}/prompt`, {
method: "POST",
headers: { ...headers, "content-type": "application/json" },
body: JSON.stringify({ id: "evt_http_prompt", prompt: { text: "hello" } }),
body: JSON.stringify({ id: "msg_http_prompt", prompt: { text: "hello" } }),
})
const first = yield* recordPrompt()
const retried = yield* recordPrompt()
@@ -604,12 +603,12 @@ describe("session HttpApi", () => {
db
.select()
.from(SessionInputTable)
.where(eq(SessionInputTable.id, SessionMessage.ID.make("evt_http_prompt")))
.where(eq(SessionInputTable.id, SessionMessage.ID.make("msg_http_prompt")))
.get()
.pipe(Effect.orDie),
)
expect(admitted).toMatchObject({
id: "evt_http_prompt",
id: "msg_http_prompt",
session_id: session.id,
delivery: "steer",
promoted_seq: null,
@@ -618,13 +617,13 @@ describe("session HttpApi", () => {
const conflict = yield* request(`/api/session/${session.id}/prompt`, {
method: "POST",
headers: { ...headers, "content-type": "application/json" },
body: JSON.stringify({ id: "evt_http_prompt", prompt: { text: "goodbye" } }),
body: JSON.stringify({ id: "msg_http_prompt", prompt: { text: "goodbye" } }),
})
expect(conflict.status).toBe(409)
expect(yield* responseJson(conflict)).toEqual({
_tag: "ConflictError",
message: "Prompt message ID conflicts with an existing durable record: evt_http_prompt",
resource: "evt_http_prompt",
message: "Prompt message ID conflicts with an existing durable record: msg_http_prompt",
resource: "msg_http_prompt",
})
}),
{ git: true, config: { formatter: false, lsp: false } },
-31
View File
@@ -1,31 +0,0 @@
import { describe, expect, test } from "bun:test"
import { Rpc } from "@/util/rpc"
type TestRpc = {
fail(input: undefined): Promise<void>
}
type Target = Parameters<typeof Rpc.client<TestRpc>>[0]
describe("Rpc", () => {
test("rejects pending calls when the worker reports an error", async () => {
const target: Target = {
postMessage(data) {
const request = JSON.parse(data)
target.onmessage?.call(
{} as Worker,
{
data: JSON.stringify({
type: "rpc.error",
id: request.id,
error: "boom",
}),
} as MessageEvent<any>,
)
},
onmessage: null,
}
await expect(Rpc.client<TestRpc>(target).call("fail", undefined)).rejects.toThrow("boom")
})
})
@@ -7,6 +7,7 @@ import { ModelV2 } from "@opencode-ai/core/model"
import { ProviderV2 } from "@opencode-ai/core/provider"
import { SessionEvent } from "@opencode-ai/core/session/event"
import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater"
import { SessionMessage } from "@opencode-ai/core/session/message"
import { ToolOutput } from "@opencode-ai/core/tool-output"
test.skip("step snapshots carry over to assistant messages", () => {
@@ -62,10 +63,11 @@ test.skip("step snapshots carry over to assistant messages", () => {
test.skip("text ended populates assistant text content", () => {
const state: SessionMessageUpdater.MemoryState = { messages: [] }
const sessionID = SessionID.make("session")
const assistantMessageID = EventV2.ID.create()
Effect.runSync(
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
id: assistantMessageID,
type: "session.next.step.started",
data: {
sessionID,
@@ -243,7 +245,7 @@ test.skip("compaction events reduce to compaction message", () => {
expect(state.messages).toHaveLength(1)
expect(state.messages[0]).toMatchObject({
id,
id: SessionMessage.ID.fromEvent(id),
type: "compaction",
reason: "auto",
summary: "final summary",
+356 -178
View File
File diff suppressed because it is too large Load Diff
+21
View File
@@ -6,6 +6,27 @@ This document covers meaningful contract changes introduced on the `feat/opencod
## Earlier Branch History
### Independent Session Message Identity
Affected schema:
- V2 Session-message IDs and pre-launch `session_input` and `session_message` rows.
Change:
- Separate mutable Session-message identity (`msg_*`) from immutable event-envelope identity (`evt_*`).
- Derive projected `msg_*` IDs reversibly from creator `evt_*` IDs without changing synchronized event payloads.
Reason:
- Event IDs identify immutable facts. Message IDs identify mutable timeline rows. Reusing one ID hid that boundary and leaked `evt_*` IDs into projected messages.
- Clients can generate user-message IDs before prompt admission for optimistic rendering and exact retry. Promotion reverses that ID into the creator event ID, and projection derives the original message ID again.
Compatibility:
- The migration resets unreleased workspace-sync events, sequence state, V2 inbox rows, and V2 timeline projections. Canonical V1 `session`, `message`, and `part` history remains untouched.
- Existing Sessions intentionally lose workspace-warp replayability across this pre-launch cutover until new replayable history is recorded.
### Replayable Session Event Refinement And Cursor Stream
Affected schema: