mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-15 17:08:21 -04:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6cb2f00d15 |
@@ -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.
|
||||||
@@ -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) =>
|
||||||
|
|||||||
@@ -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?: {
|
||||||
|
|||||||
@@ -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" },
|
||||||
|
|||||||
@@ -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
@@ -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
@@ -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
|
|
||||||
@@ -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,
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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* () {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|||||||
@@ -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). */
|
||||||
|
|||||||
@@ -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(
|
||||||
|
|||||||
@@ -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" } },
|
||||||
])
|
])
|
||||||
|
|||||||
@@ -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"])
|
|
||||||
}),
|
}),
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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 }))
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
})
|
|
||||||
@@ -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,
|
||||||
|
|||||||
+1094
-4465
File diff suppressed because it is too large
Load Diff
@@ -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",
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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),
|
||||||
|
|||||||
@@ -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",
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -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
File diff suppressed because it is too large
Load Diff
+1094
-4465
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user