From f69d2a81d89257f9d62e442e67673ecb3437c5fa Mon Sep 17 00:00:00 2001 From: Mike Bannister Date: Sun, 4 Oct 2026 12:57:32 -0400 Subject: [PATCH] fix(core): preserve active sessions during idle cleanup --- packages/core/src/location-activity.ts | 19 +---- packages/core/test/location-activity.test.ts | 90 ++++++++++---------- 2 files changed, 46 insertions(+), 63 deletions(-) diff --git a/packages/core/src/location-activity.ts b/packages/core/src/location-activity.ts index 0425e9acab04..7b58d9fddbd3 100644 --- a/packages/core/src/location-activity.ts +++ b/packages/core/src/location-activity.ts @@ -55,22 +55,9 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly 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))) { + // Process-local ownership includes silent work, human waits, and + // interruption cleanup. None of these are inactivity cancellation signals. + if (active.some((session) => session && key(session.location) === key(entry.ref))) { yield* touch(entry.ref) return } diff --git a/packages/core/test/location-activity.test.ts b/packages/core/test/location-activity.test.ts index a71b0a82f459..92c375fa88de 100644 --- a/packages/core/test/location-activity.test.ts +++ b/packages/core/test/location-activity.test.ts @@ -24,7 +24,7 @@ import { Workspace } from "@opencode/core/workspace" import { testEffect } from "./lib/effect" // Keep real execution ownership, location caching, forms, and eviction. The fixture -// runner waits on a form instead of making a model request before asking a question. +// runner waits on a form or a silent effect instead of making a model request. const locations = Layer.effect( LocationServiceMap.Service, Effect.gen(function* () { @@ -56,6 +56,7 @@ const locations = Layer.effect( }) .pipe( Effect.orDie, + Effect.andThen(sessionID.endsWith("_silent") ? Effect.never : Effect.void), Effect.as(SessionRunner.DrainResult.Complete()), Effect.onInterrupt(() => Effect.sleep("5 minutes")), ), @@ -102,15 +103,14 @@ const it = testEffect( ) describe("LocationActivity eviction", () => { - for (const [count, admission] of [ - [1, "none"], - [2, "none"], - [1, "other"], - [1, "same"], + for (const [count, completion] of [ + [1, "answer"], + [2, "answer"], + [1, "stop"], + [1, "silent"], ] as const) { - const newWork = admission !== "none" it.effect( - `interrupts ${count} waiting executions before eviction (${admission} session admitted during cleanup)`, + `retains ${count} process-owned executions across inactivity sweeps until ${completion} releases ownership`, () => Effect.gen(function* () { const db = (yield* Database.Service).db @@ -119,9 +119,8 @@ describe("LocationActivity eviction", () => { const execution = yield* SessionExecution.Service const store = yield* SessionStore.Service const sessionIDs = Array.from({ length: count }, (_, index) => - Session.ID.make(`ses_waiting_question_${index}`), + Session.ID.make(`ses_waiting_question_${index}_${completion}`), ) - const newcomer = admission === "same" ? sessionIDs[0] : Session.ID.make("ses_new_question") const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") }) const idle = Location.Ref.make({ directory: ref.directory, workspaceID: Workspace.ID.make("wrk_idle") }) yield* db @@ -132,7 +131,7 @@ describe("LocationActivity eviction", () => { yield* db .insert(SessionTable) .values( - Array.from(new Set([...sessionIDs, newcomer]), (sessionID) => ({ + sessionIDs.map((sessionID) => ({ id: sessionID, project_id: Project.ID.global, slug: "question", @@ -145,7 +144,6 @@ describe("LocationActivity eviction", () => { .pipe(Effect.orDie) 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) => @@ -156,7 +154,6 @@ describe("LocationActivity eviction", () => { 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) @@ -164,7 +161,7 @@ describe("LocationActivity eviction", () => { execution.resume(sessionID).pipe(Effect.exit, Effect.forkScoped), ) yield* Effect.addFinalizer(() => - Effect.forEach([...sessionIDs, newcomer], (sessionID) => execution.interrupt(sessionID)).pipe( + Effect.forEach(sessionIDs, (sessionID) => execution.interrupt(sessionID)).pipe( Effect.andThen(TestClock.adjust("5 minutes")), ), ) @@ -174,48 +171,47 @@ describe("LocationActivity eviction", () => { expect((yield* store.listSuspended()).toSorted()).toEqual(sessionIDs.toSorted()) yield* Location.Service.pipe(Effect.provide(map.get(idle)), Effect.scoped) + // Answering this fixture leaves a silent drain (like a tool with no progress events). + if (completion === "silent") { + yield* forms.reply({ id: pending[0].id, answer: { runtime: "bun" } }) + } + // 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. + yield* TestClock.adjust("62 minutes") 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" }) + expect(interrupted).toEqual([]) + // Resolve through the routed graph, not merely a retained reference to an evicted graph. + const retained = yield* map.contextEffect(ref).pipe(Effect.scoped) + expect(Context.get(retained, Form.Service)).toBe(forms) + expect(yield* forms.list()).toEqual(completion === "silent" ? [] : pending) - if (newWork) { - yield* execution.wake(newcomer) - if (admission === "other") yield* Deferred.await(newCreated) - } - yield* TestClock.adjust("5 minutes") - if (newWork) 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(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]]) - 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]) + if (completion === "answer") { + yield* Effect.forEach(pending, (form) => forms.reply({ id: form.id, answer: { runtime: "bun" } })) + for (const form of pending) { + expect(yield* forms.state(form.id)).toEqual({ status: "answered", answer: { runtime: "bun" } }) } - yield* execution.interrupt(newcomer) + } + if (completion !== "answer") { + // Cross the next expiry deadline while explicit Stop's finalizer is still running. + yield* TestClock.adjust("55 minutes") + expect(yield* execution.interrupt(sessionIDs[0])).toBe(true) + yield* TestClock.adjust("2 minutes") + expect(Array.from(yield* execution.active)).toEqual(sessionIDs) + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref]) + expect(Context.get(yield* map.contextEffect(ref).pipe(Effect.scoped), Form.Service)).toBe(forms) + if (completion === "stop") expect(yield* forms.state(pending[0].id)).toEqual({ status: "cancelled" }) 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([]) } + const results = yield* Effect.forEach(running, Fiber.join) + expect(results.every((exit) => exit._tag === (completion === "answer" ? "Success" : "Failure"))).toBe(true) + expect(Array.from(yield* execution.active)).toEqual([]) + expect(yield* store.listSuspended()).toEqual([]) + expect(interrupted).toEqual(completion === "answer" ? [] : [{ sessionID: sessionIDs[0], reason: "user" }]) + yield* TestClock.adjust("62 minutes") + expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([]) }), ) }