Compare commits

...
Sign in to create a new pull request.

1 commit

Author SHA1 Message Date
Kit Langton
f2d133edf7 refactor(core): separate v2 message identity 2026-06-04 12:08:00 -04:00
26 changed files with 772 additions and 289 deletions

23
packages/core/AGENTS.md Normal file
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.

View file

@ -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`;

View file

@ -31,5 +31,6 @@ export const migrations = (
import("./migration/20260603040000_session_message_projection_order"), import("./migration/20260603040000_session_message_projection_order"),
import("./migration/20260603141458_session_input_inbox"), import("./migration/20260603141458_session_input_inbox"),
import("./migration/20260603160727_jittery_ezekiel_stane"), import("./migration/20260603160727_jittery_ezekiel_stane"),
import("./migration/20260604153000_session_message_identity"),
]) ])
).map((module) => module.default) satisfies DatabaseMigration.Migration[] ).map((module) => module.default) satisfies DatabaseMigration.Migration[]

View file

@ -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

View file

@ -8,7 +8,7 @@ import { Location } from "./location"
import { externalID, type ExternalID, NonNegativeInt, withStatics } from "./schema" import { externalID, type ExternalID, NonNegativeInt, withStatics } from "./schema"
import { Identifier } from "./util/identifier" 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"), Schema.brand("Event.ID"),
withStatics((schema) => ({ withStatics((schema) => ({
create: () => schema.make("evt_" + Identifier.ascending()), create: () => schema.make("evt_" + Identifier.ascending()),

View file

@ -3,7 +3,7 @@ export * as SessionInput from "./input"
import { and, asc, eq, inArray, isNull } from "drizzle-orm" import { and, asc, eq, inArray, isNull } from "drizzle-orm"
import { DateTime, Effect, Schema } from "effect" import { DateTime, Effect, Schema } from "effect"
import type { Database } from "../database/database" import type { Database } from "../database/database"
import type { EventV2 } from "../event" import { EventV2 } from "../event"
import { EventTable } from "../event/sql" import { EventTable } from "../event/sql"
import { NonNegativeInt, PositiveInt } from "../schema" import { NonNegativeInt, PositiveInt } from "../schema"
import { V2Schema } from "../v2-schema" import { V2Schema } from "../v2-schema"
@ -66,7 +66,7 @@ export const admit = Effect.fn("SessionInput.admit")(function* (
const event = yield* db const event = yield* db
.select({ id: EventTable.id }) .select({ id: EventTable.id })
.from(EventTable) .from(EventTable)
.where(eq(EventTable.id, input.id)) .where(eq(EventTable.id, SessionMessage.ID.toEvent(input.id)))
.get() .get()
.pipe(Effect.orDie) .pipe(Effect.orDie)
const message = yield* db const message = yield* db
@ -133,7 +133,7 @@ export const guardReservedID = Effect.fn("SessionInput.guardReservedID")(functio
db: DatabaseService, db: DatabaseService,
event: EventV2.Payload, 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 (admitted === undefined) return
if (!Schema.is(SessionEvent.Prompted)(event)) if (!Schema.is(SessionEvent.Prompted)(event))
return yield* Effect.die("Durable event conflicts with admitted prompt input") 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), prompt: decodePrompt(row.prompt),
delivery: row.delivery, delivery: row.delivery,
}, },
{ id: SessionMessage.ID.make(row.id) }, { id: SessionMessage.ID.toEvent(SessionMessage.ID.make(row.id)) },
) )
} }
return rows.length return rows.length

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

View file

@ -1,5 +1,6 @@
import { castDraft, produce, type WritableDraft } from "immer" import { castDraft, produce, type WritableDraft } from "immer"
import { Effect } from "effect" import { Effect } from "effect"
import type { EventV2 } from "../event"
import { SessionEvent } from "./event" import { SessionEvent } from "./event"
import { SessionMessage } from "./message" import { SessionMessage } from "./message"
@ -112,9 +113,9 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
const latestReasoning = (assistant: DraftAssistant | undefined, reasoningID: string) => const latestReasoning = (assistant: DraftAssistant | undefined, reasoningID: string) =>
assistant?.content.findLast((item): item is DraftReasoning => item.type === "reasoning" && item.id === reasoningID) 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* () { 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)) 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) => { "session.next.agent.switched": (event) => {
return adapter.appendMessage( return adapter.appendMessage(
new SessionMessage.AgentSwitched({ new SessionMessage.AgentSwitched({
id: event.id, id: SessionMessage.ID.fromEvent(event.id),
type: "agent-switched", type: "agent-switched",
metadata: event.metadata, metadata: event.metadata,
agent: event.data.agent, agent: event.data.agent,
@ -134,7 +135,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.model.switched": (event) => { "session.next.model.switched": (event) => {
return adapter.appendMessage( return adapter.appendMessage(
new SessionMessage.ModelSwitched({ new SessionMessage.ModelSwitched({
id: event.id, id: SessionMessage.ID.fromEvent(event.id),
type: "model-switched", type: "model-switched",
metadata: event.metadata, metadata: event.metadata,
model: event.data.model, model: event.data.model,
@ -145,7 +146,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.prompted": (event) => { "session.next.prompted": (event) => {
return adapter.appendMessage( return adapter.appendMessage(
new SessionMessage.User({ new SessionMessage.User({
id: event.id, id: SessionMessage.ID.fromEvent(event.id),
type: "user", type: "user",
metadata: event.metadata, metadata: event.metadata,
text: event.data.prompt.text, text: event.data.prompt.text,
@ -161,7 +162,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
new SessionMessage.Synthetic({ new SessionMessage.Synthetic({
sessionID: event.data.sessionID, sessionID: event.data.sessionID,
text: event.data.text, text: event.data.text,
id: event.id, id: SessionMessage.ID.fromEvent(event.id),
type: "synthetic", type: "synthetic",
time: { created: event.data.timestamp }, time: { created: event.data.timestamp },
}), }),
@ -170,7 +171,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.shell.started": (event) => { "session.next.shell.started": (event) => {
return adapter.appendMessage( return adapter.appendMessage(
new SessionMessage.Shell({ new SessionMessage.Shell({
id: event.id, id: SessionMessage.ID.fromEvent(event.id),
type: "shell", type: "shell",
metadata: event.metadata, metadata: event.metadata,
callID: event.data.callID, callID: event.data.callID,
@ -205,7 +206,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
} }
yield* adapter.appendMessage( yield* adapter.appendMessage(
new SessionMessage.Assistant({ new SessionMessage.Assistant({
id: event.id, id: SessionMessage.ID.fromEvent(event.id),
type: "assistant", type: "assistant",
agent: event.data.agent, agent: event.data.agent,
model: event.data.model, model: event.data.model,
@ -419,7 +420,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.compaction.started": (event) => { "session.next.compaction.started": (event) => {
return adapter.appendMessage( return adapter.appendMessage(
new SessionMessage.Compaction({ new SessionMessage.Compaction({
id: event.id, id: SessionMessage.ID.fromEvent(event.id),
type: "compaction", type: "compaction",
metadata: event.metadata, metadata: event.metadata,
reason: event.data.reason, reason: event.data.reason,

View file

@ -2,14 +2,14 @@ export * as SessionMessage from "./message"
import { Schema } from "effect" import { Schema } from "effect"
import { ProviderMetadata } from "@opencode-ai/llm" import { ProviderMetadata } from "@opencode-ai/llm"
import { EventV2 } from "../event"
import { ModelV2 } from "../model" import { ModelV2 } from "../model"
import { ToolOutput } from "../tool-output" import { ToolOutput } from "../tool-output"
import { V2Schema } from "../v2-schema" import { V2Schema } from "../v2-schema"
import { SessionEvent } from "./event" import { SessionEvent } from "./event"
import { SessionMessageID } from "./message-id"
import { Prompt } from "./prompt" import { Prompt } from "./prompt"
export const ID = EventV2.ID export const ID = SessionMessageID.ID
export type ID = Schema.Schema.Type<typeof ID> export type ID = Schema.Schema.Type<typeof ID>
const Base = { const Base = {

View file

@ -335,10 +335,11 @@ export const layer = Layer.effectDiscard(
) )
yield* events.project(SessionEvent.Prompted, (event) => yield* events.project(SessionEvent.Prompted, (event) =>
Effect.gen(function* () { Effect.gen(function* () {
const messageID = SessionMessage.ID.fromEvent(event.id)
const existing = yield* db const existing = yield* db
.select({ id: SessionMessageTable.id }) .select({ id: SessionMessageTable.id })
.from(SessionMessageTable) .from(SessionMessageTable)
.where(eq(SessionMessageTable.id, event.id)) .where(eq(SessionMessageTable.id, messageID))
.get() .get()
.pipe(Effect.orDie) .pipe(Effect.orDie)
if (existing) return yield* Effect.die(new PromptAlreadyProjected()) if (existing) return yield* Effect.die(new PromptAlreadyProjected())
@ -346,7 +347,7 @@ export const layer = Layer.effectDiscard(
const row = yield* db const row = yield* db
.select() .select()
.from(SessionMessageTable) .from(SessionMessageTable)
.where(eq(SessionMessageTable.id, event.id)) .where(eq(SessionMessageTable.id, messageID))
.get() .get()
.pipe(Effect.orDie) .pipe(Effect.orDie)
if (!row) return yield* Effect.die("Prompt projection was not stored") if (!row) return yield* Effect.die("Prompt projection was not stored")
@ -355,7 +356,7 @@ export const layer = Layer.effectDiscard(
if (event.seq === undefined) if (event.seq === undefined)
return yield* Effect.die("Synchronized Session event is missing aggregate sequence") return yield* Effect.die("Synchronized Session event is missing aggregate sequence")
yield* SessionInput.project(db, { yield* SessionInput.project(db, {
id: SessionMessage.ID.make(event.id), id: messageID,
sessionID: event.data.sessionID, sessionID: event.data.sessionID,
prompt: event.data.prompt, prompt: event.data.prompt,
delivery: event.data.delivery, delivery: event.data.delivery,

View file

@ -13,6 +13,7 @@ import { ToolRegistry } from "../../tool-registry"
import { SessionRunnerModel } from "./model" import { SessionRunnerModel } from "./model"
import { Database } from "../../database/database" import { Database } from "../../database/database"
import { SessionInput } from "../input" import { SessionInput } from "../input"
import { SessionMessage } from "../message"
import { QuestionV2 } from "../../question" import { QuestionV2 } from "../../question"
/** /**
@ -106,7 +107,7 @@ export const layer = Layer.effect(
yield* events.publish(SessionEvent.Tool.Failed, { yield* events.publish(SessionEvent.Tool.Failed, {
sessionID, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
assistantMessageID: message.id, assistantMessageID: SessionMessage.ID.toEvent(message.id),
callID: tool.id, callID: tool.id,
error: { type: "unknown", message: "Tool execution interrupted" }, error: { type: "unknown", message: "Tool execution interrupted" },
provider: { provider: {

View file

@ -7,9 +7,11 @@ import { EffectDrizzleSqlite } from "@opencode-ai/effect-drizzle-sqlite"
import { Effect, Layer } from "effect" import { Effect, Layer } from "effect"
import { eq, inArray, sql } from "drizzle-orm" import { eq, inArray, sql } from "drizzle-orm"
import { DatabaseMigration } from "@opencode-ai/core/database/migration" 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 sessionUsageMigration from "@opencode-ai/core/database/migration/20260510033149_session_usage"
import normalizeStoragePathsMigration from "@opencode-ai/core/database/migration/20260601010001_normalize_storage_paths" 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 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 { ProjectV2 } from "@opencode-ai/core/project"
import { ProjectTable } from "@opencode-ai/core/project/sql" import { ProjectTable } from "@opencode-ai/core/project/sql"
import { AbsolutePath } from "@opencode-ai/core/schema" import { AbsolutePath } from "@opencode-ai/core/schema"
@ -62,7 +64,7 @@ describe("DatabaseMigration", () => {
expect( expect(
yield* db.get(sql`SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'session_input'`), yield* db.get(sql`SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'session_input'`),
).toEqual({ 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( expect(
yield* db.all( 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`, 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 () => { test("runs session usage backfill in order with schema changes", async () => {
await run( await run(
Effect.gen(function* () { Effect.gen(function* () {

View file

@ -7,6 +7,7 @@ import { Location } from "@opencode-ai/core/location"
import { AbsolutePath } from "@opencode-ai/core/schema" import { AbsolutePath } from "@opencode-ai/core/schema"
import { WorkspaceV2 } from "@opencode-ai/core/workspace" import { WorkspaceV2 } from "@opencode-ai/core/workspace"
import { V2Schema } from "@opencode-ai/core/v2-schema" import { V2Schema } from "@opencode-ai/core/v2-schema"
import { SessionMessage } from "@opencode-ai/core/session/message"
import { eq } from "drizzle-orm" import { eq } from "drizzle-orm"
import { location } from "./fixture/location" import { location } from "./fixture/location"
import { testEffect } from "./lib/effect" import { testEffect } from "./lib/effect"
@ -84,6 +85,21 @@ const SyncTimestamp = EventV2.define({
}) })
describe("EventV2", () => { 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", () => it.effect("derives stable namespaced external IDs", () =>
Effect.sync(() => { Effect.sync(() => {
const input = { namespace: "opencord.agent-input", key: "input-1" } const input = { namespace: "opencord.agent-input", key: "input-1" }

View file

@ -67,13 +67,23 @@ describe("SessionProjector", () => {
yield* events.publish( yield* events.publish(
SessionEvent.Prompted, 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( yield* events.publish(
SessionEvent.Prompted, 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 const sessions = yield* SessionV2.Service
@ -131,13 +141,13 @@ describe("SessionProjector", () => {
.run() .run()
.pipe(Effect.orDie) .pipe(Effect.orDie)
const events = yield* EventV2.Service 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" }) yield* SessionInput.admit(db, { id, sessionID, prompt: new Prompt({ text: "promote me" }), delivery: "steer" })
const event = yield* events.publish( const event = yield* events.publish(
SessionEvent.Prompted, SessionEvent.Prompted,
{ sessionID, timestamp: created, prompt: new Prompt({ text: "promote me" }), delivery: "steer" }, { sessionID, timestamp: created, prompt: new Prompt({ text: "promote me" }), delivery: "steer" },
{ id }, { id: SessionMessage.ID.toEvent(id) },
) )
expect( expect(
@ -168,9 +178,21 @@ describe("SessionProjector", () => {
.pipe(Effect.orDie) .pipe(Effect.orDie)
const events = yield* EventV2.Service const events = yield* EventV2.Service
yield* events.publish(SessionEvent.AgentSwitched, { sessionID, timestamp: created, agent: "build" }) yield* events.publish(SessionEvent.AgentSwitched, {
yield* events.publish(SessionEvent.ModelSwitched, { sessionID, timestamp: created, model }) sessionID,
yield* events.publish(SessionEvent.Synthetic, { sessionID, timestamp: created, text: "synthetic context" }) 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, { yield* events.publish(SessionEvent.Shell.Started, {
sessionID, sessionID,
timestamp: created, timestamp: created,
@ -183,7 +205,11 @@ describe("SessionProjector", () => {
callID: "shell-1", callID: "shell-1",
output: "/project", 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.Delta, { sessionID, timestamp: created, text: "partial" })
yield* events.publish(SessionEvent.Compaction.Ended, { yield* events.publish(SessionEvent.Compaction.Ended, {
sessionID, 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", () => it.effect("rejects a Prompted event that conflicts with an admitted inbox row", () =>
Effect.gen(function* () { Effect.gen(function* () {
const { db } = yield* Database.Service const { db } = yield* Database.Service
@ -249,14 +322,14 @@ describe("SessionProjector", () => {
.run() .run()
.pipe(Effect.orDie) .pipe(Effect.orDie)
const events = yield* EventV2.Service 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" }) yield* SessionInput.admit(db, { id, sessionID, prompt: new Prompt({ text: "admitted" }), delivery: "steer" })
const exit = yield* events const exit = yield* events
.publish( .publish(
SessionEvent.Prompted, SessionEvent.Prompted,
{ sessionID, timestamp: created, prompt: new Prompt({ text: "different" }), delivery: "steer" }, { sessionID, timestamp: created, prompt: new Prompt({ text: "different" }), delivery: "steer" },
{ id }, { id: SessionMessage.ID.toEvent(id) },
) )
.pipe(Effect.exit) .pipe(Effect.exit)
@ -288,12 +361,16 @@ describe("SessionProjector", () => {
.run() .run()
.pipe(Effect.orDie) .pipe(Effect.orDie)
const events = yield* EventV2.Service 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" }) const prompt = new Prompt({ text: "admitted" })
yield* SessionInput.admit(db, { id, sessionID, prompt, delivery: "queue" }) yield* SessionInput.admit(db, { id, sessionID, prompt, delivery: "queue" })
const exit = yield* events 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) .pipe(Effect.exit)
expect(String(exit)).toContain("Prompt projection conflicts with admitted input") 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", () => it.effect("does not revive a stale incomplete in-memory assistant projection", () =>
Effect.gen(function* () { Effect.gen(function* () {
const stale = new SessionMessage.Assistant({ const stale = new SessionMessage.Assistant({
id: SessionMessage.ID.make("evt_assistant_stale"), id: SessionMessage.ID.make("msg_assistant_stale"),
type: "assistant", type: "assistant",
agent: "build", agent: "build",
model, model,
@ -314,7 +391,7 @@ describe("SessionProjector", () => {
time: { created }, time: { created },
}) })
const completed = new SessionMessage.Assistant({ const completed = new SessionMessage.Assistant({
id: SessionMessage.ID.make("evt_assistant_completed"), id: SessionMessage.ID.make("msg_assistant_completed"),
type: "assistant", type: "assistant",
agent: "build", agent: "build",
model, model,
@ -351,8 +428,8 @@ describe("SessionProjector", () => {
yield* db yield* db
.insert(SessionMessageTable) .insert(SessionMessageTable)
.values([ .values([
assistantRow(SessionMessage.ID.make("evt_assistant_1"), 0), assistantRow(SessionMessage.ID.make("msg_assistant_1"), 0),
assistantRow(SessionMessage.ID.make("evt_assistant_2"), 1), assistantRow(SessionMessage.ID.make("msg_assistant_2"), 1),
]) ])
.run() .run()
.pipe(Effect.orDie) .pipe(Effect.orDie)
@ -361,7 +438,7 @@ describe("SessionProjector", () => {
yield* service.publish(SessionEvent.Step.Ended, { yield* service.publish(SessionEvent.Step.Ended, {
sessionID, sessionID,
timestamp: DateTime.makeUnsafe(1), timestamp: DateTime.makeUnsafe(1),
assistantMessageID: SessionMessage.ID.make("evt_assistant_2"), assistantMessageID: EventV2.ID.make("evt_assistant_2"),
finish: "stop", finish: "stop",
cost: 0, cost: 0,
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
@ -409,8 +486,8 @@ describe("SessionProjector", () => {
yield* db yield* db
.insert(SessionMessageTable) .insert(SessionMessageTable)
.values([ .values([
assistantRow(SessionMessage.ID.make("evt_assistant_stale"), 0), assistantRow(SessionMessage.ID.make("msg_assistant_stale"), 0),
assistantRow(SessionMessage.ID.make("evt_assistant_completed"), 1, { assistantRow(SessionMessage.ID.make("msg_assistant_completed"), 1, {
created: DateTime.makeUnsafe(1), created: DateTime.makeUnsafe(1),
completed: DateTime.makeUnsafe(2), completed: DateTime.makeUnsafe(2),
}), }),
@ -437,7 +514,7 @@ describe("SessionProjector", () => {
) )
expect(messages).toEqual([ expect(messages).toEqual([
new SessionMessage.Assistant({ new SessionMessage.Assistant({
id: SessionMessage.ID.make("evt_assistant_completed"), id: SessionMessage.ID.make("msg_assistant_completed"),
type: "assistant", type: "assistant",
agent: "build", agent: "build",
model, model,
@ -445,7 +522,7 @@ describe("SessionProjector", () => {
time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) }, time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
}), }),
new SessionMessage.Assistant({ new SessionMessage.Assistant({
id: SessionMessage.ID.make("evt_assistant_stale"), id: SessionMessage.ID.make("msg_assistant_stale"),
type: "assistant", type: "assistant",
agent: "build", agent: "build",
model, model,

View file

@ -310,7 +310,7 @@ describe("SessionV2.prompt", () => {
yield* events.publish( yield* events.publish(
SessionEvent.Prompted, SessionEvent.Prompted,
{ sessionID, timestamp: yield* DateTime.now, prompt, delivery: "steer" }, { 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 }) const retried = yield* session.prompt({ id: messageID, sessionID, prompt, resume: false })
@ -329,7 +329,7 @@ describe("SessionV2.prompt", () => {
yield* events.publish( yield* events.publish(
SessionEvent.Prompted, SessionEvent.Prompted,
{ sessionID, timestamp: yield* DateTime.now, prompt, delivery: "queue" }, { 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 }) 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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
@ -347,7 +347,7 @@ describe("SessionV2.prompt", () => {
yield* events.publish( yield* events.publish(
SessionEvent.Synthetic, SessionEvent.Synthetic,
{ sessionID, timestamp: yield* DateTime.now, text: "Collision" }, { sessionID, timestamp: yield* DateTime.now, text: "Collision" },
{ id: messageID }, { id: SessionMessage.ID.toEvent(messageID) },
) )
const failure = yield* session 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* () { Effect.gen(function* () {
yield* setup yield* setup
const { db } = yield* Database.Service const { db } = yield* Database.Service
@ -368,24 +368,51 @@ describe("SessionV2.prompt", () => {
const prompt = new Prompt({ text: "Reserved prompt" }) const prompt = new Prompt({ text: "Reserved prompt" })
yield* session.prompt({ id: messageID, sessionID, prompt, resume: false }) 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 const failure = yield* events
.publish( .publish(
SessionEvent.Synthetic, SessionEvent.Synthetic,
{ sessionID, timestamp: yield* DateTime.now, text: "Conflicting synthetic" }, { sessionID, timestamp: yield* DateTime.now, text: "Conflicting synthetic" },
{ id: messageID }, { id: SessionMessage.ID.toEvent(messageID) },
) )
.pipe(Effect.catchDefect(Effect.succeed)) .pipe(Effect.catchDefect(Effect.succeed))
expect(failure).toBe("Durable event conflicts with admitted prompt input") expect(failure).toBe("Durable event conflicts with admitted prompt input")
expect(yield* admitted(messageID)).not.toHaveProperty("promotedSeq") expect(yield* admitted(messageID)).not.toHaveProperty("promotedSeq")
expect(yield* session.messages({ sessionID })).toEqual([]) 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" },
])
}), }),
) )

View file

@ -1,7 +1,6 @@
import { describe, expect, test } from "bun:test" import { describe, expect, test } from "bun:test"
import { Message, Model } from "@opencode-ai/llm" import { Message, Model } from "@opencode-ai/llm"
import * as OpenAIChat from "@opencode-ai/llm/protocols/openai-chat" 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 { ModelV2 } from "@opencode-ai/core/model"
import { ProviderV2 } from "@opencode-ai/core/provider" import { ProviderV2 } from "@opencode-ai/core/provider"
import { SessionMessage } from "@opencode-ai/core/session/message" import { SessionMessage } from "@opencode-ai/core/session/message"
@ -12,7 +11,7 @@ import { ToolOutput } from "@opencode-ai/core/tool-output"
import { DateTime } from "effect" import { DateTime } from "effect"
const created = DateTime.makeUnsafe(0) 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 }) const model = Model.make({ id: "model", provider: "provider", route: OpenAIChat.route })
describe("toLLMMessages", () => { describe("toLLMMessages", () => {

View file

@ -1218,30 +1218,30 @@ describe("SessionRunnerLLM", () => {
const events = yield* EventV2.Service const events = yield* EventV2.Service
yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Recover interrupted tool" }), resume: false }) yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Recover interrupted tool" }), resume: false })
yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID) 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, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
agent: "build", agent: "build",
model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") }, model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
}) })).id
yield* events.publish(SessionEvent.Tool.Input.Started, { yield* events.publish(SessionEvent.Tool.Input.Started, {
sessionID, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
assistantMessageID: assistant.id, assistantMessageID,
callID: "call-interrupted", callID: "call-interrupted",
name: "echo", name: "echo",
}) })
yield* events.publish(SessionEvent.Tool.Input.Ended, { yield* events.publish(SessionEvent.Tool.Input.Ended, {
sessionID, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
assistantMessageID: assistant.id, assistantMessageID,
callID: "call-interrupted", callID: "call-interrupted",
text: '{"text":"stale"}', text: '{"text":"stale"}',
}) })
yield* events.publish(SessionEvent.Tool.Called, { yield* events.publish(SessionEvent.Tool.Called, {
sessionID, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
assistantMessageID: assistant.id, assistantMessageID,
callID: "call-interrupted", callID: "call-interrupted",
tool: "echo", tool: "echo",
input: { text: "stale" }, input: { text: "stale" },
@ -1280,30 +1280,30 @@ describe("SessionRunnerLLM", () => {
resume: false, resume: false,
}) })
yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID) 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, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
agent: "build", agent: "build",
model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") }, model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
}) })).id
yield* events.publish(SessionEvent.Tool.Input.Started, { yield* events.publish(SessionEvent.Tool.Input.Started, {
sessionID, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
assistantMessageID: assistant.id, assistantMessageID,
callID: "call-hosted-interrupted", callID: "call-hosted-interrupted",
name: "web_search", name: "web_search",
}) })
yield* events.publish(SessionEvent.Tool.Input.Ended, { yield* events.publish(SessionEvent.Tool.Input.Ended, {
sessionID, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
assistantMessageID: assistant.id, assistantMessageID,
callID: "call-hosted-interrupted", callID: "call-hosted-interrupted",
text: '{"query":"stale"}', text: '{"query":"stale"}',
}) })
yield* events.publish(SessionEvent.Tool.Called, { yield* events.publish(SessionEvent.Tool.Called, {
sessionID, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
assistantMessageID: assistant.id, assistantMessageID,
callID: "call-hosted-interrupted", callID: "call-hosted-interrupted",
tool: "web_search", tool: "web_search",
input: { query: "stale" }, input: { query: "stale" },
@ -1338,16 +1338,16 @@ describe("SessionRunnerLLM", () => {
resume: false, resume: false,
}) })
yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID) 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, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
agent: "build", agent: "build",
model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") }, model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
}) })).id
yield* events.publish(SessionEvent.Tool.Input.Started, { yield* events.publish(SessionEvent.Tool.Input.Started, {
sessionID, sessionID,
timestamp: yield* DateTime.now, timestamp: yield* DateTime.now,
assistantMessageID: assistant.id, assistantMessageID,
callID: "call-pending-interrupted", callID: "call-pending-interrupted",
name: "echo", name: "echo",
}) })

View file

@ -60,7 +60,7 @@ describe("Tool.Progress", () => {
const row = yield* db const row = yield* db
.select() .select()
.from(SessionMessageTable) .from(SessionMessageTable)
.where(eq(SessionMessageTable.id, assistantMessageID)) .where(eq(SessionMessageTable.id, SessionMessage.ID.fromEvent(assistantMessageID)))
.get() .get()
.pipe(Effect.orDie) .pipe(Effect.orDie)
if (!row) return yield* Effect.die("Missing projected assistant") if (!row) return yield* Effect.die("Missing projected assistant")

View file

@ -1,4 +1,5 @@
import { useEvent } from "@tui/context/event" import { useEvent } from "@tui/context/event"
import { EventV2 } from "@opencode-ai/core/event"
import type { import type {
SessionMessage, SessionMessage,
SessionMessageAssistant, SessionMessageAssistant,
@ -9,6 +10,9 @@ import type {
import { createStore, produce, reconcile } from "solid-js/store" import { createStore, produce, reconcile } from "solid-js/store"
import { createSimpleContext } from "./helper" import { createSimpleContext } from "./helper"
import { useSDK } from "./sdk" 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[]) { function activeAssistant(messages: SessionMessage[]) {
const index = messages.findIndex((message) => message.type === "assistant" && !message.time.completed) 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": case "session.next.agent.switched":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
draft.unshift({ draft.unshift({
id: event.id, id: messageID(event.id),
type: "agent-switched", type: "agent-switched",
agent: event.properties.agent, agent: event.properties.agent,
time: { created: event.properties.timestamp }, time: { created: event.properties.timestamp },
@ -92,7 +96,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.model.switched": case "session.next.model.switched":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
draft.unshift({ draft.unshift({
id: event.id, id: messageID(event.id),
type: "model-switched", type: "model-switched",
model: event.properties.model, model: event.properties.model,
time: { created: event.properties.timestamp }, time: { created: event.properties.timestamp },
@ -102,7 +106,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.prompted": { case "session.next.prompted": {
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
draft.unshift({ draft.unshift({
id: event.id, id: messageID(event.id),
type: "user", type: "user",
text: event.properties.prompt.text, text: event.properties.prompt.text,
files: event.properties.prompt.files, files: event.properties.prompt.files,
@ -116,7 +120,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.synthetic": case "session.next.synthetic":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
draft.unshift({ draft.unshift({
id: event.id, id: messageID(event.id),
type: "synthetic", type: "synthetic",
sessionID: event.properties.sessionID, sessionID: event.properties.sessionID,
text: event.properties.text, text: event.properties.text,
@ -127,7 +131,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.shell.started": case "session.next.shell.started":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
draft.unshift({ draft.unshift({
id: event.id, id: messageID(event.id),
type: "shell", type: "shell",
callID: event.properties.callID, callID: event.properties.callID,
command: event.properties.command, command: event.properties.command,
@ -149,7 +153,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
const currentAssistant = activeAssistant(draft) const currentAssistant = activeAssistant(draft)
if (currentAssistant) currentAssistant.time.completed = event.properties.timestamp if (currentAssistant) currentAssistant.time.completed = event.properties.timestamp
draft.unshift({ draft.unshift({
id: event.id, id: messageID(event.id),
type: "assistant", type: "assistant",
agent: event.properties.agent, agent: event.properties.agent,
model: event.properties.model, model: event.properties.model,
@ -161,7 +165,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
break break
case "session.next.step.ended": case "session.next.step.ended":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
const currentAssistant = ownedAssistant(draft, event.properties.assistantMessageID) const currentAssistant = ownedAssistant(draft, messageID(event.properties.assistantMessageID))
if (!currentAssistant) return if (!currentAssistant) return
currentAssistant.time.completed = event.properties.timestamp currentAssistant.time.completed = event.properties.timestamp
currentAssistant.finish = event.properties.finish currentAssistant.finish = event.properties.finish
@ -173,7 +177,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
break break
case "session.next.step.failed": case "session.next.step.failed":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
const currentAssistant = ownedAssistant(draft, event.properties.assistantMessageID) const currentAssistant = ownedAssistant(draft, messageID(event.properties.assistantMessageID))
if (!currentAssistant) return if (!currentAssistant) return
currentAssistant.time.completed = event.properties.timestamp currentAssistant.time.completed = event.properties.timestamp
currentAssistant.finish = "error" currentAssistant.finish = "error"
@ -199,7 +203,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
break break
case "session.next.tool.input.started": case "session.next.tool.input.started":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
ownedAssistant(draft, event.properties.assistantMessageID)?.content.push({ ownedAssistant(draft, messageID(event.properties.assistantMessageID))?.content.push({
type: "tool", type: "tool",
id: event.properties.callID, id: event.properties.callID,
name: event.properties.name, name: event.properties.name,
@ -211,7 +215,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.input.delta": case "session.next.tool.input.delta":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
const match = latestTool( const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID), ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID, event.properties.callID,
) )
if (match?.state.status === "pending") match.state.input += event.properties.delta 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": case "session.next.tool.input.ended":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
const match = latestTool( const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID), ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID, event.properties.callID,
) )
if (match?.state.status === "pending") match.state.input = event.properties.text 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": case "session.next.tool.called":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
const match = latestTool( const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID), ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID, event.properties.callID,
) )
if (!match) return if (!match) return
@ -241,7 +245,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.progress": case "session.next.tool.progress":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
const match = latestTool( const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID), ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID, event.properties.callID,
) )
if (match?.state.status !== "running") return if (match?.state.status !== "running") return
@ -252,7 +256,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.success": case "session.next.tool.success":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
const match = latestTool( const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID), ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID, event.properties.callID,
) )
if (match?.state.status !== "running") return if (match?.state.status !== "running") return
@ -270,7 +274,7 @@ export const { use: useSyncV2, provider: SyncProviderV2 } = createSimpleContext(
case "session.next.tool.failed": case "session.next.tool.failed":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
const match = latestTool( const match = latestTool(
ownedAssistant(draft, event.properties.assistantMessageID), ownedAssistant(draft, messageID(event.properties.assistantMessageID)),
event.properties.callID, event.properties.callID,
) )
if (!match || (match.state.status !== "pending" && match.state.status !== "running")) return 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": case "session.next.compaction.started":
update(event.properties.sessionID, (draft) => { update(event.properties.sessionID, (draft) => {
draft.unshift({ draft.unshift({
id: event.id, id: messageID(event.id),
type: "compaction", type: "compaction",
reason: event.properties.reason, reason: event.properties.reason,
summary: "", summary: "",

View file

@ -1,4 +1,3 @@
import { Config } from "@/config/config"
import { ConfigV1 } from "@opencode-ai/core/v1/config/config" import { ConfigV1 } from "@opencode-ai/core/v1/config/config"
import { EventV2 } from "@opencode-ai/core/event" import { EventV2 } from "@opencode-ai/core/event"
import { InstanceDisposed } from "@/server/event" import { InstanceDisposed } from "@/server/event"
@ -20,10 +19,10 @@ const SyncEventSchemas = EventV2.registry
return [ return [
Schema.Struct({ Schema.Struct({
type: Schema.Literal("sync"), type: Schema.Literal("sync"),
id: Schema.String, id: EventV2.ID,
syncEvent: Schema.Struct({ syncEvent: Schema.Struct({
type: Schema.Literal(EventV2.versionedType(definition.type, definition.sync.version)), type: Schema.Literal(EventV2.versionedType(definition.type, definition.sync.version)),
id: Schema.String, id: EventV2.ID,
seq: Schema.Finite, seq: Schema.Finite,
aggregateID: Schema.String, aggregateID: Schema.String,
data: definition.data, data: definition.data,
@ -41,7 +40,7 @@ const GlobalEventSchema = Schema.Struct({
...EventV2.registry ...EventV2.registry
.values() .values()
.map((definition) => .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(), .toArray(),
InstanceDisposed, InstanceDisposed,

View file

@ -1,4 +1,5 @@
import { NonNegativeInt } from "@opencode-ai/core/schema" import { NonNegativeInt } from "@opencode-ai/core/schema"
import { EventV2 } from "@opencode-ai/core/event"
import { SessionID } from "@/session/schema" import { SessionID } from "@/session/schema"
import { Schema } from "effect" import { Schema } from "effect"
import { HttpApi, HttpApiEndpoint, HttpApiError, HttpApiGroup, OpenApi } from "effect/unstable/httpapi" import { HttpApi, HttpApiEndpoint, HttpApiError, HttpApiGroup, OpenApi } from "effect/unstable/httpapi"
@ -9,7 +10,7 @@ import { described } from "./metadata"
const root = "/sync" const root = "/sync"
export const ReplayEvent = Schema.Struct({ export const ReplayEvent = Schema.Struct({
id: Schema.String, id: EventV2.ID,
aggregateID: Schema.String, aggregateID: Schema.String,
seq: NonNegativeInt, seq: NonNegativeInt,
type: Schema.String, type: Schema.String,
@ -27,7 +28,7 @@ export const SessionPayload = Schema.Struct({
}) })
export const HistoryPayload = Schema.Record(Schema.String, NonNegativeInt) export const HistoryPayload = Schema.Record(Schema.String, NonNegativeInt)
export const HistoryEvent = Schema.Struct({ export const HistoryEvent = Schema.Struct({
id: Schema.String, id: EventV2.ID,
aggregate_id: Schema.String, aggregate_id: Schema.String,
seq: NonNegativeInt, seq: NonNegativeInt,
type: Schema.String, type: Schema.String,

View file

@ -49,14 +49,14 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
await mounted await mounted
events.emit( events.emit(
global({ global({
id: "agent-1", id: "evt_agent_1",
type: "session.next.agent.switched", type: "session.next.agent.switched",
properties: { sessionID: "session-1", timestamp: 0, agent: "build" }, properties: { sessionID: "session-1", timestamp: 0, agent: "build" },
}), }),
) )
events.emit( events.emit(
global({ global({
id: "model-1", id: "evt_model_1",
type: "session.next.model.switched", type: "session.next.model.switched",
properties: { properties: {
sessionID: "session-1", sessionID: "session-1",
@ -67,7 +67,7 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
) )
events.emit( events.emit(
global({ global({
id: "assistant-1", id: "evt_assistant_1",
type: "session.next.step.started", type: "session.next.step.started",
properties: { properties: {
sessionID: "session-1", sessionID: "session-1",
@ -79,12 +79,12 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
) )
events.emit( events.emit(
global({ global({
id: "input-1", id: "evt_input_1",
type: "session.next.tool.input.started", type: "session.next.tool.input.started",
properties: { properties: {
sessionID: "session-1", sessionID: "session-1",
timestamp: 2, timestamp: 2,
assistantMessageID: "assistant-1", assistantMessageID: "evt_assistant_1",
callID: "call-1", callID: "call-1",
name: "bash", name: "bash",
}, },
@ -92,12 +92,12 @@ test("sync v2 settles pending tools when a live failure arrives", async () => {
) )
events.emit( events.emit(
global({ global({
id: "failed-1", id: "evt_failed_1",
type: "session.next.tool.failed", type: "session.next.tool.failed",
properties: { properties: {
sessionID: "session-1", sessionID: "session-1",
timestamp: 3, timestamp: 3,
assistantMessageID: "assistant-1", assistantMessageID: "evt_assistant_1",
callID: "call-1", callID: "call-1",
error: { type: "unknown", message: "aborted" }, error: { type: "unknown", message: "aborted" },
provider: { executed: false }, 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] const assistant = sync.session.message.fromSession("session-1")[0]
expect(assistant?.type).toBe("assistant") expect(assistant?.type).toBe("assistant")
if (assistant?.type !== "assistant") return if (assistant?.type !== "assistant") return
expect(assistant.id).toBe("msg_assistant_1")
const tool = assistant.content[0] const tool = assistant.content[0]
expect(tool?.type).toBe("tool") expect(tool?.type).toBe("tool")
if (tool?.type !== "tool") return if (tool?.type !== "tool") return

View file

@ -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 { SessionPaths } from "../../src/server/routes/instance/httpapi/groups/session"
import { Session } from "@/session/session" import { Session } from "@/session/session"
import { MessageID, PartID, SessionID, type SessionID as SessionIDType } from "../../src/session/schema" 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 { Database } from "@opencode-ai/core/database/database"
import { SessionInputTable, SessionMessageTable, SessionTable } from "@opencode-ai/core/session/sql" import { SessionInputTable, SessionMessageTable, SessionTable } from "@opencode-ai/core/session/sql"
import { SessionMessage } from "@opencode-ai/core/session/message" import { SessionMessage } from "@opencode-ai/core/session/message"
@ -584,7 +583,7 @@ describe("session HttpApi", () => {
request(`/api/session/${session.id}/prompt`, { request(`/api/session/${session.id}/prompt`, {
method: "POST", method: "POST",
headers: { ...headers, "content-type": "application/json" }, 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 first = yield* recordPrompt()
const retried = yield* recordPrompt() const retried = yield* recordPrompt()
@ -604,12 +603,12 @@ describe("session HttpApi", () => {
db db
.select() .select()
.from(SessionInputTable) .from(SessionInputTable)
.where(eq(SessionInputTable.id, SessionMessage.ID.make("evt_http_prompt"))) .where(eq(SessionInputTable.id, SessionMessage.ID.make("msg_http_prompt")))
.get() .get()
.pipe(Effect.orDie), .pipe(Effect.orDie),
) )
expect(admitted).toMatchObject({ expect(admitted).toMatchObject({
id: "evt_http_prompt", id: "msg_http_prompt",
session_id: session.id, session_id: session.id,
delivery: "steer", delivery: "steer",
promoted_seq: null, promoted_seq: null,
@ -618,13 +617,13 @@ describe("session HttpApi", () => {
const conflict = yield* request(`/api/session/${session.id}/prompt`, { const conflict = yield* request(`/api/session/${session.id}/prompt`, {
method: "POST", method: "POST",
headers: { ...headers, "content-type": "application/json" }, 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(conflict.status).toBe(409)
expect(yield* responseJson(conflict)).toEqual({ expect(yield* responseJson(conflict)).toEqual({
_tag: "ConflictError", _tag: "ConflictError",
message: "Prompt message ID conflicts with an existing durable record: evt_http_prompt", message: "Prompt message ID conflicts with an existing durable record: msg_http_prompt",
resource: "evt_http_prompt", resource: "msg_http_prompt",
}) })
}), }),
{ git: true, config: { formatter: false, lsp: false } }, { git: true, config: { formatter: false, lsp: false } },

View file

@ -7,6 +7,7 @@ import { ModelV2 } from "@opencode-ai/core/model"
import { ProviderV2 } from "@opencode-ai/core/provider" import { ProviderV2 } from "@opencode-ai/core/provider"
import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionEvent } from "@opencode-ai/core/session/event"
import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater" 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" import { ToolOutput } from "@opencode-ai/core/tool-output"
test.skip("step snapshots carry over to assistant messages", () => { 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", () => { test.skip("text ended populates assistant text content", () => {
const state: SessionMessageUpdater.MemoryState = { messages: [] } const state: SessionMessageUpdater.MemoryState = { messages: [] }
const sessionID = SessionID.make("session") const sessionID = SessionID.make("session")
const assistantMessageID = EventV2.ID.create()
Effect.runSync( Effect.runSync(
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(), id: assistantMessageID,
type: "session.next.step.started", type: "session.next.step.started",
data: { data: {
sessionID, sessionID,
@ -243,7 +245,7 @@ test.skip("compaction events reduce to compaction message", () => {
expect(state.messages).toHaveLength(1) expect(state.messages).toHaveLength(1)
expect(state.messages[0]).toMatchObject({ expect(state.messages[0]).toMatchObject({
id, id: SessionMessage.ID.fromEvent(id),
type: "compaction", type: "compaction",
reason: "auto", reason: "auto",
summary: "final summary", summary: "final summary",

File diff suppressed because it is too large Load diff

View file

@ -6,6 +6,27 @@ This document covers meaningful contract changes introduced on the `feat/opencod
## Earlier Branch History ## 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 ### Replayable Session Event Refinement And Cursor Stream
Affected schema: Affected schema: