Compare commits

..

1 Commits

Author SHA1 Message Date
Kit Langton 6cb2f00d15 refactor(core): simplify interrupt continuation 2026-08-15 14:57:49 -04:00
35 changed files with 3524 additions and 14144 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"@opencode-ai/core": patch
---
Simplify interrupt continuation: the steer-scoped resume decision now lives in SessionExecution as a post-cleanup inbox check, and the run coordinator drops its continuation state machine. Wakes arriving during cancellation cleanup now restart a normal full drain, and interrupting an idle session with continue now resumes pending steering input.
-14
View File
@@ -383,15 +383,6 @@ export type Endpoint5_31Output =
readonly location?: Location.Ref | undefined readonly location?: Location.Ref | undefined
readonly data: { readonly sessionID: Session.ID; readonly title: string } readonly data: { readonly sessionID: Session.ID; readonly title: string }
} }
| {
readonly id: Event.ID
readonly created: DateTime.Utc
readonly metadata?: { readonly [x: string]: unknown } | undefined
readonly type: "session.viewed"
readonly durable: { readonly aggregateID: string; readonly seq: Event.Seq; readonly version: Event.Version }
readonly location?: Location.Ref | undefined
readonly data: { readonly sessionID: Session.ID }
}
| { | {
readonly id: Event.ID readonly id: Event.ID
readonly created: DateTime.Utc readonly created: DateTime.Utc
@@ -924,10 +915,6 @@ export type Endpoint5_34Input = { readonly sessionID: Session.ID; readonly messa
export type Endpoint5_34Output = SessionMessage.Info export type Endpoint5_34Output = SessionMessage.Info
export type SessionMessageOperation<E = never> = (input: Endpoint5_34Input) => Effect.Effect<Endpoint5_34Output, E> export type SessionMessageOperation<E = never> = (input: Endpoint5_34Input) => Effect.Effect<Endpoint5_34Output, E>
export type Endpoint5_35Input = { readonly sessionID: Session.ID }
export type Endpoint5_35Output = void
export type SessionViewOperation<E = never> = (input: Endpoint5_35Input) => Effect.Effect<Endpoint5_35Output, E>
export interface SessionApi<E = never> { export interface SessionApi<E = never> {
readonly list: SessionListOperation<E> readonly list: SessionListOperation<E>
readonly create: SessionCreateOperation<E> readonly create: SessionCreateOperation<E>
@@ -972,7 +959,6 @@ export interface SessionApi<E = never> {
readonly interrupt: SessionInterruptOperation<E> readonly interrupt: SessionInterruptOperation<E>
readonly background: SessionBackgroundOperation<E> readonly background: SessionBackgroundOperation<E>
readonly message: SessionMessageOperation<E> readonly message: SessionMessageOperation<E>
readonly view: SessionViewOperation<E>
} }
export type Endpoint6_0Input = { export type Endpoint6_0Input = {
@@ -86,8 +86,6 @@ import type {
Endpoint5_33Output, Endpoint5_33Output,
Endpoint5_34Input, Endpoint5_34Input,
Endpoint5_34Output, Endpoint5_34Output,
Endpoint5_35Input,
Endpoint5_35Output,
Endpoint6_0Input, Endpoint6_0Input,
Endpoint6_0Output, Endpoint6_0Output,
Endpoint7_0Input, Endpoint7_0Input,
@@ -612,11 +610,6 @@ const Endpoint5_34 = (raw: RawClient["server.session"]) => (input: Endpoint5_34I
), ),
) )
const Endpoint5_35 = (raw: RawClient["server.session"]) => (input: Endpoint5_35Input) =>
preserveEffect<Endpoint5_35Output>()(
raw["session.view"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError)),
)
const adaptGroup5 = (raw: RawClient["server.session"]) => ({ const adaptGroup5 = (raw: RawClient["server.session"]) => ({
list: Endpoint5_0(raw), list: Endpoint5_0(raw),
create: Endpoint5_1(raw), create: Endpoint5_1(raw),
@@ -646,7 +639,6 @@ const adaptGroup5 = (raw: RawClient["server.session"]) => ({
interrupt: Endpoint5_32(raw), interrupt: Endpoint5_32(raw),
background: Endpoint5_33(raw), background: Endpoint5_33(raw),
message: Endpoint5_34(raw), message: Endpoint5_34(raw),
view: Endpoint5_35(raw),
}) })
const Endpoint6_0 = (raw: RawClient["server.message"]) => (input: Endpoint6_0Input) => const Endpoint6_0 = (raw: RawClient["server.message"]) => (input: Endpoint6_0Input) =>
@@ -80,8 +80,6 @@ import type {
SessionBackgroundOutput, SessionBackgroundOutput,
SessionMessageInput, SessionMessageInput,
SessionMessageOutput, SessionMessageOutput,
SessionViewInput,
SessionViewOutput,
MessageListInput, MessageListInput,
MessageListOutput, MessageListOutput,
ModelListInput, ModelListInput,
@@ -898,17 +896,6 @@ export function make(options: ClientOptions) {
}, },
requestOptions, requestOptions,
).then((value) => value.data), ).then((value) => value.data),
view: (input: SessionViewInput, requestOptions?: RequestOptions) =>
request<SessionViewOutput>(
{
method: "POST",
path: `/api/session/${encodeURIComponent(input.sessionID)}/view`,
successStatus: 204,
declaredStatuses: [404, 401, 400],
empty: true,
},
requestOptions,
),
}, },
message: { message: {
list: (input: MessageListInput, requestOptions?: RequestOptions) => list: (input: MessageListInput, requestOptions?: RequestOptions) =>
+4 -38
View File
@@ -479,16 +479,6 @@ export type SessionRenamed = {
data: { sessionID: string; title: string } data: { sessionID: string; title: string }
} }
export type SessionViewed = {
id: string
created: number
metadata?: { [x: string]: any }
type: "session.viewed"
durable: { aggregateID: string; seq: number; version: 1 }
location?: LocationRef
data: { sessionID: string }
}
export type SessionDeleted = { export type SessionDeleted = {
id: string id: string
created: number created: number
@@ -1520,7 +1510,7 @@ export type SessionInfo = {
model?: ModelRef model?: ModelRef
cost: MoneyUSD cost: MoneyUSD
tokens: TokenUsageInfo tokens: TokenUsageInfo
time: { created: number; updated: number; idle?: number; viewed?: number; archived?: number } time: { created: number; updated: number; archived?: number }
title?: string title?: string
location: LocationRef location: LocationRef
subpath?: string subpath?: string
@@ -1933,7 +1923,6 @@ export type SessionEventDurable =
| SessionModelSelected | SessionModelSelected
| SessionMoved | SessionMoved
| SessionRenamed | SessionRenamed
| SessionViewed
| SessionDeleted | SessionDeleted
| SessionForked | SessionForked
| SessionInboxDelivered | SessionInboxDelivered
@@ -2024,7 +2013,6 @@ export type V2Event =
| SessionModelSelected | SessionModelSelected
| SessionMoved | SessionMoved
| SessionRenamed | SessionRenamed
| SessionViewed
| SessionUsageUpdated | SessionUsageUpdated
| SessionDeleted | SessionDeleted
| SessionForked | SessionForked
@@ -2488,13 +2476,7 @@ export type SessionImportInput = {
readonly reasoning: number readonly reasoning: number
readonly cache: { readonly read: number; readonly write: number } readonly cache: { readonly read: number; readonly write: number }
} }
readonly time: { readonly time: { readonly created: number; readonly updated: number; readonly archived?: number }
readonly created: number
readonly updated: number
readonly idle?: number
readonly viewed?: number
readonly archived?: number
}
readonly title?: string readonly title?: string
readonly location: { readonly directory: string; readonly workspaceID?: string } readonly location: { readonly directory: string; readonly workspaceID?: string }
readonly subpath?: string readonly subpath?: string
@@ -2761,13 +2743,7 @@ export type SessionImportInput = {
readonly reasoning: number readonly reasoning: number
readonly cache: { readonly read: number; readonly write: number } readonly cache: { readonly read: number; readonly write: number }
} }
readonly time: { readonly time: { readonly created: number; readonly updated: number; readonly archived?: number }
readonly created: number
readonly updated: number
readonly idle?: number
readonly viewed?: number
readonly archived?: number
}
readonly title?: string readonly title?: string
readonly location: { readonly directory: string; readonly workspaceID?: string } readonly location: { readonly directory: string; readonly workspaceID?: string }
readonly subpath?: string readonly subpath?: string
@@ -3034,13 +3010,7 @@ export type SessionImportInput = {
readonly reasoning: number readonly reasoning: number
readonly cache: { readonly read: number; readonly write: number } readonly cache: { readonly read: number; readonly write: number }
} }
readonly time: { readonly time: { readonly created: number; readonly updated: number; readonly archived?: number }
readonly created: number
readonly updated: number
readonly idle?: number
readonly viewed?: number
readonly archived?: number
}
readonly title?: string readonly title?: string
readonly location: { readonly directory: string; readonly workspaceID?: string } readonly location: { readonly directory: string; readonly workspaceID?: string }
readonly subpath?: string readonly subpath?: string
@@ -3966,10 +3936,6 @@ export type SessionMessageInput = {
export type SessionMessageOutput = { data: SessionMessageInfo }["data"] export type SessionMessageOutput = { data: SessionMessageInfo }["data"]
export type SessionViewInput = { readonly sessionID: { readonly sessionID: string }["sessionID"] }
export type SessionViewOutput = void
export type MessageListInput = { export type MessageListInput = {
readonly sessionID: { readonly sessionID: string }["sessionID"] readonly sessionID: { readonly sessionID: string }["sessionID"]
readonly limit?: { readonly limit?: {
+1 -11
View File
@@ -136,10 +136,8 @@ test("event.subscribe terminates on Effect protocol decode failures", async () =
test("session methods retain decoded Effect inputs and outputs", async () => { test("session methods retain decoded Effect inputs and outputs", async () => {
const logQueries: Array<Record<string, string>> = [] const logQueries: Array<Record<string, string>> = []
const requests: Array<{ method: string; url: string }> = []
const httpClient = HttpClient.make((request) => { const httpClient = HttpClient.make((request) => {
const url = request.url const url = request.url
requests.push({ method: request.method, url })
if (url.includes("/log")) { if (url.includes("/log")) {
logQueries.push(Object.fromEntries(request.urlParams.params)) logQueries.push(Object.fromEntries(request.urlParams.params))
return Effect.succeed( return Effect.succeed(
@@ -185,7 +183,6 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
const created = yield* client.session.create({ const created = yield* client.session.create({
location: Location.Ref.make({ directory: AbsolutePath.make("/tmp/project") }), location: Location.Ref.make({ directory: AbsolutePath.make("/tmp/project") }),
}) })
yield* client.session.view({ sessionID: Session.ID.make("ses_test") })
yield* client.session.switchAgent({ sessionID: Session.ID.make("ses_test"), agent: Agent.ID.make("build") }) yield* client.session.switchAgent({ sessionID: Session.ID.make("ses_test"), agent: Agent.ID.make("build") })
yield* client.session.switchModel({ yield* client.session.switchModel({
sessionID: Session.ID.make("ses_test"), sessionID: Session.ID.make("ses_test"),
@@ -210,11 +207,7 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
return { page, active, created, admitted, context, log, message } return { page, active, created, admitted, context, log, message }
}).pipe(Effect.provideService(HttpClient.HttpClient, httpClient), Effect.runPromise) }).pipe(Effect.provideService(HttpClient.HttpClient, httpClient), Effect.runPromise)
const listed = result.page.data[0] expect(DateTime.toEpochMillis(result.page.data[0].time.created)).toBe(1_717_171_717_000)
if (!listed?.time.idle || !listed.time.viewed) throw new Error("Expected attention times")
expect(DateTime.toEpochMillis(listed.time.created)).toBe(1_717_171_717_000)
expect(DateTime.toEpochMillis(listed.time.idle)).toBe(1_717_171_717_002)
expect(DateTime.toEpochMillis(listed.time.viewed)).toBe(1_717_171_717_001)
expect(result.active).toEqual({ ses_test: { type: "running" } }) expect(result.active).toEqual({ ses_test: { type: "running" } })
expect(Object.getPrototypeOf(result.page.data[0])).toBe(Object.prototype) expect(Object.getPrototypeOf(result.page.data[0])).toBe(Object.prototype)
expect(Object.getPrototypeOf(result.created)).toBe(Object.prototype) expect(Object.getPrototypeOf(result.created)).toBe(Object.prototype)
@@ -224,7 +217,6 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
expect(DateTime.toEpochMillis(result.admitted.timeCreated)).toBe(1_717_171_717_000) expect(DateTime.toEpochMillis(result.admitted.timeCreated)).toBe(1_717_171_717_000)
expect(result.context).toEqual([]) expect(result.context).toEqual([])
expect(logQueries[0]).toEqual({ after: "0" }) expect(logQueries[0]).toEqual({ after: "0" })
expect(requests).toContainEqual({ method: "POST", url: "http://localhost:3000/api/session/ses_test/view" })
const logged = Array.from(result.log) const logged = Array.from(result.log)
expect(logged.map((item) => item.type)).toEqual(["session.model.selected", "log.synced"]) expect(logged.map((item) => item.type)).toEqual(["session.model.selected", "log.synced"])
expect(logged[0]?.type === "session.model.selected" && DateTime.toEpochMillis(logged[0].created)).toBe( expect(logged[0]?.type === "session.model.selected" && DateTime.toEpochMillis(logged[0].created)).toBe(
@@ -268,8 +260,6 @@ const session = {
time: { time: {
created: 1_717_171_717_000, created: 1_717_171_717_000,
updated: 1_717_171_717_000, updated: 1_717_171_717_000,
idle: 1_717_171_717_002,
viewed: 1_717_171_717_001,
}, },
title: "Test", title: "Test",
location: { directory: "/tmp/project" }, location: { directory: "/tmp/project" },
-5
View File
@@ -539,7 +539,6 @@ test("session methods use the public HTTP contract", async () => {
const page = await client.session.list({ limit: 10, order: "desc", parentID: null }) const page = await client.session.list({ limit: 10, order: "desc", parentID: null })
const active = await client.session.active() const active = await client.session.active()
const created = await client.session.create({ location: { directory: "/tmp/project" } }) const created = await client.session.create({ location: { directory: "/tmp/project" } })
await client.session.view({ sessionID: "ses_test" })
await client.session.switchAgent({ sessionID: "ses_test", agent: "build" }) await client.session.switchAgent({ sessionID: "ses_test", agent: "build" })
await client.session.switchModel({ await client.session.switchModel({
sessionID: "ses_test", sessionID: "ses_test",
@@ -566,7 +565,6 @@ test("session methods use the public HTTP contract", async () => {
const message = await client.session.message({ sessionID: "ses_test", messageID: "msg_model" }) const message = await client.session.message({ sessionID: "ses_test", messageID: "msg_model" })
expect(page.cursor.next).toBe("next") expect(page.cursor.next).toBe("next")
expect(page.data[0].time).toMatchObject({ idle: 1_717_171_717_002, viewed: 1_717_171_717_001 })
expect(active).toEqual({ ses_test: { type: "running" } }) expect(active).toEqual({ ses_test: { type: "running" } })
expect(created.id).toBe("ses_test") expect(created.id).toBe("ses_test")
expect(admitted.id).toBe("msg_test") expect(admitted.id).toBe("msg_test")
@@ -579,7 +577,6 @@ test("session methods use the public HTTP contract", async () => {
["GET", "http://localhost:3000/api/session?limit=10&order=desc&parentID=null"], ["GET", "http://localhost:3000/api/session?limit=10&order=desc&parentID=null"],
["GET", "http://localhost:3000/api/session/active"], ["GET", "http://localhost:3000/api/session/active"],
["POST", "http://localhost:3000/api/session"], ["POST", "http://localhost:3000/api/session"],
["POST", "http://localhost:3000/api/session/ses_test/view"],
["POST", "http://localhost:3000/api/session/ses_test/agent"], ["POST", "http://localhost:3000/api/session/ses_test/agent"],
["POST", "http://localhost:3000/api/session/ses_test/model"], ["POST", "http://localhost:3000/api/session/ses_test/model"],
["POST", "http://localhost:3000/api/session/ses_test/prompt"], ["POST", "http://localhost:3000/api/session/ses_test/prompt"],
@@ -654,8 +651,6 @@ const session = {
time: { time: {
created: 1_717_171_717_000, created: 1_717_171_717_000,
updated: 1_717_171_717_000, updated: 1_717_171_717_000,
idle: 1_717_171_717_002,
viewed: 1_717_171_717_001,
}, },
title: "Test", title: "Test",
location: { directory: "/tmp/project" }, location: { directory: "/tmp/project" },
+44 -152
View File
@@ -1,10 +1,8 @@
{ {
"version": "7", "version": "7",
"dialect": "sqlite", "dialect": "sqlite",
"id": "94b6c496-ad84-426f-9d5d-3e1ac3ebfb56", "id": "dcde8e6b-4bf4-4f6b-b2be-4030c2c3e936",
"prevIds": [ "prevIds": ["5c1aa56b-c3ee-4283-9a84-c0bf626dc604"],
"dcde8e6b-4bf4-4f6b-b2be-4030c2c3e936"
],
"ddl": [ "ddl": [
{ {
"name": "account_state", "name": "account_state",
@@ -1352,26 +1350,6 @@
"entityType": "columns", "entityType": "columns",
"table": "session_v2" "table": "session_v2"
}, },
{
"type": "integer",
"notNull": false,
"autoincrement": false,
"default": null,
"generated": null,
"name": "time_idle",
"entityType": "columns",
"table": "session_v2"
},
{
"type": "integer",
"notNull": false,
"autoincrement": false,
"default": null,
"generated": null,
"name": "time_viewed",
"entityType": "columns",
"table": "session_v2"
},
{ {
"type": "integer", "type": "integer",
"notNull": false, "notNull": false,
@@ -1503,13 +1481,9 @@
"table": "worktree" "table": "worktree"
}, },
{ {
"columns": [ "columns": ["active_account_id"],
"active_account_id"
],
"tableTo": "account", "tableTo": "account",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "SET NULL", "onDelete": "SET NULL",
"nameExplicit": false, "nameExplicit": false,
@@ -1518,13 +1492,9 @@
"table": "account_state" "table": "account_state"
}, },
{ {
"columns": [ "columns": ["aggregate_id"],
"aggregate_id"
],
"tableTo": "event_sequence", "tableTo": "event_sequence",
"columnsTo": [ "columnsTo": ["aggregate_id"],
"aggregate_id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1533,13 +1503,9 @@
"table": "event" "table": "event"
}, },
{ {
"columns": [ "columns": ["project_id"],
"project_id"
],
"tableTo": "project", "tableTo": "project",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1548,13 +1514,9 @@
"table": "permission" "table": "permission"
}, },
{ {
"columns": [ "columns": ["project_id"],
"project_id"
],
"tableTo": "project", "tableTo": "project",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1563,13 +1525,9 @@
"table": "project_directory" "table": "project_directory"
}, },
{ {
"columns": [ "columns": ["session_id"],
"session_id"
],
"tableTo": "session_v2", "tableTo": "session_v2",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1578,13 +1536,9 @@
"table": "instruction_entry" "table": "instruction_entry"
}, },
{ {
"columns": [ "columns": ["session_id"],
"session_id"
],
"tableTo": "session_v2", "tableTo": "session_v2",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1593,13 +1547,9 @@
"table": "instruction_state" "table": "instruction_state"
}, },
{ {
"columns": [ "columns": ["session_id"],
"session_id"
],
"tableTo": "session_v2", "tableTo": "session_v2",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1608,13 +1558,9 @@
"table": "session_inbox" "table": "session_inbox"
}, },
{ {
"columns": [ "columns": ["session_id"],
"session_id"
],
"tableTo": "session_v2", "tableTo": "session_v2",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1623,13 +1569,9 @@
"table": "session_message" "table": "session_message"
}, },
{ {
"columns": [ "columns": ["session_id"],
"session_id"
],
"tableTo": "session_v2", "tableTo": "session_v2",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1638,13 +1580,9 @@
"table": "session_pending" "table": "session_pending"
}, },
{ {
"columns": [ "columns": ["project_id"],
"project_id"
],
"tableTo": "project", "tableTo": "project",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1653,13 +1591,9 @@
"table": "session_v2" "table": "session_v2"
}, },
{ {
"columns": [ "columns": ["project_id"],
"project_id"
],
"tableTo": "project", "tableTo": "project",
"columnsTo": [ "columnsTo": ["id"],
"id"
],
"onUpdate": "NO ACTION", "onUpdate": "NO ACTION",
"onDelete": "CASCADE", "onDelete": "CASCADE",
"nameExplicit": false, "nameExplicit": false,
@@ -1668,175 +1602,133 @@
"table": "worktree" "table": "worktree"
}, },
{ {
"columns": [ "columns": ["email", "url"],
"email",
"url"
],
"nameExplicit": false, "nameExplicit": false,
"name": "control_account_pk", "name": "control_account_pk",
"entityType": "pks", "entityType": "pks",
"table": "control_account" "table": "control_account"
}, },
{ {
"columns": [ "columns": ["project_id", "directory"],
"project_id",
"directory"
],
"nameExplicit": false, "nameExplicit": false,
"name": "project_directory_pk", "name": "project_directory_pk",
"entityType": "pks", "entityType": "pks",
"table": "project_directory" "table": "project_directory"
}, },
{ {
"columns": [ "columns": ["session_id", "key"],
"session_id",
"key"
],
"nameExplicit": false, "nameExplicit": false,
"name": "instruction_entry_pk", "name": "instruction_entry_pk",
"entityType": "pks", "entityType": "pks",
"table": "instruction_entry" "table": "instruction_entry"
}, },
{ {
"columns": [ "columns": ["project_id", "directory"],
"project_id",
"directory"
],
"nameExplicit": false, "nameExplicit": false,
"name": "worktree_pk", "name": "worktree_pk",
"entityType": "pks", "entityType": "pks",
"table": "worktree" "table": "worktree"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "account_state_pk", "name": "account_state_pk",
"table": "account_state", "table": "account_state",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "account_pk", "name": "account_pk",
"table": "account", "table": "account",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "credential_pk", "name": "credential_pk",
"table": "credential", "table": "credential",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["aggregate_id"],
"aggregate_id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "event_sequence_pk", "name": "event_sequence_pk",
"table": "event_sequence", "table": "event_sequence",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "event_pk", "name": "event_pk",
"table": "event", "table": "event",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["key"],
"key"
],
"nameExplicit": false, "nameExplicit": false,
"name": "kv_pk", "name": "kv_pk",
"table": "kv", "table": "kv",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "permission_pk", "name": "permission_pk",
"table": "permission", "table": "permission",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "project_pk", "name": "project_pk",
"table": "project", "table": "project",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["hash"],
"hash"
],
"nameExplicit": false, "nameExplicit": false,
"name": "instruction_blob_pk", "name": "instruction_blob_pk",
"table": "instruction_blob", "table": "instruction_blob",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["session_id"],
"session_id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "instruction_state_pk", "name": "instruction_state_pk",
"table": "instruction_state", "table": "instruction_state",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "session_inbox_pk", "name": "session_inbox_pk",
"table": "session_inbox", "table": "session_inbox",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "session_message_pk", "name": "session_message_pk",
"table": "session_message", "table": "session_message",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "session_pending_pk", "name": "session_pending_pk",
"table": "session_pending", "table": "session_pending",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "session_v2_pk", "name": "session_v2_pk",
"table": "session_v2", "table": "session_v2",
"entityType": "pks" "entityType": "pks"
}, },
{ {
"columns": [ "columns": ["id"],
"id"
],
"nameExplicit": false, "nameExplicit": false,
"name": "workspace_pk", "name": "workspace_pk",
"table": "workspace", "table": "workspace",
@@ -2132,4 +2024,4 @@
} }
], ],
"renames": [] "renames": []
} }
-2
View File
@@ -43,7 +43,6 @@ import m40 from "./migration/20260808023530_workspace_domain.js"
import m41 from "./migration/20260811161259_execution_claim_attempts.js" import m41 from "./migration/20260811161259_execution_claim_attempts.js"
import m42 from "./migration/20260812181746_session_inbox.js" import m42 from "./migration/20260812181746_session_inbox.js"
import m43 from "./migration/20260812213948_worktree.js" import m43 from "./migration/20260812213948_worktree.js"
import m44 from "./migration/20260815182818_session_viewed_state.js"
export const migrations = [ export const migrations = [
m00, m00,
@@ -90,5 +89,4 @@ export const migrations = [
m41, m41,
m42, m42,
m43, m43,
m44,
] satisfies DatabaseMigration.Migration[] ] satisfies DatabaseMigration.Migration[]
@@ -1,14 +0,0 @@
import { Effect } from "effect"
import type { DatabaseMigration } from "../migration.js"
const migration: DatabaseMigration.Migration = {
id: "20260815182818_session_viewed_state",
up(tx) {
return Effect.gen(function* () {
yield* tx.run(`ALTER TABLE \`session_v2\` ADD \`time_idle\` integer;`)
yield* tx.run(`ALTER TABLE \`session_v2\` ADD \`time_viewed\` integer;`)
})
},
}
export default migration
-2
View File
@@ -209,8 +209,6 @@ const schema: Omit<DatabaseMigration.Migration, "id"> = {
\`model\` text, \`model\` text,
\`time_created\` integer NOT NULL, \`time_created\` integer NOT NULL,
\`time_updated\` integer NOT NULL, \`time_updated\` integer NOT NULL,
\`time_idle\` integer,
\`time_viewed\` integer,
\`time_compacting\` integer, \`time_compacting\` integer,
\`time_archived\` integer, \`time_archived\` integer,
\`time_suspended\` integer, \`time_suspended\` integer,
-17
View File
@@ -165,7 +165,6 @@ export interface Interface {
input: ForkInput, input: ForkInput,
) => Effect.Effect<SessionSchema.Info, NotFoundError | MessageNotFoundError | ForkEmptyError> ) => Effect.Effect<SessionSchema.Info, NotFoundError | MessageNotFoundError | ForkEmptyError>
readonly get: (sessionID: SessionSchema.ID) => Effect.Effect<SessionSchema.Info, NotFoundError> readonly get: (sessionID: SessionSchema.ID) => Effect.Effect<SessionSchema.Info, NotFoundError>
readonly view: (input: { sessionID: SessionSchema.ID }) => Effect.Effect<void, NotFoundError>
readonly remove: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError> readonly remove: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError>
readonly messages: (input: { readonly messages: (input: {
sessionID: SessionSchema.ID sessionID: SessionSchema.ID
@@ -308,7 +307,6 @@ const layer = Layer.effect(
const scope = yield* Scope.Scope const scope = yield* Scope.Scope
const activeShells = new Set<SessionSchema.ID>() const activeShells = new Set<SessionSchema.ID>()
const shellLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>() const shellLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>()
const viewLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>()
const closeTransport = Effect.fn("Session.closeTransport")(function* (session: SessionSchema.Info) { const closeTransport = Effect.fn("Session.closeTransport")(function* (session: SessionSchema.Info) {
const location = Location.Ref.make({ const location = Location.Ref.make({
directory: session.location.directory, directory: session.location.directory,
@@ -449,21 +447,6 @@ const layer = Layer.effect(
if (!session) return yield* new NotFoundError({ sessionID }) if (!session) return yield* new NotFoundError({ sessionID })
return session return session
}), }),
view: Effect.fn("Session.view")((input) =>
viewLocks.withLock(input.sessionID)(
Effect.gen(function* () {
const row = yield* db
.select({ idle: SessionTable.time_idle, viewed: SessionTable.time_viewed })
.from(SessionTable)
.where(eq(SessionTable.id, input.sessionID))
.get()
.pipe(Effect.orDie)
if (!row) return yield* new NotFoundError({ sessionID: input.sessionID })
if (row.idle === null || row.viewed === row.idle) return
yield* bus.publish(SessionEvent.Viewed, { sessionID: input.sessionID })
}),
),
),
remove: Effect.fn("Session.remove")(function* (sessionID) { remove: Effect.fn("Session.remove")(function* (sessionID) {
const session = yield* result.get(sessionID) const session = yield* result.get(sessionID)
yield* execution.interrupt(sessionID) yield* execution.interrupt(sessionID)
+7 -7
View File
@@ -137,13 +137,13 @@ export const layer = Layer.effect(
return Service.of({ return Service.of({
active: coordinator.active, active: coordinator.active,
interrupt: (sessionID, options) => interrupt: (sessionID, options) =>
coordinator.interrupt( Effect.gen(function* () {
sessionID, yield* coordinator.interrupt(sessionID, "user")
"user", if (!options?.continue) return
options?.continue // Resume only steering input from the interrupted intent. Queued next-turn work
? { continue: { request: "steer", when: SessionInbox.has(db, sessionID, "steer") } } // stays parked: a steer-scoped drain never promotes queue-delivery rows.
: undefined, if (yield* SessionInbox.has(db, sessionID, "steer")) yield* coordinator.wake(sessionID, "steer")
), }),
resume: coordinator.run, resume: coordinator.run,
wake: coordinator.wake, wake: coordinator.wake,
wakeActive: coordinator.wakeActive, wakeActive: coordinator.wakeActive,
-2
View File
@@ -53,8 +53,6 @@ export function fromRow(row: typeof SessionTable.$inferSelect): SessionSchema.In
time: { time: {
created: DateTime.makeUnsafe(row.time_created), created: DateTime.makeUnsafe(row.time_created),
updated: DateTime.makeUnsafe(row.time_updated), updated: DateTime.makeUnsafe(row.time_updated),
idle: row.time_idle === null ? undefined : DateTime.makeUnsafe(row.time_idle),
viewed: row.time_viewed === null ? undefined : DateTime.makeUnsafe(row.time_viewed),
archived: row.time_archived ? DateTime.makeUnsafe(row.time_archived) : undefined, archived: row.time_archived ? DateTime.makeUnsafe(row.time_archived) : undefined,
}, },
}) })
@@ -59,7 +59,6 @@ export function update(adapter: Adapter, event: SessionEvent.DurableEvent) {
Match.type<SessionEvent.DurableEvent>(), Match.type<SessionEvent.DurableEvent>(),
Match.discriminatorsExhaustive("type")({ Match.discriminatorsExhaustive("type")({
"session.created": () => Effect.void, "session.created": () => Effect.void,
"session.viewed": () => Effect.void,
"session.usage.recorded": () => Effect.void, "session.usage.recorded": () => Effect.void,
"session.agent.selected": (event) => { "session.agent.selected": (event) => {
return Effect.gen(function* () { return Effect.gen(function* () {
+3 -38
View File
@@ -391,30 +391,6 @@ function insertMessage(db: DatabaseService, event: SessionEvent.DurableEvent, me
.pipe(Effect.orDie) .pipe(Effect.orDie)
} }
function projectIdle(
db: DatabaseService,
event:
| typeof SessionEvent.Execution.Succeeded.Type
| typeof SessionEvent.Execution.Failed.Type
| typeof SessionEvent.Execution.Interrupted.Type,
) {
return Effect.gen(function* () {
yield* run(db, event)
if (event.type === SessionEvent.Execution.Interrupted.type && event.data.reason === "shutdown") return
const time = DateTime.toEpochMillis(event.created)
yield* db
.update(SessionTable)
.set({
// Unread uses a strict timestamp comparison, so every terminal must advance even within one millisecond.
time_idle: sql`max(${time}, coalesce(${SessionTable.time_idle} + 1, ${time}))`,
time_updated: sql`${SessionTable.time_updated}`,
})
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie)
})
}
const layer = Layer.effectDiscard( const layer = Layer.effectDiscard(
Effect.gen(function* () { Effect.gen(function* () {
const bus = yield* Bus.Service const bus = yield* Bus.Service
@@ -536,17 +512,6 @@ const layer = Layer.effectDiscard(
.run() .run()
.pipe(Effect.orDie), .pipe(Effect.orDie),
) )
yield* bus.project(SessionEvent.Viewed, (event) =>
db
.update(SessionTable)
.set({
time_viewed: sql`${SessionTable.time_idle}`,
time_updated: sql`${SessionTable.time_updated}`,
})
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie),
)
yield* bus.project(SessionEvent.UsageRecorded, (event) => applyUsage(db, event.data.sessionID, event.data)) yield* bus.project(SessionEvent.UsageRecorded, (event) => applyUsage(db, event.data.sessionID, event.data))
yield* bus.project(SessionEvent.Forked, (event) => projectFork(db, event)) yield* bus.project(SessionEvent.Forked, (event) => projectFork(db, event))
yield* bus.project(SessionEvent.InboxDelivered, (event) => yield* bus.project(SessionEvent.InboxDelivered, (event) =>
@@ -615,9 +580,9 @@ const layer = Layer.effectDiscard(
delivery: event.data.delivery, delivery: event.data.delivery,
}), }),
) )
yield* bus.project(SessionEvent.Execution.Succeeded, (event) => projectIdle(db, event)) yield* bus.project(SessionEvent.Execution.Succeeded, (event) => run(db, event))
yield* bus.project(SessionEvent.Execution.Failed, (event) => projectIdle(db, event)) yield* bus.project(SessionEvent.Execution.Failed, (event) => run(db, event))
yield* bus.project(SessionEvent.Execution.Interrupted, (event) => projectIdle(db, event)) yield* bus.project(SessionEvent.Execution.Interrupted, (event) => run(db, event))
yield* bus.project(SessionEvent.InstructionsUpdated, (event) => yield* bus.project(SessionEvent.InstructionsUpdated, (event) =>
Effect.gen(function* () { Effect.gen(function* () {
yield* run(db, event) yield* run(db, event)
+23 -67
View File
@@ -10,25 +10,19 @@ export interface Coordinator<Key, E, Reason = never> {
/** Starts an execution while idle, or joins the active execution and returns its exit. */ /** Starts an execution while idle, or joins the active execution and returns its exit. */
readonly run: (key: Key) => Effect.Effect<void, E> readonly run: (key: Key) => Effect.Effect<void, E>
/** Rings the doorbell: an idle key starts an execution; an active one drains again before settling. */ /** Rings the doorbell: an idle key starts an execution; an active one drains again before settling. */
readonly wake: (key: Key, request?: Request) => Effect.Effect<void> readonly wake: (key: Key, scope?: Promotable) => Effect.Effect<void>
/** Rings the current execution's doorbell with its existing request. Idle keys remain idle. */ /** Rings the current execution's doorbell with its existing scope. Idle keys remain idle. */
readonly wakeActive: (key: Key) => Effect.Effect<void> readonly wakeActive: (key: Key) => Effect.Effect<void>
/** Stops the active execution, clears its doorbell, and waits for cleanup. No-op when idle. */ /** Stops the active execution, clears its doorbell, and waits for cleanup. No-op when idle. */
readonly interrupt: ( readonly interrupt: (key: Key, reason?: Reason) => Effect.Effect<void>
key: Key,
reason?: Reason,
options?: { readonly continue?: { readonly request: Request; readonly when: Effect.Effect<boolean> } },
) => Effect.Effect<void>
/** Resolves once no execution is active for the key. Returns immediately when already idle and never starts work. */ /** Resolves once no execution is active for the key. Returns immediately when already idle and never starts work. */
readonly awaitIdle: (key: Key) => Effect.Effect<void> readonly awaitIdle: (key: Key) => Effect.Effect<void>
} }
export type Request = Promotable
/** /**
* One execution is a busy period for one key: one fiber that drains from the first wake * One execution is a busy period for one key: one fiber that drains from the first wake
* until the key would stay idle. `pendingWake` is the doorbell: work recorded during the * until the key would stay idle. `pendingWake` is the doorbell: work recorded during the
* execution rings it with its eligibility request, and the execution loop drains again * execution rings it with the scope that work needs, and the execution loop drains again
* instead of ending. The doorbell closes the gap between a drain's last eligibility check * instead of ending. The doorbell closes the gap between a drain's last eligibility check
* and the idle transition, since those cannot be one atomic step. `done` resolves joiners * and the idle transition, since those cannot be one atomic step. `done` resolves joiners
* with this execution's exit. * with this execution's exit.
@@ -36,15 +30,10 @@ export type Request = Promotable
type Execution<E, Reason> = { type Execution<E, Reason> = {
readonly done: Deferred.Deferred<void, E> readonly done: Deferred.Deferred<void, E>
owner?: Fiber.Fiber<void> owner?: Fiber.Fiber<void>
request: Request scope: Promotable
pendingWake?: Request pendingWake?: Promotable
stopping: boolean stopping: boolean
interruptionReason?: Reason interruptionReason?: Reason
continuation?: {
readonly request: Request
readonly when: Effect.Effect<boolean>
signaled: boolean
}
} }
/** /**
@@ -59,7 +48,7 @@ type Execution<E, Reason> = {
* ``` * ```
*/ */
export const make = <Key, E, Reason = never>(options: { export const make = <Key, E, Reason = never>(options: {
readonly drain: (key: Key, force: boolean, request: Request) => Effect.Effect<void, E> readonly drain: (key: Key, force: boolean, scope: Promotable) => Effect.Effect<void, E>
/** Runs once when a process-local busy period begins, before its first drain. */ /** Runs once when a process-local busy period begins, before its first drain. */
readonly started?: (key: Key) => Effect.Effect<void> readonly started?: (key: Key) => Effect.Effect<void>
/** /**
@@ -73,11 +62,11 @@ export const make = <Key, E, Reason = never>(options: {
const fork = yield* FiberSet.makeRuntime<never, void, never>() const fork = yield* FiberSet.makeRuntime<never, void, never>()
const loop = (key: Key, execution: Execution<E, Reason>, force: boolean): Effect.Effect<void, E> => const loop = (key: Key, execution: Execution<E, Reason>, force: boolean): Effect.Effect<void, E> =>
Effect.suspend(() => options.drain(key, force, execution.request)).pipe( Effect.suspend(() => options.drain(key, force, execution.scope)).pipe(
Effect.flatMap(() => Effect.flatMap(() =>
Effect.suspend(() => { Effect.suspend(() => {
if (execution.stopping || execution.pendingWake === undefined) return Effect.void if (execution.stopping || execution.pendingWake === undefined) return Effect.void
execution.request = execution.pendingWake execution.scope = execution.pendingWake
execution.pendingWake = undefined execution.pendingWake = undefined
// Trampoline so drains that complete synchronously cannot grow the stack. // Trampoline so drains that complete synchronously cannot grow the stack.
return Effect.yieldNow.pipe(Effect.andThen(loop(key, execution, false))) return Effect.yieldNow.pipe(Effect.andThen(loop(key, execution, false)))
@@ -85,10 +74,10 @@ export const make = <Key, E, Reason = never>(options: {
), ),
) )
const start = (key: Key, force: boolean, request: Request) => { const start = (key: Key, force: boolean, scope: Promotable) => {
const execution: Execution<E, Reason> = { const execution: Execution<E, Reason> = {
done: Deferred.makeUnsafe<void, E>(), done: Deferred.makeUnsafe<void, E>(),
request, scope,
stopping: false, stopping: false,
} }
executions.set(key, execution) executions.set(key, execution)
@@ -104,7 +93,7 @@ export const make = <Key, E, Reason = never>(options: {
execution.owner = undefined execution.owner = undefined
}).pipe(Effect.andThen(options.settled?.(key, exit, execution.interruptionReason) ?? Effect.void)), }).pipe(Effect.andThen(options.settled?.(key, exit, execution.interruptionReason) ?? Effect.void)),
), ),
Effect.onExit((exit) => finish(key, execution, exit)), Effect.onExit((exit) => Effect.sync(() => settle(key, execution, exit))),
Effect.exit, Effect.exit,
Effect.asVoid, Effect.asVoid,
), ),
@@ -114,22 +103,12 @@ export const make = <Key, E, Reason = never>(options: {
// A doorbell that survives the execution loop (rung after the loop decided to end, or // A doorbell that survives the execution loop (rung after the loop decided to end, or
// during failure or interruption cleanup) starts a fresh execution for the remaining work. // during failure or interruption cleanup) starts a fresh execution for the remaining work.
const settle = (key: Key, execution: Execution<E, Reason>, exit: Exit.Exit<void, E>, resume: boolean) => { const settle = (key: Key, execution: Execution<E, Reason>, exit: Exit.Exit<void, E>) => {
if (resume && execution.continuation) start(key, false, execution.continuation.request) if (execution.pendingWake) start(key, false, execution.pendingWake)
else if (execution.pendingWake) start(key, false, execution.pendingWake)
else executions.delete(key) else executions.delete(key)
Deferred.doneUnsafe(execution.done, exit) Deferred.doneUnsafe(execution.done, exit)
} }
const finish = (key: Key, execution: Execution<E, Reason>, exit: Exit.Exit<void, E>) => {
if (!execution.continuation) return Effect.sync(() => settle(key, execution, exit, false))
return execution.continuation.when.pipe(
Effect.flatMap((ready) =>
Effect.sync(() => settle(key, execution, exit, ready || execution.continuation?.signaled === true)),
),
)
}
const run = (key: Key): Effect.Effect<void, E> => const run = (key: Key): Effect.Effect<void, E> =>
Effect.suspend(() => { Effect.suspend(() => {
const execution = executions.get(key) const execution = executions.get(key)
@@ -141,55 +120,32 @@ export const make = <Key, E, Reason = never>(options: {
return Deferred.await(start(key, true, "input").done) return Deferred.await(start(key, true, "input").done)
}) })
const wake = (key: Key, request: Request = "input") => const wake = (key: Key, scope: Promotable = "input") =>
Effect.sync(() => { Effect.sync(() => {
const execution = executions.get(key) const execution = executions.get(key)
if (execution !== undefined) { if (execution !== undefined) {
if (execution.stopping) { // Coalesced wakes keep the widest scope: "input" subsumes "steer".
if (execution.continuation) execution.continuation.signaled = true execution.pendingWake = execution.pendingWake === "input" ? "input" : scope
else execution.continuation = { request, when: Effect.succeed(true), signaled: true }
return
}
// Coalesced wakes keep the widest request: "input" subsumes "steer".
execution.pendingWake = execution.pendingWake === "input" ? "input" : request
return return
} }
start(key, false, request) start(key, false, scope)
}) })
const wakeActive = (key: Key) => const wakeActive = (key: Key) =>
Effect.suspend(() => { Effect.suspend(() => {
const execution = executions.get(key) const execution = executions.get(key)
return execution ? wake(key, execution.request) : Effect.void return execution ? wake(key, execution.scope) : Effect.void
}) })
const interrupt = ( const interrupt = (key: Key, reason?: Reason): Effect.Effect<void> =>
key: Key,
reason?: Reason,
options?: { readonly continue?: { readonly request: Request; readonly when: Effect.Effect<boolean> } },
): Effect.Effect<void> =>
Effect.suspend(() => { Effect.suspend(() => {
const execution = executions.get(key) const execution = executions.get(key)
if (execution === undefined) return Effect.void if (execution?.owner === undefined || execution.stopping) return Effect.void
if (execution.stopping) {
if (options?.continue)
execution.continuation = {
...options.continue,
signaled: execution.continuation?.signaled ?? false,
}
return Deferred.await(execution.done).pipe(Effect.exit, Effect.asVoid)
}
if (execution.owner === undefined) {
if (!options?.continue) return Effect.void
execution.stopping = true
execution.pendingWake = undefined
execution.continuation = { ...options.continue, signaled: false }
return Deferred.await(execution.done).pipe(Effect.exit, Effect.asVoid)
}
execution.stopping = true execution.stopping = true
// Wakes recorded so far belong to the interrupted intent; the interrupt claims them.
// Wakes arriving during cleanup are new admissions and restart normally at settle.
execution.pendingWake = undefined execution.pendingWake = undefined
execution.interruptionReason = reason execution.interruptionReason = reason
if (options?.continue) execution.continuation = { ...options.continue, signaled: false }
return Fiber.interrupt(execution.owner) return Fiber.interrupt(execution.owner)
}) })
-2
View File
@@ -56,8 +56,6 @@ export const SessionTable = sqliteTable(
variant?: string variant?: string
}>(), }>(),
...Timestamps, ...Timestamps,
time_idle: integer(),
time_viewed: integer(),
time_compacting: integer(), time_compacting: integer(),
time_archived: integer(), time_archived: integer(),
/** The execution claim timestamp (historical column name; see SessionStore.claim). */ /** The execution claim timestamp (historical column name; see SessionStore.claim). */
-4
View File
@@ -117,10 +117,6 @@ const layer = Layer.effect(
tokens_cache_write: input.data.info.tokens.cache.write, tokens_cache_write: input.data.info.tokens.cache.write,
time_created: DateTime.toEpochMillis(input.data.info.time.created), time_created: DateTime.toEpochMillis(input.data.info.time.created),
time_updated: DateTime.toEpochMillis(input.data.info.time.updated), time_updated: DateTime.toEpochMillis(input.data.info.time.updated),
time_idle: input.data.info.time.idle ? DateTime.toEpochMillis(input.data.info.time.idle) : null,
time_viewed: input.data.info.time.viewed
? DateTime.toEpochMillis(input.data.info.time.viewed)
: null,
time_archived: input.data.info.time.archived time_archived: input.data.info.time.archived
? DateTime.toEpochMillis(input.data.info.time.archived) ? DateTime.toEpochMillis(input.data.info.time.archived)
: null, : null,
@@ -13,7 +13,6 @@ import { tmpdir } from "./fixture/tmpdir"
import type { SqlClient } from "effect/unstable/sql/SqlClient" import type { SqlClient } from "effect/unstable/sql/SqlClient"
import legacyCredentialsMigration from "@opencode-ai/core/database/migration/20260805200742_import_legacy_credentials" import legacyCredentialsMigration from "@opencode-ai/core/database/migration/20260805200742_import_legacy_credentials"
import worktreeMigration from "@opencode-ai/core/database/migration/20260812213948_worktree" import worktreeMigration from "@opencode-ai/core/database/migration/20260812213948_worktree"
import sessionViewedStateMigration from "@opencode-ai/core/database/migration/20260815182818_session_viewed_state"
import { Global } from "@opencode-ai/util/global" import { Global } from "@opencode-ai/util/global"
const run = <A, E>( const run = <A, E>(
@@ -74,27 +73,6 @@ describe("DatabaseMigration", () => {
) )
}) })
test("adds nullable attention state to existing sessions", async () => {
await run(
Effect.gen(function* () {
const db = yield* makeDb
yield* db.run(sql`CREATE TABLE session_v2 (id text PRIMARY KEY, title text)`)
yield* db.run(sql`INSERT INTO session_v2 (id, title) VALUES ('ses_existing', 'Existing')`)
yield* DatabaseMigration.applyOnly(db, [sessionViewedStateMigration])
yield* DatabaseMigration.applyOnly(db, [sessionViewedStateMigration])
expect(yield* db.get(sql`SELECT id, title, time_idle, time_viewed FROM session_v2`)).toEqual({
id: "ses_existing",
title: "Existing",
time_idle: null,
time_viewed: null,
})
expect(yield* db.get(sql`SELECT count(*) AS count FROM migration`)).toEqual({ count: 1 })
}),
)
})
test("rejects a non-empty database without a session table", async () => { test("rejects a non-empty database without a session table", async () => {
await expect( await expect(
run( run(
+3 -16
View File
@@ -840,15 +840,7 @@ describe("SessionTransfer", () => {
const imported = yield* transfer.import({ const imported = yield* transfer.import({
data: { data: {
info: { info: { ...template, id: sessionID },
...template,
id: sessionID,
time: {
...template.time,
idle: DateTime.makeUnsafe(200),
viewed: DateTime.makeUnsafe(150),
},
},
messages: [ messages: [
{ {
id: sourceMessageID, id: sourceMessageID,
@@ -871,18 +863,13 @@ describe("SessionTransfer", () => {
const messages = yield* session.messages({ sessionID, order: "asc" }) const messages = yield* session.messages({ sessionID, order: "asc" })
expect(imported).toMatchObject({ id: sessionID, title: "Exported", location }) expect(imported).toMatchObject({ id: sessionID, title: "Exported", location })
expect(imported.time).toMatchObject({ idle: DateTime.makeUnsafe(200), viewed: DateTime.makeUnsafe(150) })
expect(messages).toMatchObject([ expect(messages).toMatchObject([
{ id: sourceMessageID, type: "user", text: "Imported message" }, { id: sourceMessageID, type: "user", text: "Imported message" },
{ id: errorMessageID, type: "compaction", error: { type: "test_error", message: "Original error" } }, { id: errorMessageID, type: "compaction", error: { type: "test_error", message: "Original error" } },
]) ])
expect(yield* Bus.latestSequence(db, sessionID)).toBe(2) expect(yield* Bus.latestSequence(db, sessionID)).toBe(2)
const exported = yield* transfer.export({ sessionID }) expect((yield* transfer.export({ sessionID })).messages).toEqual(messages)
expect(exported.info.time).toMatchObject({ idle: DateTime.makeUnsafe(200), viewed: DateTime.makeUnsafe(150) }) expect((yield* transfer.export({ sessionID, sanitize: true })).messages).toMatchObject([
expect(exported.messages).toEqual(messages)
const sanitized = yield* transfer.export({ sessionID, sanitize: true })
expect(sanitized.info.time).toMatchObject({ idle: DateTime.makeUnsafe(200), viewed: DateTime.makeUnsafe(150) })
expect(sanitized.messages).toMatchObject([
{ id: sourceMessageID, text: `[redacted:text:${sourceMessageID}]` }, { id: sourceMessageID, text: `[redacted:text:${sourceMessageID}]` },
{ id: errorMessageID, error: { type: "test_error", message: "Original error" } }, { id: errorMessageID, error: { type: "test_error", message: "Original error" } },
]) ])
+110 -1
View File
@@ -14,8 +14,10 @@ import { SessionExecution } from "@opencode-ai/core/session/execution"
import { SessionRestart } from "@opencode-ai/core/session/execution/restart" import { SessionRestart } from "@opencode-ai/core/session/execution/restart"
import { UserInterruptedError } from "@opencode-ai/core/session/error" import { UserInterruptedError } from "@opencode-ai/core/session/error"
import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionEvent } from "@opencode-ai/core/session/event"
import { SessionInbox } from "@opencode-ai/core/session/inbox"
import { SessionMessage } from "@opencode-ai/core/session/message"
import { SessionRunner } from "@opencode-ai/core/session/runner/index" import { SessionRunner } from "@opencode-ai/core/session/runner/index"
import { SessionTable } from "@opencode-ai/core/session/sql" import { SessionInboxTable, SessionTable } from "@opencode-ai/core/session/sql"
import { SessionStore } from "@opencode-ai/core/session/store" import { SessionStore } from "@opencode-ai/core/session/store"
import { Context, Deferred, Effect, Exit, Fiber, Layer, LayerMap, Scope } from "effect" import { Context, Deferred, Effect, Exit, Fiber, Layer, LayerMap, Scope } from "effect"
import { eq } from "drizzle-orm" import { eq } from "drizzle-orm"
@@ -290,6 +292,113 @@ describe("SessionExecution lifecycle", () => {
) )
}) })
describe("SessionExecution interrupt continuation", () => {
it.effect("resumes only steering input after an interrupt with continue", () =>
Effect.gen(function* () {
const database = yield* Database.Service
const sessionID = Session.ID.make("ses_continue_steer")
yield* seedSessions(database, [sessionID])
yield* seedInbox(database, sessionID, ["steer", "queue"])
const draining = yield* Deferred.make<void>()
const drains: Array<{ force: boolean; promotable?: SessionInbox.Promotable }> = []
const scope = yield* Scope.make()
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
const context = yield* buildExecution(scope, (input) =>
Effect.suspend(() => {
drains.push({ force: input.force, promotable: input.promotable })
if (drains.length > 1) return Effect.void
return Deferred.succeed(draining, undefined).pipe(Effect.andThen(Effect.never))
}),
)
const execution = Context.get(context, SessionExecution.Service)
yield* execution.resume(sessionID).pipe(Effect.forkScoped)
yield* Deferred.await(draining)
yield* execution.interrupt(sessionID, { continue: true })
yield* execution.awaitIdle(sessionID)
// The successor drain is steer-scoped: queued next-turn work stays parked.
expect(drains).toEqual([
{ force: true, promotable: "input" },
{ force: false, promotable: "steer" },
])
}),
)
it.effect("stays parked after an interrupt with continue when only queued work remains", () =>
Effect.gen(function* () {
const database = yield* Database.Service
const sessionID = Session.ID.make("ses_continue_parked")
yield* seedSessions(database, [sessionID])
yield* seedInbox(database, sessionID, ["queue"])
const draining = yield* Deferred.make<void>()
const drains: Array<SessionInbox.Promotable | undefined> = []
const scope = yield* Scope.make()
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
const context = yield* buildExecution(scope, (input) =>
Effect.suspend(() => {
drains.push(input.promotable)
return Deferred.succeed(draining, undefined).pipe(Effect.andThen(Effect.never))
}),
)
const execution = Context.get(context, SessionExecution.Service)
yield* execution.resume(sessionID).pipe(Effect.forkScoped)
yield* Deferred.await(draining)
yield* execution.interrupt(sessionID, { continue: true })
yield* execution.awaitIdle(sessionID)
expect(drains).toEqual(["input"])
expect(yield* execution.active).toEqual(new Set())
}),
)
it.effect("an idle interrupt with continue resumes pending steers", () =>
Effect.gen(function* () {
const database = yield* Database.Service
const sessionID = Session.ID.make("ses_continue_idle")
yield* seedSessions(database, [sessionID])
yield* seedInbox(database, sessionID, ["steer"])
const drains: Array<{ force: boolean; promotable?: SessionInbox.Promotable }> = []
const scope = yield* Scope.make()
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
const context = yield* buildExecution(scope, (input) =>
Effect.sync(() => void drains.push({ force: input.force, promotable: input.promotable })),
)
const execution = Context.get(context, SessionExecution.Service)
yield* execution.interrupt(sessionID, { continue: true })
yield* execution.awaitIdle(sessionID)
expect(drains).toEqual([{ force: false, promotable: "steer" }])
}),
)
})
function seedInbox(
database: Database.Service["Service"],
sessionID: Session.ID,
deliveries: ReadonlyArray<SessionInbox.Delivery>,
) {
return database.db
.insert(SessionInboxTable)
.values(
deliveries.map((delivery, index) => ({
id: SessionMessage.ID.create(),
session_id: sessionID,
type: "compaction" as const,
payload: {},
delivery,
enqueued_seq: index + 1,
})),
)
.run()
.pipe(Effect.orDie)
}
function seedSessions( function seedSessions(
database: Database.Service["Service"], database: Database.Service["Service"],
sessionIDs: ReadonlyArray<Session.ID>, sessionIDs: ReadonlyArray<Session.ID>,
@@ -1,5 +1,6 @@
import { describe, expect } from "bun:test" import { describe, expect } from "bun:test"
import { Cause, Deferred, Effect, Exit, Fiber, Layer } from "effect" import { Cause, Deferred, Effect, Exit, Fiber, Layer } from "effect"
import { SessionInbox } from "@opencode-ai/core/session/inbox"
import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator" import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator"
import { testEffect } from "./lib/effect" import { testEffect } from "./lib/effect"
@@ -269,31 +270,24 @@ describe("SessionRunCoordinator", () => {
), ),
) )
it.effect("replaces a settlement-window wake with a steer continuation", () => it.effect("a settlement-window wake starts a fresh execution with its own scope", () =>
Effect.scoped( Effect.scoped(
Effect.gen(function* () { Effect.gen(function* () {
const settling = yield* Deferred.make<void>() const settling = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>() const release = yield* Deferred.make<void>()
const requests: SessionRunCoordinator.Request[] = [] const scopes: SessionInbox.Promotable[] = []
const coordinator = yield* SessionRunCoordinator.make({ const coordinator = yield* SessionRunCoordinator.make({
drain: (_key, _force, request) => Effect.sync(() => requests.push(request)), drain: (_key, _force, scope) => Effect.sync(() => scopes.push(scope)),
settled: () => Deferred.succeed(settling, undefined).pipe(Effect.andThen(Deferred.await(release))), settled: () => Deferred.succeed(settling, undefined).pipe(Effect.andThen(Deferred.await(release))),
}) })
yield* coordinator.wake("session", "input") yield* coordinator.wake("session", "steer")
yield* Deferred.await(settling) yield* Deferred.await(settling)
yield* coordinator.wake("session", "input") yield* coordinator.wake("session", "input")
const interrupted = yield* coordinator
.interrupt("session", undefined, {
continue: { request: "steer", when: Effect.succeed(true) },
})
.pipe(Effect.forkChild)
yield* Effect.yieldNow
yield* Deferred.succeed(release, undefined) yield* Deferred.succeed(release, undefined)
yield* Fiber.join(interrupted)
yield* coordinator.awaitIdle("session") yield* coordinator.awaitIdle("session")
expect(requests).toEqual(["input", "steer"]) expect(scopes).toEqual(["steer", "input"])
}), }),
), ),
) )
@@ -371,17 +365,17 @@ describe("SessionRunCoordinator", () => {
), ),
) )
it.effect("coalesces drain requests with input taking precedence", () => it.effect("coalesces drain scopes with input taking precedence", () =>
Effect.scoped( Effect.scoped(
Effect.gen(function* () { Effect.gen(function* () {
const firstStarted = yield* Deferred.make<void>() const firstStarted = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>() const release = yield* Deferred.make<void>()
const requests: SessionRunCoordinator.Request[] = [] const scopes: SessionInbox.Promotable[] = []
const coordinator = yield* SessionRunCoordinator.make({ const coordinator = yield* SessionRunCoordinator.make({
drain: (_key, _force, request) => drain: (_key, _force, scope) =>
Effect.gen(function* () { Effect.gen(function* () {
requests.push(request) scopes.push(scope)
if (requests.length !== 1) return if (scopes.length !== 1) return
yield* Deferred.succeed(firstStarted, undefined) yield* Deferred.succeed(firstStarted, undefined)
yield* Deferred.await(release) yield* Deferred.await(release)
}), }),
@@ -394,22 +388,22 @@ describe("SessionRunCoordinator", () => {
yield* Deferred.succeed(release, undefined) yield* Deferred.succeed(release, undefined)
yield* coordinator.awaitIdle("session") yield* coordinator.awaitIdle("session")
expect(requests).toEqual(["steer", "input"]) expect(scopes).toEqual(["steer", "input"])
}), }),
), ),
) )
it.effect("does not carry a completed input request into a steer drain", () => it.effect("does not carry a completed input scope into a steer drain", () =>
Effect.scoped( Effect.scoped(
Effect.gen(function* () { Effect.gen(function* () {
const firstStarted = yield* Deferred.make<void>() const firstStarted = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>() const release = yield* Deferred.make<void>()
const requests: SessionRunCoordinator.Request[] = [] const scopes: SessionInbox.Promotable[] = []
const coordinator = yield* SessionRunCoordinator.make({ const coordinator = yield* SessionRunCoordinator.make({
drain: (_key, _force, request) => drain: (_key, _force, scope) =>
Effect.gen(function* () { Effect.gen(function* () {
requests.push(request) scopes.push(scope)
if (requests.length !== 1) return if (scopes.length !== 1) return
yield* Deferred.succeed(firstStarted, undefined) yield* Deferred.succeed(firstStarted, undefined)
yield* Deferred.await(release) yield* Deferred.await(release)
}), }),
@@ -421,7 +415,7 @@ describe("SessionRunCoordinator", () => {
yield* Deferred.succeed(release, undefined) yield* Deferred.succeed(release, undefined)
yield* coordinator.awaitIdle("session") yield* coordinator.awaitIdle("session")
expect(requests).toEqual(["input", "steer"]) expect(scopes).toEqual(["input", "steer"])
}), }),
), ),
) )
@@ -431,12 +425,12 @@ describe("SessionRunCoordinator", () => {
Effect.gen(function* () { Effect.gen(function* () {
const firstStarted = yield* Deferred.make<void>() const firstStarted = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>() const release = yield* Deferred.make<void>()
const requests: SessionRunCoordinator.Request[] = [] const scopes: SessionInbox.Promotable[] = []
const coordinator = yield* SessionRunCoordinator.make({ const coordinator = yield* SessionRunCoordinator.make({
drain: (_key, _force, request) => drain: (_key, _force, scope) =>
Effect.gen(function* () { Effect.gen(function* () {
requests.push(request) scopes.push(scope)
if (requests.length !== 1) return if (scopes.length !== 1) return
yield* Deferred.succeed(firstStarted, undefined) yield* Deferred.succeed(firstStarted, undefined)
yield* Deferred.await(release) yield* Deferred.await(release)
}), }),
@@ -449,61 +443,23 @@ describe("SessionRunCoordinator", () => {
yield* Deferred.succeed(release, undefined) yield* Deferred.succeed(release, undefined)
yield* coordinator.awaitIdle("session") yield* coordinator.awaitIdle("session")
expect(requests).toEqual(["steer", "steer"]) expect(scopes).toEqual(["steer", "steer"])
}), }),
), ),
) )
it.effect("coalesces overlapping interrupt continuations into one steer successor", () => it.effect("a cleanup-era wake starts a successor with its own scope", () =>
Effect.scoped( Effect.scoped(
Effect.gen(function* () { Effect.gen(function* () {
const firstStarted = yield* Deferred.make<void>() const firstStarted = yield* Deferred.make<void>()
const cleanupStarted = yield* Deferred.make<void>() const cleanupStarted = yield* Deferred.make<void>()
const cleanupGate = yield* Deferred.make<void>() const cleanupGate = yield* Deferred.make<void>()
const requests: SessionRunCoordinator.Request[] = [] const scopes: SessionInbox.Promotable[] = []
const coordinator = yield* SessionRunCoordinator.make({ const coordinator = yield* SessionRunCoordinator.make({
drain: (_key, _force, request) => drain: (_key, _force, scope) =>
Effect.gen(function* () { Effect.gen(function* () {
requests.push(request) scopes.push(scope)
if (requests.length !== 1) return if (scopes.length !== 1) return
yield* Deferred.succeed(firstStarted, undefined)
yield* Effect.never.pipe(
Effect.onInterrupt(() =>
Deferred.succeed(cleanupStarted, undefined).pipe(Effect.andThen(Deferred.await(cleanupGate))),
),
)
}),
})
const continuation = { continue: { request: "steer" as const, when: Effect.succeed(false) } }
yield* coordinator.wake("session")
yield* Deferred.await(firstStarted)
const first = yield* coordinator.interrupt("session", undefined, continuation).pipe(Effect.forkChild)
yield* Deferred.await(cleanupStarted)
const second = yield* coordinator.interrupt("session", undefined, continuation).pipe(Effect.forkChild)
yield* Effect.yieldNow
yield* coordinator.wake("session", "input")
yield* Deferred.succeed(cleanupGate, undefined)
yield* Effect.all([Fiber.join(first), Fiber.join(second)])
yield* coordinator.awaitIdle("session")
expect(requests).toEqual(["input", "steer"])
}),
),
)
it.effect("a continuing interrupt replaces a cleanup-era input wake", () =>
Effect.scoped(
Effect.gen(function* () {
const firstStarted = yield* Deferred.make<void>()
const cleanupStarted = yield* Deferred.make<void>()
const cleanupGate = yield* Deferred.make<void>()
const requests: SessionRunCoordinator.Request[] = []
const coordinator = yield* SessionRunCoordinator.make({
drain: (_key, _force, request) =>
Effect.gen(function* () {
requests.push(request)
if (requests.length !== 1) return
yield* Deferred.succeed(firstStarted, undefined) yield* Deferred.succeed(firstStarted, undefined)
yield* Effect.never.pipe( yield* Effect.never.pipe(
Effect.onInterrupt(() => Effect.onInterrupt(() =>
@@ -515,45 +471,16 @@ describe("SessionRunCoordinator", () => {
yield* coordinator.wake("session", "input") yield* coordinator.wake("session", "input")
yield* Deferred.await(firstStarted) yield* Deferred.await(firstStarted)
const plain = yield* coordinator.interrupt("session").pipe(Effect.forkChild) const interrupt = yield* coordinator.interrupt("session").pipe(Effect.forkChild)
yield* Deferred.await(cleanupStarted) yield* Deferred.await(cleanupStarted)
// A new admission during cancellation restarts normally: interruption only
// claims the wakes recorded before it.
yield* coordinator.wake("session", "input") yield* coordinator.wake("session", "input")
const continuing = yield* coordinator
.interrupt("session", undefined, {
continue: { request: "steer", when: Effect.succeed(false) },
})
.pipe(Effect.forkChild)
yield* Effect.yieldNow
yield* Deferred.succeed(cleanupGate, undefined) yield* Deferred.succeed(cleanupGate, undefined)
yield* Effect.all([Fiber.join(plain), Fiber.join(continuing)]) yield* Fiber.join(interrupt)
yield* coordinator.awaitIdle("session") yield* coordinator.awaitIdle("session")
expect(requests).toEqual(["input", "steer"]) expect(scopes).toEqual(["input", "input"])
}),
),
)
it.effect("does not start a conditional continuation without eligible work", () =>
Effect.scoped(
Effect.gen(function* () {
const started = yield* Deferred.make<void>()
const requests: SessionRunCoordinator.Request[] = []
const coordinator = yield* SessionRunCoordinator.make({
drain: (_key, _force, request) =>
Effect.sync(() => requests.push(request)).pipe(
Effect.andThen(Deferred.succeed(started, undefined)),
Effect.andThen(Effect.never),
),
})
yield* coordinator.wake("session")
yield* Deferred.await(started)
yield* coordinator.interrupt("session", undefined, {
continue: { request: "steer", when: Effect.succeed(false) },
})
yield* coordinator.awaitIdle("session")
expect(requests).toEqual(["input"])
}), }),
), ),
) )
-114
View File
@@ -1,114 +0,0 @@
import { describe, expect } from "bun:test"
import { Bus } from "@opencode-ai/core/bus"
import { Database } from "@opencode-ai/core/database/database"
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
import { EventTable } from "@opencode-ai/core/event/sql"
import { Location } from "@opencode-ai/core/location"
import { Project } from "@opencode-ai/core/project"
import { AbsolutePath } from "@opencode-ai/core/schema"
import { Session } from "@opencode-ai/core/session"
import { SessionEvent } from "@opencode-ai/core/session/event"
import { SessionExecution } from "@opencode-ai/core/session/execution"
import { SessionProjector } from "@opencode-ai/core/session/projector"
import { SessionTable } from "@opencode-ai/core/session/sql"
import { SessionStore } from "@opencode-ai/core/session/store"
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
import { DateTime, Effect } from "effect"
import { eq } from "drizzle-orm"
import { testEffect } from "./lib/effect"
import { globalProjectLayer } from "./lib/project"
const it = testEffect(
AppNodeBuilder.build(
LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node, Session.node]),
[
[Bus.node, Bus.configured({ persist: true })],
[Project.node, globalProjectLayer],
[SessionExecution.node, SessionExecution.noopLayer],
],
),
)
const location = Location.Ref.make({ directory: AbsolutePath.make("/project") })
describe("Session.view", () => {
it.effect("copies the latest idle time without changing session recency", () =>
Effect.gen(function* () {
const session = yield* Session.Service
const bus = yield* Bus.Service
const { db } = yield* Database.Service
const created = yield* session.create({ location })
expect(created.time.idle).toBeUndefined()
expect(created.time.viewed).toBeUndefined()
yield* Effect.all([session.view({ sessionID: created.id }), session.view({ sessionID: created.id })], {
concurrency: "unbounded",
discard: true,
})
expect((yield* session.get(created.id)).time.viewed).toBeUndefined()
yield* bus.publish(SessionEvent.Execution.Succeeded, { sessionID: created.id })
const idle = yield* session.get(created.id)
expect(idle.time.idle).toBeDefined()
expect(idle.time.viewed).toBeUndefined()
expect(idle.time.updated).toEqual(created.time.updated)
yield* session.view({ sessionID: created.id })
const viewed = yield* session.get(created.id)
if (!viewed.time.idle || !viewed.time.viewed) return yield* Effect.die(new Error("Expected attention times"))
expect(viewed.time.viewed).toEqual(viewed.time.idle)
expect(viewed.time.updated).toEqual(created.time.updated)
expect(
yield* db
.select({ idle: SessionTable.time_idle, viewed: SessionTable.time_viewed })
.from(SessionTable)
.where(eq(SessionTable.id, created.id))
.get(),
).toEqual({
idle: DateTime.toEpochMillis(viewed.time.idle),
viewed: DateTime.toEpochMillis(viewed.time.viewed),
})
expect((yield* session.list()).data.find((item) => item.id === created.id)?.time).toEqual(viewed.time)
yield* session.view({ sessionID: created.id })
expect((yield* session.get(created.id)).time).toEqual(viewed.time)
yield* bus.publish(SessionEvent.Execution.Failed, {
sessionID: created.id,
error: { type: "unknown", message: "failed" },
})
const unread = yield* session.get(created.id)
if (!unread.time.idle || !unread.time.viewed) return yield* Effect.die(new Error("Expected attention times"))
expect(DateTime.toEpochMillis(unread.time.idle)).toBeGreaterThan(DateTime.toEpochMillis(unread.time.viewed))
yield* session.view({ sessionID: created.id })
expect((yield* session.get(created.id)).time.viewed).toEqual(unread.time.idle)
yield* bus.publish(SessionEvent.Execution.Interrupted, { sessionID: created.id, reason: "shutdown" })
expect((yield* session.get(created.id)).time.idle).toEqual(unread.time.idle)
yield* bus.publish(SessionEvent.Execution.Interrupted, { sessionID: created.id, reason: "user" })
const interrupted = yield* session.get(created.id)
if (!interrupted.time.idle || !interrupted.time.viewed)
return yield* Effect.die(new Error("Expected attention times"))
expect(DateTime.toEpochMillis(interrupted.time.idle)).toBeGreaterThan(
DateTime.toEpochMillis(interrupted.time.viewed),
)
expect(
(yield* db
.select({ type: EventTable.type })
.from(EventTable)
.where(eq(EventTable.aggregate_id, created.id))
.all()).filter((event) => event.type === "session.viewed.1"),
).toHaveLength(2)
}),
)
it.effect("rejects an unknown session", () =>
Effect.gen(function* () {
const session = yield* Session.Service
const sessionID = Session.ID.make("ses_missing_view")
expect(yield* Effect.flip(session.view({ sessionID }))).toEqual(new Session.NotFoundError({ sessionID }))
}),
)
})
-2
View File
@@ -60,8 +60,6 @@ const session = (
model: null, model: null,
time_created: 1, time_created: 1,
time_updated: 2, time_updated: 2,
time_idle: null,
time_viewed: null,
time_compacting: 3, time_compacting: 3,
time_archived: null, time_archived: null,
time_suspended: null, time_suspended: null,
File diff suppressed because it is too large Load Diff
-13
View File
@@ -693,19 +693,6 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(sessionLo
}), }),
), ),
) )
.add(
HttpApiEndpoint.post("session.view", "/api/session/:sessionID/view", {
params: { sessionID: Session.ID },
success: HttpApiSchema.NoContent,
error: SessionNotFoundError,
}).annotateMerge(
OpenApi.annotations({
identifier: "v2.session.view",
summary: "View session",
description: "Mark the latest recorded idle transition as viewed.",
}),
),
)
.annotateMerge( .annotateMerge(
OpenApi.annotations({ OpenApi.annotations({
title: "session", title: "session",
-8
View File
@@ -105,13 +105,6 @@ export const Renamed = Event.durable({
}) })
export type Renamed = typeof Renamed.Type export type Renamed = typeof Renamed.Type
export const Viewed = Event.durable({
type: "session.viewed",
...options,
schema: Base,
})
export type Viewed = typeof Viewed.Type
export const UsageRecorded = Event.durable({ export const UsageRecorded = Event.durable({
type: "session.usage.recorded", type: "session.usage.recorded",
...options, ...options,
@@ -587,7 +580,6 @@ export const Definitions = Event.inventory(
ModelSelected, ModelSelected,
Moved, Moved,
Renamed, Renamed,
Viewed,
UsageUpdated, UsageUpdated,
Deleted, Deleted,
Forked, Forked,
-2
View File
@@ -40,8 +40,6 @@ export const Info = Schema.Struct({
time: Schema.Struct({ time: Schema.Struct({
created: DateTimeUtcFromMillis, created: DateTimeUtcFromMillis,
updated: DateTimeUtcFromMillis, updated: DateTimeUtcFromMillis,
idle: DateTimeUtcFromMillis.pipe(optional),
viewed: DateTimeUtcFromMillis.pipe(optional),
archived: DateTimeUtcFromMillis.pipe(optional), archived: DateTimeUtcFromMillis.pipe(optional),
}), }),
title: Schema.String.pipe(optional), title: Schema.String.pipe(optional),
+9 -21
View File
@@ -54,29 +54,17 @@ describe("contract hygiene", () => {
}), }),
).toEqual({ text: "completed" }) ).toEqual({ text: "completed" })
const info = Session.Info.make({
id: Session.ID.make("ses_untitled"),
projectID: Project.ID.make("global"),
cost: Money.USD.zero,
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
time: {
created: DateTime.makeUnsafe(0),
updated: DateTime.makeUnsafe(0),
idle: undefined,
viewed: undefined,
},
title: undefined,
location: { directory: AbsolutePath.make("/project") },
})
const encoded = Schema.encodeSync(Session.Info)(info)
expect(encoded).not.toHaveProperty("title")
expect(encoded.time).toEqual({ created: 0, updated: 0 })
expect( expect(
Schema.encodeSync(Session.Info)({ Schema.encodeSync(Session.Info)({
...info, id: Session.ID.make("ses_untitled"),
time: { ...info.time, idle: DateTime.makeUnsafe(2), viewed: DateTime.makeUnsafe(1) }, projectID: Project.ID.make("global"),
}).time, cost: Money.USD.zero,
).toEqual({ created: 0, updated: 0, idle: 2, viewed: 1 }) tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
time: { created: DateTime.makeUnsafe(0), updated: DateTime.makeUnsafe(0) },
title: undefined,
location: { directory: AbsolutePath.make("/project") },
}),
).not.toHaveProperty("title")
}) })
test("session inbox items omit the internal enqueue sequence", () => { test("session inbox items omit the internal enqueue sequence", () => {
@@ -83,7 +83,6 @@ describe("public event manifest", () => {
"session.model.selected.1", "session.model.selected.1",
"session.moved.1", "session.moved.1",
"session.renamed.1", "session.renamed.1",
"session.viewed.1",
"session.usage.recorded.1", "session.usage.recorded.1",
"session.forked.2", "session.forked.2",
"session.inbox.delivered.1", "session.inbox.delivered.1",
-16
View File
@@ -180,22 +180,6 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
} }
}), }),
) )
.handle(
"session.view",
Effect.fn(function* (ctx) {
yield* session.view({ sessionID: ctx.params.sessionID }).pipe(
Effect.catchTag(
"Session.NotFoundError",
(error) =>
new SessionNotFoundError({
sessionID: error.sessionID,
message: `Session not found: ${error.sessionID}`,
}),
),
)
return HttpApiSchema.NoContent.make()
}),
)
.handle( .handle(
"session.remove", "session.remove",
Effect.fn(function* (ctx) { Effect.fn(function* (ctx) {
-30
View File
@@ -52,36 +52,6 @@ it.live("serves unauthenticated and answers CORS preflight when no password is c
}).pipe(Effect.scoped), }).pipe(Effect.scoped),
) )
it.live("serves the session view operation and missing-session error", () =>
Effect.gen(function* () {
const handler = yield* ServerFetch.make(options)
const created = yield* Effect.promise(() =>
handler(
new Request("http://opencode.local/api/session", {
method: "POST",
headers: { "content-type": "application/json" },
body: "{}",
}),
).then((response) => response.json()),
)
if (typeof created !== "object" || created === null || !("data" in created))
return yield* Effect.die(new Error("Expected a session response"))
const data = created.data
if (typeof data !== "object" || data === null || !("id" in data) || typeof data.id !== "string")
return yield* Effect.die(new Error("Expected a session ID"))
const viewed = yield* Effect.promise(() =>
handler(new Request(`http://opencode.local/api/session/${data.id}/view`, { method: "POST" })),
)
expect(viewed.status).toBe(204)
const missing = yield* Effect.promise(() =>
handler(new Request("http://opencode.local/api/session/ses_missing_view/view", { method: "POST" })),
)
expect(missing.status).toBe(404)
}).pipe(Effect.scoped),
)
// Pins the eager-boot guarantee: the application layer is built before the handler returns, so // Pins the eager-boot guarantee: the application layer is built before the handler returns, so
// an aborted first request cannot interrupt layer construction and wedge every later request // an aborted first request cannot interrupt layer construction and wedge every later request
// (the Effect-TS/effect#6319 failure class that lazy first-request builds are prone to). // (the Effect-TS/effect#6319 failure class that lazy first-request builds are prone to).
+1094 -4465
View File
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff