From 0fa120d1b956319d87b045568190d522cc095ddc Mon Sep 17 00:00:00 2001 From: Khoa Huynh Date: Sun, 27 Sep 2026 00:44:12 -0400 Subject: [PATCH 1/3] fix(core): preserve progressing sessions during location cleanup --- packages/core/src/job.ts | 8 + packages/core/src/location-activity.ts | 141 +++++++--- packages/core/test/browser-idle.test.ts | 3 +- packages/core/test/location-activity.test.ts | 262 +++++++++++++++++-- 4 files changed, 359 insertions(+), 55 deletions(-) diff --git a/packages/core/src/job.ts b/packages/core/src/job.ts index 234f47d1752d..7ab0093152d7 100644 --- a/packages/core/src/job.ts +++ b/packages/core/src/job.ts @@ -130,6 +130,8 @@ export interface Interface { readonly start: (input: StartInput) => Effect.Effect readonly wait: (input: WaitInput) => Effect.Effect readonly block: (input: BlockInput) => Effect.Effect + /** Whether a running job is currently blocking this Session. Background jobs do not count. */ + readonly isBlocking: (input: BlockInput) => Effect.Effect readonly background: (id: string) => Effect.Effect readonly backgroundAll: (input: BackgroundAllInput) => Effect.Effect readonly cancel: (id: string) => Effect.Effect @@ -356,6 +358,11 @@ export const make = Effect.gen(function* () { ) }) + const isBlocking: Interface["isBlocking"] = Effect.fnUntraced(function* (input) { + const job = (yield* SynchronizedRef.get(state.jobs)).get(input.id) + return job?.blockingSessions.has(input.sessionID) ?? false + }) + const markBackground = Effect.fnUntraced(function* (job: Active) { const next = { ...job, @@ -469,6 +476,7 @@ export const make = Effect.gen(function* () { start, wait, block, + isBlocking, background, backgroundAll, cancel, diff --git a/packages/core/src/location-activity.ts b/packages/core/src/location-activity.ts index 0425e9acab04..2e629cad27eb 100644 --- a/packages/core/src/location-activity.ts +++ b/packages/core/src/location-activity.ts @@ -1,15 +1,21 @@ export * as LocationActivity from "./location-activity.js" import { Clock, Context, Duration, Effect, Layer, RcMap, Schema } from "effect" +import { Event } from "@opencode/schema/event" +import { Permission } from "@opencode/schema/permission" +import { Form } from "@opencode/schema/form" import { Bus } from "./bus.js" +import { Job } from "./job.js" import { Location } from "./location.js" import { LocationServiceMap } from "./location-service-map.js" import { SessionEvent } from "./session/event.js" import { SessionExecution } from "./session/execution.js" import { SessionStore } from "./session/store.js" +import { SessionSchema } from "./session/schema.js" import { makeGlobalNode } from "@opencode/util/effect/app-node" -const isSessionEvent = Schema.is(SessionEvent.Durable) +const isSessionEvent = (event: Event.Payload): event is SessionEvent.Event => + Object.hasOwn(SessionEvent.All.cases, event.type) export class Service extends Context.Service()("@opencode/LocationActivity") {} @@ -21,23 +27,83 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly const bus = yield* Bus.Service const locations = yield* LocationServiceMap.Service const execution = yield* SessionExecution.Service + const jobs = yield* Job.Service const sessions = yield* SessionStore.Service const timeToLive = Duration.toMillis(options.timeToLive ?? "60 minutes") const entries = new Map() + const progress = new Map() + const parents = new Map() + const waits = new Map>() + const trackWait = (sessionID: SessionSchema.ID, id: string) => { + if (!progress.has(sessionID)) return + const pending = waits.get(sessionID) ?? new Map() + pending.set(id, clock.currentTimeMillisUnsafe() + timeToLive) + waits.set(sessionID, pending) + } + const clearWait = (sessionID: SessionSchema.ID, id: string) => { + const pending = waits.get(sessionID) + pending?.delete(id) + if (pending?.size === 0) waits.delete(sessionID) + if (progress.has(sessionID)) progress.set(sessionID, clock.currentTimeMillisUnsafe() + timeToLive) + } const key = (ref: Location.Ref) => `${LocationServiceMap.canonical(ref).directory}\0${ref.workspaceID ?? ""}` const touch = (ref: Location.Ref) => Effect.sync(() => { entries.set(key(ref), { ref, expiresAt: clock.currentTimeMillisUnsafe() + timeToLive }) }) - const unsubscribe = yield* bus.listen((event) => { - if (!isSessionEvent(event)) return Effect.void - const location = event.location - if (!location) return Effect.void - return RcMap.has(locations.rcMap, location).pipe( - Effect.flatMap((active) => (active ? touch(location) : Effect.void)), - ) - }) + const unsubscribe = yield* bus.listen((event) => + Effect.gen(function* () { + if (event.type === Permission.Event.Asked.type && Schema.is(Permission.Event.Asked)(event)) { + trackWait(event.data.sessionID, event.data.id) + return + } + if (event.type === Permission.Event.Replied.type && Schema.is(Permission.Event.Replied)(event)) { + clearWait(event.data.sessionID, event.data.requestID) + return + } + if (event.type === Form.Event.Created.type && Schema.is(Form.Event.Created)(event)) { + const sessionID = event.data.form.sessionID + if (Schema.is(SessionSchema.ID)(sessionID)) trackWait(sessionID, event.data.form.id) + return + } + if ( + (event.type === Form.Event.Replied.type && Schema.is(Form.Event.Replied)(event)) || + (event.type === Form.Event.Cancelled.type && Schema.is(Form.Event.Cancelled)(event)) + ) { + const sessionID = event.data.sessionID + if (Schema.is(SessionSchema.ID)(sessionID)) clearWait(sessionID, event.data.id) + return + } + if (!isSessionEvent(event)) return + const sessionID = event.data.sessionID + if (event.type !== SessionEvent.Viewed.type) { + if (event.type === SessionEvent.Execution.Started.type || progress.has(sessionID)) + progress.set(sessionID, clock.currentTimeMillisUnsafe() + timeToLive) + // Parentage alone also includes background jobs; only a blocking chain carries progress. + for ( + let child = sessionID, parent = parents.get(child); + parent; + child = parent, parent = parents.get(child) + ) { + if (!progress.has(parent) || !(yield* jobs.isBlocking({ id: child, sessionID: parent }))) break + progress.set(parent, clock.currentTimeMillisUnsafe() + timeToLive) + } + } + if ( + event.type === SessionEvent.Execution.Succeeded.type || + event.type === SessionEvent.Execution.Failed.type || + event.type === SessionEvent.Execution.Interrupted.type + ) { + progress.delete(sessionID) + parents.delete(sessionID) + waits.delete(sessionID) + } + if (!event.durable) return + const location = event.location + if (location && (yield* RcMap.has(locations.rcMap, location))) yield* touch(location) + }), + ) yield* Effect.addFinalizer(() => unsubscribe) yield* Effect.gen(function* () { yield* Effect.sleep(options.sweepInterval ?? "1 minute") @@ -48,32 +114,47 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly if (!cached.has(id)) entries.delete(id) } const now = clock.currentTimeMillisUnsafe() + const activeIDs = yield* execution.active + yield* Effect.forEach( + Array.from(activeIDs).filter((id) => !parents.has(id)), + (id) => + Effect.gen(function* () { + const session = yield* sessions.get(id) + parents.set(id, session?.parentID ?? null) + if ( + session?.parentID && + progress.has(session.parentID) && + (yield* jobs.isBlocking({ id, sessionID: session.parentID })) + ) + progress.set(session.parentID, now + timeToLive) + }), + { discard: true }, + ) + yield* Effect.forEach( + activeIDs, + (sessionID) => { + if (!progress.has(sessionID)) progress.set(sessionID, now + timeToLive) + const waiting = Array.from(waits.get(sessionID)?.values() ?? []).reduce( + (earliest, wait) => (earliest === undefined || wait < earliest ? wait : earliest), + undefined, + ) + if ((waiting ?? progress.get(sessionID) ?? 0) > now) return Effect.void + return execution.interrupt(sessionID, { reason: "inactivity" }).pipe(Effect.asVoid) + }, + { discard: true, concurrency: "unbounded" }, + ) const expired = Array.from(entries.values()).filter((entry) => entry.expiresAt <= now) if (expired.length === 0) return - const active = yield* Effect.forEach(yield* execution.active, (sessionID) => sessions.get(sessionID)) yield* Effect.forEach( expired, (entry) => Effect.gen(function* () { - const owners = active.flatMap((session) => - session && key(session.location) === key(entry.ref) ? [session] : [], - ) - // Invalidation only detaches the cache entry; borrowers retain the old - // graph. Stop its executions and settle tool cleanup before detaching it. - yield* Effect.forEach( - owners, - (session) => execution.interrupt(session.id, { reason: "inactivity", awaitSettlement: true }), - { - discard: true, - concurrency: "unbounded", - }, - ) - const remaining = yield* Effect.forEach(yield* execution.active, (sessionID) => sessions.get(sessionID)) - // New work admitted during cleanup may now own the cached graph. - if (remaining.some((session) => session && key(session.location) === key(entry.ref))) { - yield* touch(entry.ref) - return - } + // Invalidation detaches the graph while borrowers still hold it. Active + // executions retain their Location until they settle, even after its idle deadline. + const currentIDs = yield* execution.active + const remaining = yield* Effect.forEach(currentIDs, (sessionID) => sessions.get(sessionID)) + if (remaining.some((session) => session && key(session.location) === key(entry.ref))) return + if ((entries.get(key(entry.ref))?.expiresAt ?? 0) > now) return entries.delete(key(entry.ref)) yield* Effect.logInfo("location services evicted", { directory: entry.ref.directory, @@ -92,5 +173,5 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly export const node = makeGlobalNode({ service: Service, layer: layer(), - deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node], + deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node, Job.node], }) diff --git a/packages/core/test/browser-idle.test.ts b/packages/core/test/browser-idle.test.ts index a7db22b729d1..4165b50a929c 100644 --- a/packages/core/test/browser-idle.test.ts +++ b/packages/core/test/browser-idle.test.ts @@ -13,6 +13,7 @@ import { LocationServiceMap } from "@opencode/core/location-services" import { Plugin } from "@opencode/core/plugin" import { Rpc } from "@opencode/core/rpc" import { AbsolutePath } from "@opencode/core/schema" +import { Job } from "@opencode/core/job" import { Session } from "@opencode/core/session" import { SessionExecution } from "@opencode/core/session/execution" import { SessionStore } from "@opencode/core/session/store" @@ -31,7 +32,7 @@ const it = testEffect( makeGlobalNode({ service: LocationActivity.Service, layer: LocationActivity.layer({ timeToLive: "2 seconds", sweepInterval: "100 millis" }), - deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node], + deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node, Job.node], }), ), ], diff --git a/packages/core/test/location-activity.test.ts b/packages/core/test/location-activity.test.ts index 4ca873afbc01..8016aaca1e6a 100644 --- a/packages/core/test/location-activity.test.ts +++ b/packages/core/test/location-activity.test.ts @@ -4,9 +4,11 @@ import { TestClock } from "effect/testing" import { AppNodeBuilder } from "@opencode/core/effect/app-node-builder" import { LayerNode } from "@opencode/util/effect/layer-node" import { makeGlobalNode } from "@opencode/util/effect/app-node" +import { Permission } from "@opencode/schema/permission" import { Bus } from "@opencode/core/bus" import { Database } from "@opencode/core/database/database" import { Form } from "@opencode/core/form" +import { Job } from "@opencode/core/job" import { Location } from "@opencode/core/location" import { LocationActivity } from "@opencode/core/location-activity" import { LocationServiceMap, type LocationServices } from "@opencode/core/location-services" @@ -16,6 +18,7 @@ import { AbsolutePath } from "@opencode/core/schema" import { Session } from "@opencode/core/session" import { SessionExecution } from "@opencode/core/session/execution" import { SessionEvent } from "@opencode/core/session/event" +import { SessionMessage } from "@opencode/core/session/message" import { SessionRunner } from "@opencode/core/session/runner/index" import { SessionTable } from "@opencode/core/session/sql" import { SessionStore } from "@opencode/core/session/store" @@ -47,17 +50,26 @@ const locations = Layer.effect( const forms = yield* Form.Service return SessionRunner.Service.of({ drain: ({ sessionID }) => - forms - .ask({ - sessionID, - title: "Questions", - fields: [{ key: "runtime", type: "string" }], - }) - .pipe( - Effect.orDie, - Effect.as(SessionRunner.DrainResult.Complete()), - Effect.onInterrupt(() => Effect.sleep("5 minutes")), - ), + (sessionID.startsWith("ses_active_work") || + sessionID === Session.ID.make("ses_quiet_work") || + sessionID === Session.ID.make("ses_permission_work") + ? Effect.never + : forms + .ask({ + sessionID, + title: "Questions", + fields: [{ key: "runtime", type: "string" }], + }) + .pipe( + Effect.andThen( + sessionID === Session.ID.make("ses_answered_work") ? Effect.never : Effect.void, + ), + ) + ).pipe( + Effect.orDie, + Effect.as(SessionRunner.DrainResult.Complete()), + Effect.onInterrupt(() => Effect.sleep("5 minutes")), + ), }) }), ), @@ -85,6 +97,7 @@ const it = testEffect( Bus.node, SessionStore.node, LocationServiceMap.node, + Job.node, SessionExecution.node, LocationActivity.node, ]), @@ -102,14 +115,13 @@ const it = testEffect( describe("LocationActivity eviction", () => { for (const [count, admission] of [ - [1, "none"], [2, "none"], [1, "other"], [1, "same"], ] as const) { const newWork = admission !== "none" it.effect( - `interrupts ${count} waiting executions before eviction (${admission} session admitted during cleanup)`, + `expires ${count} waiting executions before eviction (${admission} session admitted during cleanup)`, () => Effect.gen(function* () { const db = (yield* Database.Service).db @@ -179,7 +191,6 @@ describe("LocationActivity eviction", () => { // Interruption has cancelled each question, but slow cleanup still owns the graph. expect(Array.from(yield* execution.active).toSorted()).toEqual(sessionIDs.toSorted()) expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) - expect(yield* forms.list()).toEqual([]) for (const form of pending) expect(yield* forms.state(form.id)).toEqual({ status: "cancelled" }) if (newWork) { @@ -198,16 +209,6 @@ describe("LocationActivity eviction", () => { expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual(newWork ? [ref] : []) if (newWork) { expect(yield* forms.list({ sessionID: newcomer })).toEqual([pending[count]]) - if (admission === "same") { - const later = LocationServiceMap.canonical({ directory: AbsolutePath.make("/later") }) - yield* Location.Service.pipe(Effect.provide(map.get(later)), Effect.scoped) - yield* TestClock.adjust("30 minutes") - // Keep fresh work active while a different graph reaches its own deadline. - yield* bus.publish(SessionEvent.Execution.Started, { sessionID: newcomer }, { location: ref }) - yield* TestClock.adjust("32 minutes") - expect(Array.from(yield* execution.active)).toEqual([newcomer]) - expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) - } yield* execution.interrupt(newcomer) yield* TestClock.adjust("5 minutes") yield* execution.awaitIdle(newcomer) @@ -218,4 +219,217 @@ describe("LocationActivity eviction", () => { }), ) } + + it.effect("expires unanswered and quiet sessions while a neighbor progresses", () => + Effect.gen(function* () { + const db = (yield* Database.Service).db + const bus = yield* Bus.Service + const map = yield* LocationServiceMap.Service + const execution = yield* SessionExecution.Service + const waiting = Session.ID.make("ses_waiting_question") + const quiet = Session.ID.make("ses_quiet_work") + const working = Session.ID.make("ses_active_work") + const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) + yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] }).run() + yield* db + .insert(SessionTable) + .values( + [waiting, quiet, working].map((id) => ({ + id, + project_id: Project.ID.global, + slug: "question", + directory: ref.directory, + title: "Session", + version: "test", + })), + ) + .run() + const created = yield* Deferred.make() + const interrupted: SessionEvent.Execution.Interrupted["data"][] = [] + const unsubscribe = yield* bus.listen((event) => { + if (event.type === Form.Event.Created.type) return Deferred.succeed(created, undefined).pipe(Effect.asVoid) + if (event.type === SessionEvent.Execution.Interrupted.type) + return Effect.sync(() => + interrupted.push(Schema.decodeUnknownSync(SessionEvent.Execution.Interrupted.data)(event.data)), + ) + return Effect.void + }) + yield* Effect.addFinalizer(() => unsubscribe) + yield* execution.resume(waiting).pipe(Effect.exit, Effect.forkScoped) + yield* execution.resume(quiet).pipe(Effect.exit, Effect.forkScoped) + yield* execution.resume(working).pipe(Effect.exit, Effect.forkScoped) + yield* Effect.addFinalizer(() => + Effect.forEach([waiting, quiet, working], (id) => execution.interrupt(id)).pipe( + Effect.andThen(TestClock.adjust("5 minutes")), + ), + ) + yield* Deferred.await(created) + yield* TestClock.adjust("30 minutes") + yield* bus.publish(SessionEvent.Tool.Progress, { + sessionID: working, + assistantMessageID: SessionMessage.ID.make("msg_progress"), + id: "tool-running", + metadata: { status: "still working" }, + }) + yield* bus.publish(SessionEvent.Text.Delta, { + sessionID: waiting, + assistantMessageID: SessionMessage.ID.make("msg_waiting"), + ordinal: 0, + delta: "a parallel tool is still working", + }) + yield* TestClock.adjust("33 minutes") + yield* TestClock.adjust("5 minutes") + expect(Array.from(yield* execution.active)).toEqual([working]) + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) + expect(interrupted.toSorted((a, b) => a.sessionID.localeCompare(b.sessionID))).toEqual([ + { sessionID: quiet, reason: "inactivity" }, + { sessionID: waiting, reason: "inactivity" }, + ]) + }), + ) + + it.effect("starts a permission response window when permission is asked", () => + Effect.gen(function* () { + const db = (yield* Database.Service).db + const bus = yield* Bus.Service + const execution = yield* SessionExecution.Service + const sessionID = Session.ID.make("ses_permission_work") + const directory = AbsolutePath.make("/project") + yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: directory, sandboxes: [] }).run() + yield* db + .insert(SessionTable) + .values({ + id: sessionID, + project_id: Project.ID.global, + slug: "permission", + directory, + title: "Session", + version: "test", + }) + .run() + yield* execution.resume(sessionID).pipe(Effect.exit, Effect.forkScoped) + yield* Effect.addFinalizer(() => + execution.interrupt(sessionID).pipe(Effect.andThen(TestClock.adjust("5 minutes"))), + ) + yield* TestClock.adjust("30 minutes") + yield* bus.publish(Permission.Event.Asked, { + id: Permission.ID.create("per_waiting"), + sessionID, + action: "read", + resources: ["file"], + }) + yield* TestClock.adjust("32 minutes") + expect((yield* execution.active).has(sessionID)).toBe(true) + yield* TestClock.adjust("34 minutes") + expect((yield* execution.active).has(sessionID)).toBe(false) + }), + ) + + it.effect("gives an answered request a fresh inactivity window", () => + Effect.gen(function* () { + const db = (yield* Database.Service).db + const bus = yield* Bus.Service + const map = yield* LocationServiceMap.Service + const execution = yield* SessionExecution.Service + const sessionID = Session.ID.make("ses_answered_work") + const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) + yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] }).run() + yield* db + .insert(SessionTable) + .values({ + id: sessionID, + project_id: Project.ID.global, + slug: "answer", + directory: ref.directory, + title: "Session", + version: "test", + }) + .run() + const created = yield* Deferred.make() + const unsubscribe = yield* bus.listen((event) => + event.type === Form.Event.Created.type + ? Deferred.succeed(created, Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form) + : Effect.void, + ) + yield* Effect.addFinalizer(() => unsubscribe) + yield* execution.resume(sessionID).pipe(Effect.exit, Effect.forkScoped) + yield* Effect.addFinalizer(() => + execution.interrupt(sessionID).pipe(Effect.andThen(TestClock.adjust("5 minutes"))), + ) + const request = yield* Deferred.await(created) + yield* TestClock.adjust("59 minutes") + const context = yield* map.contextEffect(ref).pipe(Effect.scoped) + yield* Context.get(context, Form.Service).reply({ id: request.id, answer: { runtime: "yes" } }) + yield* TestClock.adjust("2 minutes") + expect((yield* execution.active).has(sessionID)).toBe(true) + yield* TestClock.adjust("59 minutes") + yield* TestClock.adjust("5 minutes") + expect((yield* execution.active).has(sessionID)).toBe(false) + }), + ) + + it.effect("counts foreground child progress for its parent, but not background progress", () => + Effect.gen(function* () { + const db = (yield* Database.Service).db + const bus = yield* Bus.Service + const execution = yield* SessionExecution.Service + const jobs = yield* Job.Service + const parent = Session.ID.make("ses_quiet_work") + const child = Session.ID.make("ses_active_work_child") + const directory = AbsolutePath.make("/project") + yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: directory, sandboxes: [] }).run() + yield* db + .insert(SessionTable) + .values([ + { + id: parent, + project_id: Project.ID.global, + slug: "parent", + directory, + title: "Parent", + version: "test", + }, + { + id: child, + parent_id: parent, + project_id: Project.ID.global, + slug: "child", + directory, + title: "Child", + version: "test", + }, + ]) + .run() + yield* execution.resume(parent).pipe(Effect.exit, Effect.forkScoped) + yield* execution.resume(child).pipe(Effect.exit, Effect.forkScoped) + yield* jobs.start({ id: child, type: "subagent", run: Effect.never }) + yield* jobs.block({ id: child, sessionID: parent }).pipe(Effect.forkScoped) + yield* Effect.addFinalizer(() => + Effect.forEach([parent, child], (id) => execution.interrupt(id)).pipe( + Effect.andThen(TestClock.adjust("5 minutes")), + ), + ) + yield* TestClock.adjust("30 minutes") + yield* bus.publish(SessionEvent.Text.Delta, { + sessionID: child, + assistantMessageID: SessionMessage.ID.make("msg_child"), + ordinal: 0, + delta: "still working", + }) + yield* TestClock.adjust("38 minutes") + expect((yield* execution.active).has(parent)).toBe(true) + expect((yield* execution.active).has(child)).toBe(true) + yield* jobs.background(child) + yield* TestClock.adjust("7 minutes") + yield* bus.publish(SessionEvent.Text.Delta, { + sessionID: child, + assistantMessageID: SessionMessage.ID.make("msg_child"), + ordinal: 0, + delta: "still working", + }) + yield* TestClock.adjust("25 minutes") + expect((yield* execution.active).has(parent)).toBe(false) + expect((yield* execution.active).has(child)).toBe(true) + }), + ) }) From 21b930bfc0c7f6d7273366eb9e31d82ffe77e90d Mon Sep 17 00:00:00 2001 From: Khoa Huynh Date: Thu, 1 Oct 2026 17:20:12 -0400 Subject: [PATCH 2/3] fix(core): share location activity with pending requests --- packages/core/src/location-activity.ts | 56 +++--- packages/core/test/location-activity.test.ts | 198 +++++++------------ 2 files changed, 102 insertions(+), 152 deletions(-) diff --git a/packages/core/src/location-activity.ts b/packages/core/src/location-activity.ts index 2e629cad27eb..7df11d58ab16 100644 --- a/packages/core/src/location-activity.ts +++ b/packages/core/src/location-activity.ts @@ -75,20 +75,18 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly if (Schema.is(SessionSchema.ID)(sessionID)) clearWait(sessionID, event.data.id) return } - if (!isSessionEvent(event)) return + if (!isSessionEvent(event) || event.type === SessionEvent.Viewed.type) return const sessionID = event.data.sessionID - if (event.type !== SessionEvent.Viewed.type) { - if (event.type === SessionEvent.Execution.Started.type || progress.has(sessionID)) - progress.set(sessionID, clock.currentTimeMillisUnsafe() + timeToLive) - // Parentage alone also includes background jobs; only a blocking chain carries progress. - for ( - let child = sessionID, parent = parents.get(child); - parent; - child = parent, parent = parents.get(child) - ) { - if (!progress.has(parent) || !(yield* jobs.isBlocking({ id: child, sessionID: parent }))) break - progress.set(parent, clock.currentTimeMillisUnsafe() + timeToLive) - } + if (event.type === SessionEvent.Execution.Started.type || progress.has(sessionID)) + progress.set(sessionID, clock.currentTimeMillisUnsafe() + timeToLive) + // Parentage alone also includes background jobs; only a blocking chain carries progress. + for ( + let child = sessionID, parent = parents.get(child); + parent; + child = parent, parent = parents.get(child) + ) { + if (!progress.has(parent) || !(yield* jobs.isBlocking({ id: child, sessionID: parent }))) break + progress.set(parent, clock.currentTimeMillisUnsafe() + timeToLive) } if ( event.type === SessionEvent.Execution.Succeeded.type || @@ -99,8 +97,9 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly parents.delete(sessionID) waits.delete(sessionID) } - if (!event.durable) return - const location = event.location + // Automatic cleanup must not renew other pending requests in this location. + if (event.type === SessionEvent.Execution.Interrupted.type && event.data.reason === "inactivity") return + const location = event.location ?? (yield* sessions.get(sessionID))?.location if (location && (yield* RcMap.has(locations.rcMap, location))) yield* touch(location) }), ) @@ -132,15 +131,21 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly ) yield* Effect.forEach( activeIDs, - (sessionID) => { - if (!progress.has(sessionID)) progress.set(sessionID, now + timeToLive) - const waiting = Array.from(waits.get(sessionID)?.values() ?? []).reduce( - (earliest, wait) => (earliest === undefined || wait < earliest ? wait : earliest), - undefined, - ) - if ((waiting ?? progress.get(sessionID) ?? 0) > now) return Effect.void - return execution.interrupt(sessionID, { reason: "inactivity" }).pipe(Effect.asVoid) - }, + (sessionID) => + Effect.gen(function* () { + if (!progress.has(sessionID)) progress.set(sessionID, now + timeToLive) + const session = waits.has(sessionID) ? yield* sessions.get(sessionID) : undefined + const waiting = waits.get(sessionID) + // Pending input follows location inactivity; other stalled work keeps its own deadline. + const deadline = waiting + ? Math.max( + Math.min(...waiting.values()), + session ? (entries.get(key(session.location))?.expiresAt ?? 0) : 0, + ) + : (progress.get(sessionID) ?? 0) + if (deadline > now) return + yield* execution.interrupt(sessionID, { reason: "inactivity" }) + }), { discard: true, concurrency: "unbounded" }, ) const expired = Array.from(entries.values()).filter((entry) => entry.expiresAt <= now) @@ -151,8 +156,7 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly Effect.gen(function* () { // Invalidation detaches the graph while borrowers still hold it. Active // executions retain their Location until they settle, even after its idle deadline. - const currentIDs = yield* execution.active - const remaining = yield* Effect.forEach(currentIDs, (sessionID) => sessions.get(sessionID)) + const remaining = yield* Effect.forEach(yield* execution.active, (sessionID) => sessions.get(sessionID)) if (remaining.some((session) => session && key(session.location) === key(entry.ref))) return if ((entries.get(key(entry.ref))?.expiresAt ?? 0) > now) return entries.delete(key(entry.ref)) diff --git a/packages/core/test/location-activity.test.ts b/packages/core/test/location-activity.test.ts index 8016aaca1e6a..39b94c786683 100644 --- a/packages/core/test/location-activity.test.ts +++ b/packages/core/test/location-activity.test.ts @@ -54,17 +54,11 @@ const locations = Layer.effect( sessionID === Session.ID.make("ses_quiet_work") || sessionID === Session.ID.make("ses_permission_work") ? Effect.never - : forms - .ask({ - sessionID, - title: "Questions", - fields: [{ key: "runtime", type: "string" }], - }) - .pipe( - Effect.andThen( - sessionID === Session.ID.make("ses_answered_work") ? Effect.never : Effect.void, - ), - ) + : forms.ask({ + sessionID, + title: "Questions", + fields: [{ key: "runtime", type: "string" }], + }) ).pipe( Effect.orDie, Effect.as(SessionRunner.DrainResult.Complete()), @@ -114,12 +108,7 @@ const it = testEffect( ) describe("LocationActivity eviction", () => { - for (const [count, admission] of [ - [2, "none"], - [1, "other"], - [1, "same"], - ] as const) { - const newWork = admission !== "none" + for (const [count, admission] of [[1, "same"]] as const) { it.effect( `expires ${count} waiting executions before eviction (${admission} session admitted during cleanup)`, () => @@ -132,7 +121,7 @@ describe("LocationActivity eviction", () => { const sessionIDs = Array.from({ length: count }, (_, index) => Session.ID.make(`ses_waiting_question_${index}`), ) - const newcomer = admission === "same" ? sessionIDs[0] : Session.ID.make("ses_new_question") + const newcomer = sessionIDs[0] const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) const idle = Location.Ref.make({ directory: ref.directory, workspaceID: Workspace.ID.make("wrk_idle") }) yield* db @@ -193,61 +182,63 @@ describe("LocationActivity eviction", () => { expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) for (const form of pending) expect(yield* forms.state(form.id)).toEqual({ status: "cancelled" }) - if (newWork) { - yield* execution.wake(newcomer) - if (admission === "other") yield* Deferred.await(newCreated) - } + yield* execution.wake(newcomer) yield* TestClock.adjust("5 minutes") - if (newWork) yield* Deferred.await(newCreated) + yield* Deferred.await(newCreated) const results = yield* Effect.forEach(running, Fiber.join) expect(results.every((exit) => exit._tag === "Failure")).toBe(true) - expect(Array.from(yield* execution.active)).toEqual(newWork ? [newcomer] : []) - expect(yield* store.listSuspended()).toEqual(newWork ? [newcomer] : []) + expect(Array.from(yield* execution.active)).toEqual([newcomer]) + expect(yield* store.listSuspended()).toEqual([newcomer]) expect(interrupted.toSorted((a, b) => a.sessionID.localeCompare(b.sessionID))).toEqual( sessionIDs.map((sessionID) => ({ sessionID, reason: "inactivity" })), ) - expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual(newWork ? [ref] : []) - if (newWork) { - expect(yield* forms.list({ sessionID: newcomer })).toEqual([pending[count]]) - yield* execution.interrupt(newcomer) - yield* TestClock.adjust("5 minutes") - yield* execution.awaitIdle(newcomer) - yield* TestClock.adjust("62 minutes") - expect(yield* store.listSuspended()).toEqual([]) - expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([]) - } + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) + expect(yield* forms.list({ sessionID: newcomer })).toEqual([pending[count]]) + yield* execution.interrupt(newcomer) + yield* TestClock.adjust("5 minutes") + yield* execution.awaitIdle(newcomer) + yield* TestClock.adjust("62 minutes") + expect(yield* store.listSuspended()).toEqual([]) + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([]) }), ) } - it.effect("expires unanswered and quiet sessions while a neighbor progresses", () => + it.effect("keeps pending requests alive while their location progresses, then expires only its idle sessions", () => Effect.gen(function* () { const db = (yield* Database.Service).db const bus = yield* Bus.Service const map = yield* LocationServiceMap.Service const execution = yield* SessionExecution.Service const waiting = Session.ID.make("ses_waiting_question") + const permission = Session.ID.make("ses_permission_work") const quiet = Session.ID.make("ses_quiet_work") const working = Session.ID.make("ses_active_work") + const elsewhere = Session.ID.make("ses_active_work_elsewhere") const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) + const other = Location.Ref.make({ directory: ref.directory, workspaceID: Workspace.ID.make("wrk_other") }) yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] }).run() yield* db .insert(SessionTable) .values( - [waiting, quiet, working].map((id) => ({ + [waiting, permission, quiet, working, elsewhere].map((id) => ({ id, project_id: Project.ID.global, slug: "question", directory: ref.directory, + workspace_id: id === elsewhere ? other.workspaceID : undefined, title: "Session", version: "test", })), ) .run() - const created = yield* Deferred.make() + const created = yield* Deferred.make() const interrupted: SessionEvent.Execution.Interrupted["data"][] = [] const unsubscribe = yield* bus.listen((event) => { - if (event.type === Form.Event.Created.type) return Deferred.succeed(created, undefined).pipe(Effect.asVoid) + if (event.type === Form.Event.Created.type) + return Deferred.succeed(created, Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form).pipe( + Effect.asVoid, + ) if (event.type === SessionEvent.Execution.Interrupted.type) return Effect.sync(() => interrupted.push(Schema.decodeUnknownSync(SessionEvent.Execution.Interrupted.data)(event.data)), @@ -256,14 +247,25 @@ describe("LocationActivity eviction", () => { }) yield* Effect.addFinalizer(() => unsubscribe) yield* execution.resume(waiting).pipe(Effect.exit, Effect.forkScoped) + yield* execution.resume(permission).pipe(Effect.exit, Effect.forkScoped) yield* execution.resume(quiet).pipe(Effect.exit, Effect.forkScoped) yield* execution.resume(working).pipe(Effect.exit, Effect.forkScoped) + yield* execution.resume(elsewhere).pipe(Effect.exit, Effect.forkScoped) yield* Effect.addFinalizer(() => - Effect.forEach([waiting, quiet, working], (id) => execution.interrupt(id)).pipe( + Effect.forEach([waiting, permission, quiet, working, elsewhere], (id) => execution.interrupt(id)).pipe( Effect.andThen(TestClock.adjust("5 minutes")), ), ) - yield* Deferred.await(created) + const request = yield* Deferred.await(created) + const context = yield* map.contextEffect(ref).pipe(Effect.scoped) + const forms = Context.get(context, Form.Service) + yield* TestClock.adjust("1 minute") + yield* bus.publish(Permission.Event.Asked, { + id: Permission.ID.create("per_waiting"), + sessionID: permission, + action: "read", + resources: ["file"], + }) yield* TestClock.adjust("30 minutes") yield* bus.publish(SessionEvent.Tool.Progress, { sessionID: working, @@ -272,102 +274,46 @@ describe("LocationActivity eviction", () => { metadata: { status: "still working" }, }) yield* bus.publish(SessionEvent.Text.Delta, { - sessionID: waiting, - assistantMessageID: SessionMessage.ID.make("msg_waiting"), + sessionID: elsewhere, + assistantMessageID: SessionMessage.ID.make("msg_elsewhere"), ordinal: 0, - delta: "a parallel tool is still working", + delta: "another workspace is still working", }) yield* TestClock.adjust("33 minutes") yield* TestClock.adjust("5 minutes") - expect(Array.from(yield* execution.active)).toEqual([working]) - expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) + expect(Array.from(yield* execution.active).toSorted()).toEqual( + [waiting, permission, working, elsewhere].toSorted(), + ) + expect(yield* forms.state(request.id)).toEqual({ status: "pending" }) + expect(interrupted).toEqual([{ sessionID: quiet, reason: "inactivity" }]) + + yield* execution.interrupt(working) + yield* TestClock.adjust("5 minutes") + yield* execution.awaitIdle(working) + yield* TestClock.adjust("10 minutes") + // Neither activity in another workspace nor viewing the question renews this location. + yield* bus.publish(SessionEvent.Text.Delta, { + sessionID: elsewhere, + assistantMessageID: SessionMessage.ID.make("msg_elsewhere"), + ordinal: 0, + delta: "still working in the other workspace", + }) + yield* bus.publish(SessionEvent.Viewed, { sessionID: waiting, idle: 0 }, { location: ref }) + expect(yield* forms.state(request.id)).toEqual({ status: "pending" }) + yield* TestClock.adjust("53 minutes") + yield* TestClock.adjust("5 minutes") + expect(Array.from(yield* execution.active)).toEqual([elsewhere]) + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([other]) + expect(yield* forms.state(request.id)).toEqual({ status: "cancelled" }) expect(interrupted.toSorted((a, b) => a.sessionID.localeCompare(b.sessionID))).toEqual([ + { sessionID: working, reason: "user" }, + { sessionID: permission, reason: "inactivity" }, { sessionID: quiet, reason: "inactivity" }, { sessionID: waiting, reason: "inactivity" }, ]) }), ) - it.effect("starts a permission response window when permission is asked", () => - Effect.gen(function* () { - const db = (yield* Database.Service).db - const bus = yield* Bus.Service - const execution = yield* SessionExecution.Service - const sessionID = Session.ID.make("ses_permission_work") - const directory = AbsolutePath.make("/project") - yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: directory, sandboxes: [] }).run() - yield* db - .insert(SessionTable) - .values({ - id: sessionID, - project_id: Project.ID.global, - slug: "permission", - directory, - title: "Session", - version: "test", - }) - .run() - yield* execution.resume(sessionID).pipe(Effect.exit, Effect.forkScoped) - yield* Effect.addFinalizer(() => - execution.interrupt(sessionID).pipe(Effect.andThen(TestClock.adjust("5 minutes"))), - ) - yield* TestClock.adjust("30 minutes") - yield* bus.publish(Permission.Event.Asked, { - id: Permission.ID.create("per_waiting"), - sessionID, - action: "read", - resources: ["file"], - }) - yield* TestClock.adjust("32 minutes") - expect((yield* execution.active).has(sessionID)).toBe(true) - yield* TestClock.adjust("34 minutes") - expect((yield* execution.active).has(sessionID)).toBe(false) - }), - ) - - it.effect("gives an answered request a fresh inactivity window", () => - Effect.gen(function* () { - const db = (yield* Database.Service).db - const bus = yield* Bus.Service - const map = yield* LocationServiceMap.Service - const execution = yield* SessionExecution.Service - const sessionID = Session.ID.make("ses_answered_work") - const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) - yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] }).run() - yield* db - .insert(SessionTable) - .values({ - id: sessionID, - project_id: Project.ID.global, - slug: "answer", - directory: ref.directory, - title: "Session", - version: "test", - }) - .run() - const created = yield* Deferred.make() - const unsubscribe = yield* bus.listen((event) => - event.type === Form.Event.Created.type - ? Deferred.succeed(created, Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form) - : Effect.void, - ) - yield* Effect.addFinalizer(() => unsubscribe) - yield* execution.resume(sessionID).pipe(Effect.exit, Effect.forkScoped) - yield* Effect.addFinalizer(() => - execution.interrupt(sessionID).pipe(Effect.andThen(TestClock.adjust("5 minutes"))), - ) - const request = yield* Deferred.await(created) - yield* TestClock.adjust("59 minutes") - const context = yield* map.contextEffect(ref).pipe(Effect.scoped) - yield* Context.get(context, Form.Service).reply({ id: request.id, answer: { runtime: "yes" } }) - yield* TestClock.adjust("2 minutes") - expect((yield* execution.active).has(sessionID)).toBe(true) - yield* TestClock.adjust("59 minutes") - yield* TestClock.adjust("5 minutes") - expect((yield* execution.active).has(sessionID)).toBe(false) - }), - ) - it.effect("counts foreground child progress for its parent, but not background progress", () => Effect.gen(function* () { const db = (yield* Database.Service).db From 5d689043b01fe139c94824aee25f917bb3950a31 Mon Sep 17 00:00:00 2001 From: Khoa Huynh Date: Fri, 2 Oct 2026 14:51:49 -0400 Subject: [PATCH 3/3] fix(core): tie inactivity to session-owned work --- packages/core/src/inactivity.ts | 92 +++++ packages/core/src/job.ts | 20 +- packages/core/src/location-activity.ts | 192 +++------- packages/core/src/shell.ts | 11 +- packages/core/test/browser-idle.test.ts | 4 +- packages/core/test/job.test.ts | 3 +- packages/core/test/location-activity.test.ts | 378 ++++++++----------- packages/core/test/session-execution.test.ts | 12 +- packages/core/test/session-shell.test.ts | 75 +++- 9 files changed, 405 insertions(+), 382 deletions(-) create mode 100644 packages/core/src/inactivity.ts diff --git a/packages/core/src/inactivity.ts b/packages/core/src/inactivity.ts new file mode 100644 index 000000000000..d48572e7a48d --- /dev/null +++ b/packages/core/src/inactivity.ts @@ -0,0 +1,92 @@ +export * as Inactivity from "./inactivity.js" + +import { Clock, Context, Effect, Layer, Scope } from "effect" +import { makeGlobalNode } from "@opencode/util/effect/app-node" +import type { Location } from "@opencode/schema/location" +import type { SessionSchema } from "./session/schema.js" + +type Owner = + | { readonly sessionID: SessionSchema.ID; readonly location?: Location.Ref } + | { readonly sessionID?: undefined; readonly location: Location.Ref } +type Activity = { since: number; working: number } + +/** One process-local inactivity policy for Sessions and their shared Location resources. */ +export class Service extends Context.Service< + Service, + { + /** Record progress or human input, without borrowing a Location graph. */ + readonly touch: (owner: Owner) => Effect.Effect + /** Owned work prevents expiration until its scope closes, then starts a fresh idle window. */ + readonly hold: (owner: Owner) => Effect.Effect + /** Find expired live executions and cached graphs; prune observations that no longer have owners. */ + readonly expired: (input: { + readonly sessions: ReadonlySet + readonly locations: readonly Location.Ref[] + readonly timeToLive: number + }) => Effect.Effect<{ sessions: SessionSchema.ID[]; locations: Location.Ref[] }> + } +>()("@opencode/Inactivity") {} + +const layer = Layer.effect( + Service, + Effect.gen(function* () { + const clock = yield* Clock.Clock + const sessions = new Map() + const locations = new Map() + const key = (ref: Location.Ref) => `${ref.directory}\0${ref.workspaceID ?? ""}` + const entries = (owner: Owner) => { + const session = owner.sessionID + ? (sessions.get(owner.sessionID) ?? { since: clock.currentTimeMillisUnsafe(), working: 0 }) + : undefined + const location = owner.location + ? (locations.get(key(owner.location)) ?? { since: clock.currentTimeMillisUnsafe(), working: 0 }) + : undefined + if (session && owner.sessionID) sessions.set(owner.sessionID, session) + if (location && owner.location) locations.set(key(owner.location), location) + return [session, location].flatMap((entry) => (entry ? [entry] : [])) + } + return Service.of({ + touch: (owner) => + Effect.sync(() => entries(owner).forEach((entry) => (entry.since = clock.currentTimeMillisUnsafe()))), + hold: (owner) => + Effect.acquireRelease( + Effect.sync(() => + entries(owner).map((entry) => { + entry.working++ + return entry + }), + ), + (held) => + Effect.sync(() => + held.forEach((entry) => { + entry.working-- + entry.since = clock.currentTimeMillisUnsafe() + }), + ), + ).pipe(Effect.asVoid), + expired: (input) => + Effect.sync(() => { + input.sessions.forEach((sessionID) => entries({ sessionID })) + input.locations.forEach((location) => entries({ location })) + sessions.forEach((entry, id) => { + if (!input.sessions.has(id) && entry.working === 0) sessions.delete(id) + }) + const cached = new Set(input.locations.map(key)) + locations.forEach((entry, id) => { + if (!cached.has(id) && entry.working === 0) locations.delete(id) + }) + const idle = (entry: Activity | undefined) => + entry !== undefined && + entry.working === 0 && + entry.since + input.timeToLive <= clock.currentTimeMillisUnsafe() + return { + sessions: Array.from(input.sessions).filter((id) => idle(sessions.get(id))), + locations: input.locations.filter((ref) => idle(locations.get(key(ref)))), + } + }), + }) + }), +) + +/** Shared lifetime accounting; never recovered from durable "running" markers. */ +export const node = makeGlobalNode({ service: Service, layer, deps: [] }) diff --git a/packages/core/src/job.ts b/packages/core/src/job.ts index 7bee22493b08..5bad23106eea 100644 --- a/packages/core/src/job.ts +++ b/packages/core/src/job.ts @@ -5,6 +5,7 @@ import { makeGlobalNode } from "@opencode/util/effect/app-node" import { KV } from "./kv.js" import { SessionMessage } from "./session/message.js" import { SessionSchema } from "./session/schema.js" +import { Inactivity } from "./inactivity.js" const Background = Schema.Struct({ id: Schema.String, @@ -129,8 +130,6 @@ export interface Interface { readonly start: (input: StartInput) => Effect.Effect readonly wait: (input: WaitInput) => Effect.Effect readonly block: (input: BlockInput) => Effect.Effect - /** Whether a running job is currently blocking this Session. Background jobs do not count. */ - readonly isBlocking: (input: BlockInput) => Effect.Effect readonly background: (id: string) => Effect.Effect readonly backgroundAll: (input: BackgroundAllInput) => Effect.Effect readonly cancel: (id: string) => Effect.Effect @@ -179,6 +178,7 @@ function decrementSession(input: Map, sessionID: Sessi */ export const make = Effect.gen(function* () { const kv = yield* KV.Service + const inactivity = yield* Inactivity.Service const state: State = { jobs: yield* SynchronizedRef.make(new Map()), scope: yield* Scope.Scope, @@ -287,7 +287,13 @@ export const make = Effect.gen(function* () { }), ) if ("scope" in result) - yield* restore(input.run).pipe( + yield* restore( + Effect.gen(function* () { + if (input.recovery?.kind === "subagent") + yield* inactivity.hold({ sessionID: input.recovery.parentSessionID }) + return yield* input.run + }).pipe(Effect.scoped), + ).pipe( Effect.exit, Effect.flatMap((exit) => settle(input.id, result.scope, exit)), Effect.asVoid, @@ -356,11 +362,6 @@ export const make = Effect.gen(function* () { ) }) - const isBlocking: Interface["isBlocking"] = Effect.fnUntraced(function* (input) { - const job = (yield* SynchronizedRef.get(state.jobs)).get(input.id) - return job?.blockingSessions.has(input.sessionID) ?? false - }) - const markBackground = Effect.fnUntraced(function* (job: Active) { const next = { ...job, @@ -474,7 +475,6 @@ export const make = Effect.gen(function* () { start, wait, block, - isBlocking, background, backgroundAll, cancel, @@ -485,4 +485,4 @@ export const make = Effect.gen(function* () { const layer = Layer.effect(Service, make) -export const node = makeGlobalNode({ service: Service, layer, deps: [KV.node] }) +export const node = makeGlobalNode({ service: Service, layer, deps: [KV.node, Inactivity.node] }) diff --git a/packages/core/src/location-activity.ts b/packages/core/src/location-activity.ts index 7df11d58ab16..90fde6b27f4c 100644 --- a/packages/core/src/location-activity.ts +++ b/packages/core/src/location-activity.ts @@ -1,12 +1,11 @@ export * as LocationActivity from "./location-activity.js" -import { Clock, Context, Duration, Effect, Layer, RcMap, Schema } from "effect" +import { Context, Duration, Effect, Layer, RcMap, Schema } from "effect" import { Event } from "@opencode/schema/event" import { Permission } from "@opencode/schema/permission" import { Form } from "@opencode/schema/form" import { Bus } from "./bus.js" -import { Job } from "./job.js" -import { Location } from "./location.js" +import { Inactivity } from "./inactivity.js" import { LocationServiceMap } from "./location-service-map.js" import { SessionEvent } from "./session/event.js" import { SessionExecution } from "./session/execution.js" @@ -17,165 +16,92 @@ import { makeGlobalNode } from "@opencode/util/effect/app-node" const isSessionEvent = (event: Event.Payload): event is SessionEvent.Event => Object.hasOwn(SessionEvent.All.cases, event.type) +/** Observe Session progress and expire process-local execution separately from shared Location resources. */ export class Service extends Context.Service()("@opencode/LocationActivity") {} +/** Run the inactivity sweep; shorter intervals are also used by isolated lifecycle fixtures. */ export function layer(options: { readonly timeToLive?: Duration.Input; readonly sweepInterval?: Duration.Input } = {}) { return Layer.effect( Service, Effect.gen(function* () { - const clock = yield* Clock.Clock const bus = yield* Bus.Service const locations = yield* LocationServiceMap.Service const execution = yield* SessionExecution.Service - const jobs = yield* Job.Service const sessions = yield* SessionStore.Service + const inactivity = yield* Inactivity.Service const timeToLive = Duration.toMillis(options.timeToLive ?? "60 minutes") - const entries = new Map() - const progress = new Map() - const parents = new Map() - const waits = new Map>() - const trackWait = (sessionID: SessionSchema.ID, id: string) => { - if (!progress.has(sessionID)) return - const pending = waits.get(sessionID) ?? new Map() - pending.set(id, clock.currentTimeMillisUnsafe() + timeToLive) - waits.set(sessionID, pending) - } - const clearWait = (sessionID: SessionSchema.ID, id: string) => { - const pending = waits.get(sessionID) - pending?.delete(id) - if (pending?.size === 0) waits.delete(sessionID) - if (progress.has(sessionID)) progress.set(sessionID, clock.currentTimeMillisUnsafe() + timeToLive) - } - const key = (ref: Location.Ref) => `${LocationServiceMap.canonical(ref).directory}\0${ref.workspaceID ?? ""}` - const touch = (ref: Location.Ref) => - Effect.sync(() => { - entries.set(key(ref), { ref, expiresAt: clock.currentTimeMillisUnsafe() + timeToLive }) - }) - const unsubscribe = yield* bus.listen((event) => Effect.gen(function* () { - if (event.type === Permission.Event.Asked.type && Schema.is(Permission.Event.Asked)(event)) { - trackWait(event.data.sessionID, event.data.id) - return - } - if (event.type === Permission.Event.Replied.type && Schema.is(Permission.Event.Replied)(event)) { - clearWait(event.data.sessionID, event.data.requestID) - return - } - if (event.type === Form.Event.Created.type && Schema.is(Form.Event.Created)(event)) { - const sessionID = event.data.form.sessionID - if (Schema.is(SessionSchema.ID)(sessionID)) trackWait(sessionID, event.data.form.id) - return - } + if (event.type === SessionEvent.Viewed.type) return if ( - (event.type === Form.Event.Replied.type && Schema.is(Form.Event.Replied)(event)) || - (event.type === Form.Event.Cancelled.type && Schema.is(Form.Event.Cancelled)(event)) - ) { - const sessionID = event.data.sessionID - if (Schema.is(SessionSchema.ID)(sessionID)) clearWait(sessionID, event.data.id) + isSessionEvent(event) && + event.type === SessionEvent.Execution.Interrupted.type && + event.data.reason === "inactivity" + ) return - } - if (!isSessionEvent(event) || event.type === SessionEvent.Viewed.type) return - const sessionID = event.data.sessionID - if (event.type === SessionEvent.Execution.Started.type || progress.has(sessionID)) - progress.set(sessionID, clock.currentTimeMillisUnsafe() + timeToLive) - // Parentage alone also includes background jobs; only a blocking chain carries progress. - for ( - let child = sessionID, parent = parents.get(child); - parent; - child = parent, parent = parents.get(child) - ) { - if (!progress.has(parent) || !(yield* jobs.isBlocking({ id: child, sessionID: parent }))) break - progress.set(parent, clock.currentTimeMillisUnsafe() + timeToLive) - } - if ( - event.type === SessionEvent.Execution.Succeeded.type || - event.type === SessionEvent.Execution.Failed.type || - event.type === SessionEvent.Execution.Interrupted.type - ) { - progress.delete(sessionID) - parents.delete(sessionID) - waits.delete(sessionID) - } - // Automatic cleanup must not renew other pending requests in this location. - if (event.type === SessionEvent.Execution.Interrupted.type && event.data.reason === "inactivity") return + const sessionID = isSessionEvent(event) + ? event.data.sessionID + : Schema.is(Permission.Event.Asked)(event) + ? event.data.sessionID + : Schema.is(Permission.Event.Replied)(event) + ? event.data.sessionID + : Schema.is(Form.Event.Created)(event) + ? event.data.form.sessionID + : Schema.is(Form.Event.Replied)(event) || Schema.is(Form.Event.Cancelled)(event) + ? event.data.sessionID + : undefined + if (!Schema.is(SessionSchema.ID)(sessionID)) return const location = event.location ?? (yield* sessions.get(sessionID))?.location - if (location && (yield* RcMap.has(locations.rcMap, location))) yield* touch(location) + yield* inactivity.touch({ sessionID, ...(location ? { location } : {}) }) }), ) yield* Effect.addFinalizer(() => unsubscribe) yield* Effect.gen(function* () { yield* Effect.sleep(options.sweepInterval ?? "1 minute") - const refs = Array.from(yield* RcMap.keys(locations.rcMap)) - const cached = new Set(refs.map(key)) - yield* Effect.forEach(refs, (ref) => (entries.has(key(ref)) ? Effect.void : touch(ref)), { discard: true }) - for (const id of entries.keys()) { - if (!cached.has(id)) entries.delete(id) + const expired = yield* inactivity.expired({ + sessions: yield* execution.active, + locations: Array.from(yield* RcMap.keys(locations.rcMap)), + timeToLive, + }) + yield* Effect.forEach(expired.sessions, (id) => execution.interrupt(id, { reason: "inactivity" }), { + discard: true, + concurrency: "unbounded", + }) + // Expiring one execution never shuts down a graph still owned by another execution or process. + for (const ref of expired.locations) { + const remaining = yield* Effect.forEach(yield* execution.active, (id) => sessions.get(id)) + if ( + remaining.some( + (session) => + session && + session.location.directory === ref.directory && + session.location.workspaceID === ref.workspaceID, + ) + ) + continue + const current = yield* inactivity.expired({ + sessions: yield* execution.active, + locations: Array.from(yield* RcMap.keys(locations.rcMap)), + timeToLive, + }) + if ( + !current.locations.some( + (location) => location.directory === ref.directory && location.workspaceID === ref.workspaceID, + ) + ) + continue + yield* Effect.logInfo("location services evicted", { directory: ref.directory, workspaceID: ref.workspaceID }) + yield* locations.invalidate(ref) } - const now = clock.currentTimeMillisUnsafe() - const activeIDs = yield* execution.active - yield* Effect.forEach( - Array.from(activeIDs).filter((id) => !parents.has(id)), - (id) => - Effect.gen(function* () { - const session = yield* sessions.get(id) - parents.set(id, session?.parentID ?? null) - if ( - session?.parentID && - progress.has(session.parentID) && - (yield* jobs.isBlocking({ id, sessionID: session.parentID })) - ) - progress.set(session.parentID, now + timeToLive) - }), - { discard: true }, - ) - yield* Effect.forEach( - activeIDs, - (sessionID) => - Effect.gen(function* () { - if (!progress.has(sessionID)) progress.set(sessionID, now + timeToLive) - const session = waits.has(sessionID) ? yield* sessions.get(sessionID) : undefined - const waiting = waits.get(sessionID) - // Pending input follows location inactivity; other stalled work keeps its own deadline. - const deadline = waiting - ? Math.max( - Math.min(...waiting.values()), - session ? (entries.get(key(session.location))?.expiresAt ?? 0) : 0, - ) - : (progress.get(sessionID) ?? 0) - if (deadline > now) return - yield* execution.interrupt(sessionID, { reason: "inactivity" }) - }), - { discard: true, concurrency: "unbounded" }, - ) - const expired = Array.from(entries.values()).filter((entry) => entry.expiresAt <= now) - if (expired.length === 0) return - yield* Effect.forEach( - expired, - (entry) => - Effect.gen(function* () { - // Invalidation detaches the graph while borrowers still hold it. Active - // executions retain their Location until they settle, even after its idle deadline. - const remaining = yield* Effect.forEach(yield* execution.active, (sessionID) => sessions.get(sessionID)) - if (remaining.some((session) => session && key(session.location) === key(entry.ref))) return - if ((entries.get(key(entry.ref))?.expiresAt ?? 0) > now) return - entries.delete(key(entry.ref)) - yield* Effect.logInfo("location services evicted", { - directory: entry.ref.directory, - workspaceID: entry.ref.workspaceID, - }).pipe(Effect.andThen(locations.invalidate(entry.ref))) - }), - { discard: true, concurrency: "unbounded" }, - ) }).pipe(Effect.forever, Effect.forkScoped) - return Service.of({}) }), ) } +/** Process-global observation and cleanup, using Session-owned work accounting. */ export const node = makeGlobalNode({ service: Service, layer: layer(), - deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node, Job.node], + deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node, Inactivity.node], }) diff --git a/packages/core/src/shell.ts b/packages/core/src/shell.ts index 393275517141..7c339451f709 100644 --- a/packages/core/src/shell.ts +++ b/packages/core/src/shell.ts @@ -10,6 +10,7 @@ import { CrossSpawnSpawner } from "@opencode/util/cross-spawn-spawner" import { makeGlobalNode, makeLocationNode } from "@opencode/util/effect/app-node" import { FSUtil } from "@opencode/util/fs-util" import { Bus } from "./bus.js" +import { Inactivity } from "./inactivity.js" import { Environment } from "./environment/index.js" import { FileRetention } from "./file-retention.js" import { Location } from "./location.js" @@ -121,6 +122,7 @@ const layer = () => Service, Effect.gen(function* () { const bus = yield* Bus.Service + const inactivity = yield* Inactivity.Service const location = yield* Location.Service const global = yield* Global.Service const shell = yield* ShellSelect.Service @@ -255,11 +257,9 @@ const layer = () => input: CreateInput, before?: (input: ShellCreateBefore) => Effect.Effect, ) { - const sessionID = input.metadata?.sessionID + const sessionID = Schema.is(SessionSchema.ID)(input.metadata?.sessionID) ? input.metadata.sessionID : undefined const sessionEnvironment = - location.workspaceID === undefined && Schema.is(SessionSchema.ID)(sessionID) - ? yield* environments.get(sessionID) - : undefined + location.workspaceID === undefined && sessionID !== undefined ? yield* environments.get(sessionID) : undefined const invocation: ShellCreateBefore = { command: input.command, cwd: input.cwd ?? location.directory, @@ -296,6 +296,8 @@ const layer = () => runFork( Effect.scoped( Effect.gen(function* () { + // Registered before spawn so ownership outlasts the handle's cleanup, even on failure. + yield* inactivity.hold({ location, sessionID }) const handle = yield* environment.spawner .spawn( ChildProcess.make(invocation.shell, args, { @@ -443,6 +445,7 @@ export const node = makeLocationNode({ layer: layer(), deps: [ Bus.node, + Inactivity.node, Location.node, Global.node, ShellSelect.node, diff --git a/packages/core/test/browser-idle.test.ts b/packages/core/test/browser-idle.test.ts index 4165b50a929c..4e8aaf00271f 100644 --- a/packages/core/test/browser-idle.test.ts +++ b/packages/core/test/browser-idle.test.ts @@ -13,7 +13,7 @@ import { LocationServiceMap } from "@opencode/core/location-services" import { Plugin } from "@opencode/core/plugin" import { Rpc } from "@opencode/core/rpc" import { AbsolutePath } from "@opencode/core/schema" -import { Job } from "@opencode/core/job" +import { Inactivity } from "@opencode/core/inactivity" import { Session } from "@opencode/core/session" import { SessionExecution } from "@opencode/core/session/execution" import { SessionStore } from "@opencode/core/session/store" @@ -32,7 +32,7 @@ const it = testEffect( makeGlobalNode({ service: LocationActivity.Service, layer: LocationActivity.layer({ timeToLive: "2 seconds", sweepInterval: "100 millis" }), - deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node, Job.node], + deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node, Inactivity.node], }), ), ], diff --git a/packages/core/test/job.test.ts b/packages/core/test/job.test.ts index 400a2adb4aaa..23deb3ed7f33 100644 --- a/packages/core/test/job.test.ts +++ b/packages/core/test/job.test.ts @@ -1,5 +1,6 @@ import { describe, expect } from "bun:test" import { Job } from "@opencode/core/job" +import { Inactivity } from "@opencode/core/inactivity" import { KV } from "@opencode/core/kv" import { AppNodeBuilder } from "@opencode/core/effect/app-node-builder" import { Integration } from "@opencode/core/integration" @@ -8,7 +9,7 @@ import { Cause, Deferred, Effect, Exit, Fiber, Scope } from "effect" import { SessionSchema } from "@opencode/core/session/schema" import { testEffect } from "./lib/effect" -const it = testEffect(AppNodeBuilder.build(LayerNode.group([Job.node, KV.node]))) +const it = testEffect(AppNodeBuilder.build(LayerNode.group([Job.node, KV.node, Inactivity.node]))) describe("Job", () => { it.live("tracks process-local work through explicit observation", () => diff --git a/packages/core/test/location-activity.test.ts b/packages/core/test/location-activity.test.ts index 39b94c786683..0aba76001c77 100644 --- a/packages/core/test/location-activity.test.ts +++ b/packages/core/test/location-activity.test.ts @@ -50,9 +50,7 @@ const locations = Layer.effect( const forms = yield* Form.Service return SessionRunner.Service.of({ drain: ({ sessionID }) => - (sessionID.startsWith("ses_active_work") || - sessionID === Session.ID.make("ses_quiet_work") || - sessionID === Session.ID.make("ses_permission_work") + (sessionID.startsWith("ses_active_work") || sessionID === Session.ID.make("ses_permission_work") ? Effect.never : forms.ask({ sessionID, @@ -107,131 +105,104 @@ const it = testEffect( ), ) +const seed = Effect.fn("LocationActivityTest.seed")(function* ( + directory: AbsolutePath, + rows: readonly { id: Session.ID; parent_id?: Session.ID }[], +) { + const database = yield* Database.Service + yield* database.db.insert(ProjectTable).values({ id: Project.ID.global, worktree: directory, sandboxes: [] }).run() + yield* database.db + .insert(SessionTable) + .values( + rows.map((row) => ({ + project_id: Project.ID.global, + directory, + slug: "activity", + title: "Session", + version: "test", + ...row, + })), + ) + .run() +}) + describe("LocationActivity eviction", () => { - for (const [count, admission] of [[1, "same"]] as const) { - it.effect( - `expires ${count} waiting executions before eviction (${admission} session admitted during cleanup)`, - () => - Effect.gen(function* () { - const db = (yield* Database.Service).db - const bus = yield* Bus.Service - const map = yield* LocationServiceMap.Service - const execution = yield* SessionExecution.Service - const store = yield* SessionStore.Service - const sessionIDs = Array.from({ length: count }, (_, index) => - Session.ID.make(`ses_waiting_question_${index}`), - ) - const newcomer = sessionIDs[0] - const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) - const idle = Location.Ref.make({ directory: ref.directory, workspaceID: Workspace.ID.make("wrk_idle") }) - yield* db - .insert(ProjectTable) - .values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] }) - .run() - .pipe(Effect.orDie) - yield* db - .insert(SessionTable) - .values( - Array.from(new Set([...sessionIDs, newcomer]), (sessionID) => ({ - id: sessionID, - project_id: Project.ID.global, - slug: "question", - directory: ref.directory, - title: "Waiting question", - version: "test", - })), - ) - .run() - .pipe(Effect.orDie) + it.effect("expires a waiting execution without evicting work admitted during its cleanup", () => + Effect.gen(function* () { + const bus = yield* Bus.Service + const map = yield* LocationServiceMap.Service + const execution = yield* SessionExecution.Service + const store = yield* SessionStore.Service + const sessionID = Session.ID.make("ses_waiting_question") + const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) + const idle = Location.Ref.make({ directory: ref.directory, workspaceID: Workspace.ID.make("wrk_idle") }) + yield* seed(ref.directory, [{ id: sessionID }]) - const created = yield* Deferred.make() - const newCreated = yield* Deferred.make() - const pending: Form.Info[] = [] - const interrupted: SessionEvent.Execution.Interrupted["data"][] = [] - const unsubscribe = yield* bus.listen((event) => - Effect.gen(function* () { - if (event.type === SessionEvent.Execution.Interrupted.type) { - interrupted.push(Schema.decodeUnknownSync(SessionEvent.Execution.Interrupted.data)(event.data)) - } - if (event.type !== Form.Event.Created.type) return - pending.push(Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form) - if (pending.length === count) yield* Deferred.succeed(created, undefined) - if (pending.length > count) yield* Deferred.succeed(newCreated, undefined) - }), - ) - yield* Effect.addFinalizer(() => unsubscribe) - const running = yield* Effect.forEach(sessionIDs, (sessionID) => - execution.resume(sessionID).pipe(Effect.exit, Effect.forkScoped), - ) - yield* Effect.addFinalizer(() => - Effect.forEach([...sessionIDs, newcomer], (sessionID) => execution.interrupt(sessionID)).pipe( - Effect.andThen(TestClock.adjust("5 minutes")), - ), - ) - yield* Deferred.await(created) - const context = yield* map.contextEffect(ref).pipe(Effect.scoped) - const forms = Context.get(context, Form.Service) - expect((yield* store.listSuspended()).toSorted()).toEqual(sessionIDs.toSorted()) - yield* Location.Service.pipe(Effect.provide(map.get(idle)), Effect.scoped) + const created = yield* Deferred.make() + const newCreated = yield* Deferred.make() + const pending: Form.Info[] = [] + const interrupted: SessionEvent.Execution.Interrupted["data"][] = [] + const unsubscribe = yield* bus.listen((event) => + Effect.gen(function* () { + if (event.type === SessionEvent.Execution.Interrupted.type) { + interrupted.push(Schema.decodeUnknownSync(SessionEvent.Execution.Interrupted.data)(event.data)) + } + if (event.type !== Form.Event.Created.type) return + pending.push(Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form) + if (pending.length === 1) yield* Deferred.succeed(created, undefined) + if (pending.length === 2) yield* Deferred.succeed(newCreated, undefined) + }), + ) + yield* Effect.addFinalizer(() => unsubscribe) + const running = yield* execution.resume(sessionID).pipe(Effect.exit, Effect.forkScoped) + yield* Effect.addFinalizer(() => + execution.interrupt(sessionID).pipe(Effect.andThen(TestClock.adjust("5 minutes"))), + ) + yield* Deferred.await(created) + const context = yield* map.contextEffect(ref).pipe(Effect.scoped) + const forms = Context.get(context, Form.Service) + expect(yield* store.listSuspended()).toEqual([sessionID]) + yield* Location.Service.pipe(Effect.provide(map.get(idle)), Effect.scoped) - // Human input produces no durable activity while the question is pending. - yield* TestClock.adjust("1 minute") - yield* TestClock.adjust("62 minutes") - // Interruption has cancelled each question, but slow cleanup still owns the graph. - expect(Array.from(yield* execution.active).toSorted()).toEqual(sessionIDs.toSorted()) - expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) - for (const form of pending) expect(yield* forms.state(form.id)).toEqual({ status: "cancelled" }) + // Human input produces no durable activity while the question is pending. + yield* TestClock.adjust("1 minute") + yield* TestClock.adjust("62 minutes") + // Interruption has cancelled the question, but slow cleanup still owns the graph. + expect(Array.from(yield* execution.active)).toEqual([sessionID]) + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) + for (const form of pending) expect(yield* forms.state(form.id)).toEqual({ status: "cancelled" }) - yield* execution.wake(newcomer) - yield* TestClock.adjust("5 minutes") - yield* Deferred.await(newCreated) - const results = yield* Effect.forEach(running, Fiber.join) - expect(results.every((exit) => exit._tag === "Failure")).toBe(true) - expect(Array.from(yield* execution.active)).toEqual([newcomer]) - expect(yield* store.listSuspended()).toEqual([newcomer]) - expect(interrupted.toSorted((a, b) => a.sessionID.localeCompare(b.sessionID))).toEqual( - sessionIDs.map((sessionID) => ({ sessionID, reason: "inactivity" })), - ) - expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) - expect(yield* forms.list({ sessionID: newcomer })).toEqual([pending[count]]) - yield* execution.interrupt(newcomer) - yield* TestClock.adjust("5 minutes") - yield* execution.awaitIdle(newcomer) - yield* TestClock.adjust("62 minutes") - expect(yield* store.listSuspended()).toEqual([]) - expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([]) - }), - ) - } + yield* execution.wake(sessionID) + yield* TestClock.adjust("5 minutes") + yield* Deferred.await(newCreated) + expect((yield* Fiber.join(running))._tag).toBe("Failure") + expect(Array.from(yield* execution.active)).toEqual([sessionID]) + expect(yield* store.listSuspended()).toEqual([sessionID]) + expect(interrupted).toEqual([{ sessionID, reason: "inactivity" }]) + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) + expect(yield* forms.list({ sessionID })).toEqual([pending[1]]) + yield* execution.interrupt(sessionID) + yield* TestClock.adjust("5 minutes") + yield* execution.awaitIdle(sessionID) + yield* TestClock.adjust("62 minutes") + expect(yield* store.listSuspended()).toEqual([]) + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([]) + }), + ) - it.effect("keeps pending requests alive while their location progresses, then expires only its idle sessions", () => + it.effect("expires unanswered requests without interrupting another progressing session in the same location", () => Effect.gen(function* () { - const db = (yield* Database.Service).db const bus = yield* Bus.Service const map = yield* LocationServiceMap.Service const execution = yield* SessionExecution.Service const waiting = Session.ID.make("ses_waiting_question") const permission = Session.ID.make("ses_permission_work") - const quiet = Session.ID.make("ses_quiet_work") const working = Session.ID.make("ses_active_work") - const elsewhere = Session.ID.make("ses_active_work_elsewhere") const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) - const other = Location.Ref.make({ directory: ref.directory, workspaceID: Workspace.ID.make("wrk_other") }) - yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] }).run() - yield* db - .insert(SessionTable) - .values( - [waiting, permission, quiet, working, elsewhere].map((id) => ({ - id, - project_id: Project.ID.global, - slug: "question", - directory: ref.directory, - workspace_id: id === elsewhere ? other.workspaceID : undefined, - title: "Session", - version: "test", - })), - ) - .run() + yield* seed( + ref.directory, + [waiting, permission, working].map((id) => ({ id })), + ) const created = yield* Deferred.make() const interrupted: SessionEvent.Execution.Interrupted["data"][] = [] const unsubscribe = yield* bus.listen((event) => { @@ -248,11 +219,9 @@ describe("LocationActivity eviction", () => { yield* Effect.addFinalizer(() => unsubscribe) yield* execution.resume(waiting).pipe(Effect.exit, Effect.forkScoped) yield* execution.resume(permission).pipe(Effect.exit, Effect.forkScoped) - yield* execution.resume(quiet).pipe(Effect.exit, Effect.forkScoped) yield* execution.resume(working).pipe(Effect.exit, Effect.forkScoped) - yield* execution.resume(elsewhere).pipe(Effect.exit, Effect.forkScoped) yield* Effect.addFinalizer(() => - Effect.forEach([waiting, permission, quiet, working, elsewhere], (id) => execution.interrupt(id)).pipe( + Effect.forEach([waiting, permission, working], (id) => execution.interrupt(id)).pipe( Effect.andThen(TestClock.adjust("5 minutes")), ), ) @@ -273,109 +242,96 @@ describe("LocationActivity eviction", () => { id: "tool-running", metadata: { status: "still working" }, }) - yield* bus.publish(SessionEvent.Text.Delta, { - sessionID: elsewhere, - assistantMessageID: SessionMessage.ID.make("msg_elsewhere"), - ordinal: 0, - delta: "another workspace is still working", - }) - yield* TestClock.adjust("33 minutes") - yield* TestClock.adjust("5 minutes") - expect(Array.from(yield* execution.active).toSorted()).toEqual( - [waiting, permission, working, elsewhere].toSorted(), - ) - expect(yield* forms.state(request.id)).toEqual({ status: "pending" }) - expect(interrupted).toEqual([{ sessionID: quiet, reason: "inactivity" }]) - - yield* execution.interrupt(working) - yield* TestClock.adjust("5 minutes") - yield* execution.awaitIdle(working) - yield* TestClock.adjust("10 minutes") - // Neither activity in another workspace nor viewing the question renews this location. - yield* bus.publish(SessionEvent.Text.Delta, { - sessionID: elsewhere, - assistantMessageID: SessionMessage.ID.make("msg_elsewhere"), - ordinal: 0, - delta: "still working in the other workspace", - }) + // Neither another Session's progress nor viewing this question renews its own deadline. yield* bus.publish(SessionEvent.Viewed, { sessionID: waiting, idle: 0 }, { location: ref }) - expect(yield* forms.state(request.id)).toEqual({ status: "pending" }) - yield* TestClock.adjust("53 minutes") + yield* TestClock.adjust("33 minutes") yield* TestClock.adjust("5 minutes") - expect(Array.from(yield* execution.active)).toEqual([elsewhere]) - expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([other]) + expect(Array.from(yield* execution.active)).toEqual([working]) expect(yield* forms.state(request.id)).toEqual({ status: "cancelled" }) + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) expect(interrupted.toSorted((a, b) => a.sessionID.localeCompare(b.sessionID))).toEqual([ - { sessionID: working, reason: "user" }, { sessionID: permission, reason: "inactivity" }, - { sessionID: quiet, reason: "inactivity" }, { sessionID: waiting, reason: "inactivity" }, ]) + + yield* execution.interrupt(working) + yield* TestClock.adjust("5 minutes") + yield* execution.awaitIdle(working) + yield* TestClock.adjust("62 minutes") + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([]) }), ) - it.effect("counts foreground child progress for its parent, but not background progress", () => - Effect.gen(function* () { - const db = (yield* Database.Service).db - const bus = yield* Bus.Service - const execution = yield* SessionExecution.Service - const jobs = yield* Job.Service - const parent = Session.ID.make("ses_quiet_work") - const child = Session.ID.make("ses_active_work_child") - const directory = AbsolutePath.make("/project") - yield* db.insert(ProjectTable).values({ id: Project.ID.global, worktree: directory, sandboxes: [] }).run() - yield* db - .insert(SessionTable) - .values([ - { - id: parent, - project_id: Project.ID.global, - slug: "parent", - directory, - title: "Parent", - version: "test", - }, - { - id: child, - parent_id: parent, - project_id: Project.ID.global, - slug: "child", - directory, - title: "Child", - version: "test", + it.effect( + "keeps a waiting parent alive through foreground and background child work, then starts its idle window", + () => + Effect.gen(function* () { + const bus = yield* Bus.Service + const execution = yield* SessionExecution.Service + const jobs = yield* Job.Service + const map = yield* LocationServiceMap.Service + const parent = Session.ID.make("ses_waiting_parent") + const child = Session.ID.make("ses_active_work_child") + const directory = AbsolutePath.make("/project") + yield* seed(directory, [{ id: parent }, { id: child, parent_id: parent }]) + const created = yield* Deferred.make() + const unsubscribe = yield* bus.listen((event) => + event.type === Form.Event.Created.type + ? Deferred.succeed(created, Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form).pipe( + Effect.asVoid, + ) + : Effect.void, + ) + yield* Effect.addFinalizer(() => unsubscribe) + yield* execution.resume(parent).pipe(Effect.exit, Effect.forkScoped) + yield* jobs.start({ + id: child, + type: "subagent", + recovery: { + kind: "subagent", + parentSessionID: parent, + childSessionID: child, + agent: "explore", + description: "Work", }, - ]) - .run() - yield* execution.resume(parent).pipe(Effect.exit, Effect.forkScoped) - yield* execution.resume(child).pipe(Effect.exit, Effect.forkScoped) - yield* jobs.start({ id: child, type: "subagent", run: Effect.never }) - yield* jobs.block({ id: child, sessionID: parent }).pipe(Effect.forkScoped) - yield* Effect.addFinalizer(() => - Effect.forEach([parent, child], (id) => execution.interrupt(id)).pipe( - Effect.andThen(TestClock.adjust("5 minutes")), - ), - ) - yield* TestClock.adjust("30 minutes") - yield* bus.publish(SessionEvent.Text.Delta, { - sessionID: child, - assistantMessageID: SessionMessage.ID.make("msg_child"), - ordinal: 0, - delta: "still working", - }) - yield* TestClock.adjust("38 minutes") - expect((yield* execution.active).has(parent)).toBe(true) - expect((yield* execution.active).has(child)).toBe(true) - yield* jobs.background(child) - yield* TestClock.adjust("7 minutes") - yield* bus.publish(SessionEvent.Text.Delta, { - sessionID: child, - assistantMessageID: SessionMessage.ID.make("msg_child"), - ordinal: 0, - delta: "still working", - }) - yield* TestClock.adjust("25 minutes") - expect((yield* execution.active).has(parent)).toBe(false) - expect((yield* execution.active).has(child)).toBe(true) - }), + run: execution.resume(child).pipe(Effect.as("done")), + }) + yield* jobs.block({ id: child, sessionID: parent }).pipe(Effect.forkScoped) + yield* Effect.addFinalizer(() => + Effect.forEach([parent, child], (id) => execution.interrupt(id)).pipe( + Effect.andThen(TestClock.adjust("5 minutes")), + ), + ) + const request = yield* Deferred.await(created) + const context = yield* map.contextEffect(Location.Ref.make({ directory })).pipe(Effect.scoped) + const forms = Context.get(context, Form.Service) + yield* TestClock.adjust("30 minutes") + yield* bus.publish(SessionEvent.Text.Delta, { + sessionID: child, + assistantMessageID: SessionMessage.ID.make("msg_child"), + ordinal: 0, + delta: "still working", + }) + yield* TestClock.adjust("38 minutes") + expect((yield* execution.active).has(parent)).toBe(true) + expect((yield* execution.active).has(child)).toBe(true) + yield* jobs.background(child) + yield* TestClock.adjust("7 minutes") + expect((yield* execution.active).has(parent)).toBe(true) + expect((yield* execution.active).has(child)).toBe(true) + expect(yield* forms.state(request.id)).toEqual({ status: "pending" }) + // The child retains its own stall timeout: it cannot pin the parent forever without progress. + yield* TestClock.adjust("20 minutes") + yield* jobs.wait({ id: child }) + expect((yield* execution.active).has(child)).toBe(false) + expect(yield* forms.state(request.id)).toEqual({ status: "pending" }) + yield* TestClock.adjust("59 minutes") + expect((yield* execution.active).has(parent)).toBe(true) + expect(yield* forms.state(request.id)).toEqual({ status: "pending" }) + yield* TestClock.adjust("2 minutes") + expect(yield* forms.state(request.id)).toEqual({ status: "cancelled" }) + yield* TestClock.adjust("5 minutes") + expect((yield* execution.active).has(parent)).toBe(false) + }), ) }) diff --git a/packages/core/test/session-execution.test.ts b/packages/core/test/session-execution.test.ts index 063c5fedfbc9..9e48865072ab 100644 --- a/packages/core/test/session-execution.test.ts +++ b/packages/core/test/session-execution.test.ts @@ -6,6 +6,7 @@ import { LayerNode } from "@opencode/util/effect/layer-node" import { Bus } from "@opencode/core/bus" import { Instance } from "@opencode/core/instance/service" import { Job } from "@opencode/core/job" +import { Inactivity } from "@opencode/core/inactivity" import { KV } from "@opencode/core/kv" import { LocationServiceMap } from "@opencode/core/location-service-map" import type { LocationServices } from "@opencode/core/location-services" @@ -27,7 +28,16 @@ import { testEffect } from "./lib/effect" const it = testEffect( AppNodeBuilder.build( - LayerNode.group([Database.node, Bus.node, SessionStore.node, SessionInbox.node, Job.node, KV.node, Session.node]), + LayerNode.group([ + Database.node, + Bus.node, + SessionStore.node, + SessionInbox.node, + Job.node, + KV.node, + Session.node, + Inactivity.node, + ]), ), ) diff --git a/packages/core/test/session-shell.test.ts b/packages/core/test/session-shell.test.ts index 1340c0f67ed8..7fc6e0f27676 100644 --- a/packages/core/test/session-shell.test.ts +++ b/packages/core/test/session-shell.test.ts @@ -1,10 +1,12 @@ import { describe, expect, setDefaultTimeout } from "bun:test" import fs from "fs/promises" import path from "path" -import { Cause, Context, Deferred, Effect, Exit, Fiber, Layer, Option, Schedule, Stream } from "effect" +import { Cause, Context, Deferred, Effect, Exit, Fiber, Layer, Option, RcMap, Schedule, Stream } from "effect" +import { TestClock } from "effect/testing" import { Bus } from "@opencode/core/bus" import { AppNodeBuilder } from "@opencode/core/effect/app-node-builder" import { Location } from "@opencode/core/location" +import { LocationActivity } from "@opencode/core/location-activity" import { LocationServiceMap } from "@opencode/core/location-service-map" import { AbsolutePath } from "@opencode/core/schema" import { Session } from "@opencode/core/session" @@ -13,6 +15,8 @@ import { SessionExecution } from "@opencode/core/session/execution" import { SessionRunCoordinator } from "@opencode/core/session/run-coordinator" import { Shell } from "@opencode/core/shell" import { LayerNode } from "@opencode/util/effect/layer-node" +import { Global } from "@opencode/util/global" +import { tempGlobalLayer } from "./fixture/global" import { location } from "./fixture/location" import { offlineModels } from "./fixture/models" import { tmpdirScoped } from "./fixture/tmpdir" @@ -64,6 +68,18 @@ const it = testEffect( ]).pipe(Layer.provideMerge(controlLayer)), ) +const idleIt = testEffect( + AppNodeBuilder.build( + LayerNode.group([Bus.node, Session.node, SessionExecution.node, LocationServiceMap.node, LocationActivity.node]), + [ + Global.node.replace(tempGlobalLayer), + Bus.node.replace(Bus.configured({ persist: true })), + SessionExecution.node.replace(executionLayer.pipe(Layer.provide(controlLayer))), + offlineModels, + ], + ).pipe(Layer.provideMerge(controlLayer)), +) + const setup = Effect.gen(function* () { const tmp = yield* tmpdirScoped() const session = yield* Session.Service @@ -195,27 +211,46 @@ describe("Session.shell", () => { }), ) - it.live("does not suppress a normal prompt wake while a shell is running or wake again when it completes", () => - Effect.gen(function* () { - const fixture = yield* setup - const command = yield* launch(fixture, "prompt") + idleIt.effect( + "keeps a quiet shell's Session and location alive without suppressing prompt wakes or waking on completion", + () => + Effect.gen(function* () { + const fixture = yield* setup + const locations = yield* LocationServiceMap.Service + const ref = LocationServiceMap.canonical(fixture.created.location) + const command = yield* launch(fixture, "prompt").pipe(TestClock.withLive) - const prompt = yield* fixture.session.prompt({ sessionID: fixture.created.id, text: "Continue while this runs" }) - yield* Deferred.await(fixture.control.started).pipe(Effect.timeout("5 seconds")) - expect(fixture.control.wakes).toEqual([fixture.created.id]) - expect(command.caller.pollUnsafe()).toBeUndefined() - yield* Deferred.succeed(fixture.control.release, undefined) - yield* fixture.execution.awaitIdle(fixture.created.id) + const prompt = yield* fixture.session.prompt({ + sessionID: fixture.created.id, + text: "Continue while this runs", + }) + yield* Deferred.await(fixture.control.started).pipe(Effect.timeout("5 seconds")) + expect(fixture.control.wakes).toEqual([fixture.created.id]) + expect(command.caller.pollUnsafe()).toBeUndefined() + // The real process emits nothing while waiting on the file gate; advance only inactivity time. + yield* TestClock.adjust("1 minute") + yield* TestClock.adjust("62 minutes") + expect((yield* fixture.execution.active).has(fixture.created.id)).toBe(true) + expect(yield* RcMap.has(locations.rcMap, ref)).toBe(true) + yield* Deferred.succeed(fixture.control.release, undefined) + yield* fixture.execution.awaitIdle(fixture.created.id) + yield* TestClock.adjust("62 minutes") + expect(yield* RcMap.has(locations.rcMap, ref)).toBe(true) + expect((yield* fixture.shell.get(command.shellID)).status).toBe("running") - yield* command.release - yield* Fiber.join(command.caller).pipe(Effect.timeout("5 seconds")) - expect(fixture.control.wakes).toEqual([fixture.created.id]) - expect(yield* fixture.execution.active).toEqual(new Set()) - expect(yield* fixture.session.inbox(fixture.created.id)).toMatchObject([ - { id: prompt.id, type: "user" }, - { type: "synthetic", payload: { metadata: { source: "shell" } } }, - ]) - }), + yield* command.release + yield* Fiber.join(command.caller).pipe(Effect.timeout("5 seconds"), TestClock.withLive) + expect(fixture.control.wakes).toEqual([fixture.created.id]) + expect(yield* fixture.execution.active).toEqual(new Set()) + expect(yield* fixture.session.inbox(fixture.created.id)).toMatchObject([ + { id: prompt.id, type: "user" }, + { type: "synthetic", payload: { metadata: { source: "shell" } } }, + ]) + yield* TestClock.adjust("59 minutes") + expect(yield* RcMap.has(locations.rcMap, ref)).toBe(true) + yield* TestClock.adjust("2 minutes") + expect(yield* RcMap.has(locations.rcMap, ref)).toBe(false) + }), ) for (const exit of [0, 7]) {