Repository navigation
fix(core): preserve active sessions during idle cleanup - #53238
possibilities wants to merge 1 commit into
Conversation
|
Thanks for your contribution! This PR doesn't have a linked issue. All PRs must reference an existing issue. Please:
See CONTRIBUTING.md for details. |
|
The following comment was made by an LLM, it may be inaccurate: |
|
The linked bug is #51343, already referenced as The standards run checked only Could a maintainer clear the false-positive |
|
It seems to work quite well! For me, the question dialogs kept closing quite quickly, and with the fix, they stay open like in v1. |
Interesting, sounds completely different than what's fixed here, no? |
|
I don't know if it's completely different, but with the release version v2.0.20 the questions dialog got interrupted quite quickly (sometimes max. 1 minute timeout? after they pop up) with "aborted"/"Step interrupted". So this doesn't quite match the 60m timeout from #51343, but it is still a timeout issue that apparently isn't present anymore. But I used the v2.0.23 branch source and cherry-picked your fix. |
Very cool, thanks for explaining. <3 Now just gotta find out if PRs ever get reviewed and/or merged here. I'm new and that 1.7k PR count is intimidating :) |
|
^^ haha that's true. But, it's more like the disappearance of the question dialog was very annoying, and I wanted to fix it 😅 And tbf, I didn't check if my bug was not already resolved in v2.0.20..v2.0.23 :/ |
Might be worth dropping the cherry picked patch and see how it does. It could be interesting/edifying. |
|
Additional v2.0.23 regression evidence for this fix: I tested a local active-retention variant against upstream v2.0.23 ( The regressions cover active native question owners remaining cached across repeated expiry, admitting new work, eventual idle eviction after completion, separate workspace keys sharing a directory, and cleanup after a cancelled question has settled. These use the native test clock; a real >60-minute production acceptance run has not yet been completed. One source-level difference in this local variant: The tested v2.0.23 patch is included below for review; it is not a request to replace the broader coverage already in this PR. Tested v2.0.23 source and regression patch--- a/packages/core/src/location-activity.ts
+++ b/packages/core/src/location-activity.ts
@@ -50,26 +50,12 @@
const now = clock.currentTimeMillisUnsafe()
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.
+ // Work may have been admitted since expiration was observed.
if (remaining.some((session) => session && key(session.location) === key(entry.ref))) {
yield* touch(entry.ref)
return
--- a/packages/core/test/location-activity.test.ts
+++ b/packages/core/test/location-activity.test.ts
@@ -101,121 +101,172 @@
)
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)`,
- () =>
+ it.effect("keeps active question owners after repeated expiry and admits new work", () =>
+ 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 ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") })
+ const owners = [Session.ID.make("ses_question_owner_0"), Session.ID.make("ses_question_owner_1")]
+ const newcomer = Session.ID.make("ses_question_newcomer")
+ const sessionIDs = [...owners, newcomer]
+ yield* db
+ .insert(ProjectTable)
+ .values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] })
+ .run()
+ .pipe(Effect.orDie)
+ yield* db
+ .insert(SessionTable)
+ .values(
+ sessionIDs.map((id) => ({
+ id,
+ project_id: Project.ID.global,
+ slug: "question",
+ directory: ref.directory,
+ title: "Waiting question",
+ version: "test",
+ })),
+ )
+ .run()
+ .pipe(Effect.orDie)
+
+ const ownersCreated = yield* Deferred.make<void>()
+ const newcomerCreated = yield* Deferred.make<Form.Info>()
+ const pending: Form.Info[] = []
+ const interrupted: SessionEvent.Execution.Interrupted["data"][] = []
+ const unsubscribe = yield* bus.listen((event) =>
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 = 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
- .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)
-
- const created = yield* Deferred.make<void>()
- const newCreated = yield* Deferred.make<void>()
- 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)
-
- // 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])
- expect(yield* forms.list()).toEqual([])
- 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* 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])
- }
- 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([])
- }
+ 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
+ const form = Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form
+ pending.push(form)
+ if (pending.length === owners.length) yield* Deferred.succeed(ownersCreated, undefined)
+ if (form.sessionID === newcomer) yield* Deferred.succeed(newcomerCreated, form)
}),
- )
- }
+ )
+ yield* Effect.addFinalizer(() => unsubscribe)
+ const running = yield* Effect.forEach(owners, (id) => execution.resume(id).pipe(Effect.exit, Effect.forkScoped))
+ yield* Effect.addFinalizer(() =>
+ Effect.forEach(sessionIDs, (id) => execution.interrupt(id)).pipe(Effect.andThen(TestClock.adjust("5 minutes"))),
+ )
+ yield* Deferred.await(ownersCreated)
+ const context = yield* map.contextEffect(ref).pipe(Effect.scoped)
+ const forms = Context.get(context, Form.Service)
+
+ yield* TestClock.adjust("125 minutes")
+ expect(Array.from(yield* execution.active).toSorted()).toEqual(owners.toSorted())
+ expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref])
+ expect(yield* forms.list()).toEqual(pending)
+ expect(interrupted).toEqual([])
+
+ yield* execution.wake(newcomer)
+ yield* Deferred.await(newcomerCreated)
+ expect(Array.from(yield* execution.active).toSorted()).toEqual(sessionIDs.toSorted())
+ expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref])
+ expect(interrupted).toEqual([])
+
+ for (const form of pending) yield* forms.reply({ id: form.id, answer: { runtime: "bun" } })
+ expect((yield* Effect.forEach(running, Fiber.join)).every((exit) => exit._tag === "Success")).toBe(true)
+ yield* execution.awaitIdle(newcomer)
+ yield* TestClock.adjust("62 minutes")
+ expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([])
+ }),
+ )
+
+ it.effect("evicts an idle workspace without evicting a separate active workspace", () =>
+ 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 directory = AbsolutePath.make("/project")
+ const activeRef = LocationServiceMap.canonical({ directory, workspaceID: Workspace.ID.make("wrk_active") })
+ const idleRef = LocationServiceMap.canonical({ directory, workspaceID: Workspace.ID.make("wrk_idle") })
+ const id = Session.ID.make("ses_question_workspace")
+ yield* db
+ .insert(ProjectTable)
+ .values({ id: Project.ID.global, worktree: directory, sandboxes: [] })
+ .run()
+ .pipe(Effect.orDie)
+ yield* db
+ .insert(SessionTable)
+ .values({
+ id,
+ project_id: Project.ID.global,
+ slug: "question",
+ directory,
+ workspace_id: activeRef.workspaceID,
+ title: "Waiting question",
+ version: "test",
+ })
+ .run()
+ .pipe(Effect.orDie)
+ const formReady = yield* Deferred.make<Form.Info>()
+ const unsubscribe = yield* bus.listen((event) =>
+ event.type === Form.Event.Created.type
+ ? Deferred.succeed(formReady, Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form)
+ : Effect.void,
+ )
+ yield* Effect.addFinalizer(() => unsubscribe)
+ const running = yield* execution.resume(id).pipe(Effect.exit, Effect.forkScoped)
+ yield* Effect.addFinalizer(() => execution.interrupt(id).pipe(Effect.andThen(TestClock.adjust("5 minutes"))))
+ const form = yield* Deferred.await(formReady)
+ yield* Location.Service.pipe(Effect.provide(map.get(idleRef)), Effect.scoped)
+ yield* TestClock.adjust("62 minutes")
+ expect(Array.from(yield* execution.active)).toEqual([id])
+ expect(Array.from(yield* RcMap.keys(map.rcMap)).toSorted()).toEqual([activeRef])
+ const context = yield* map.contextEffect(activeRef).pipe(Effect.scoped)
+ yield* Context.get(context, Form.Service).reply({ id: form.id, answer: { runtime: "bun" } })
+ expect((yield* Fiber.join(running))._tag).toBe("Success")
+ }),
+ )
+
+ it.effect("evicts after a cancelled question has settled", () =>
+ 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 ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/project") })
+ const id = Session.ID.make("ses_question_cancelled_idle")
+ yield* db
+ .insert(ProjectTable)
+ .values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] })
+ .run()
+ .pipe(Effect.orDie)
+ yield* db
+ .insert(SessionTable)
+ .values({
+ id,
+ project_id: Project.ID.global,
+ slug: "question",
+ directory: ref.directory,
+ title: "Waiting question",
+ version: "test",
+ })
+ .run()
+ .pipe(Effect.orDie)
+ const formReady = yield* Deferred.make<Form.Info>()
+ const unsubscribe = yield* bus.listen((event) =>
+ event.type === Form.Event.Created.type
+ ? Deferred.succeed(formReady, Schema.decodeUnknownSync(Form.Event.Created.data)(event.data).form)
+ : Effect.void,
+ )
+ yield* Effect.addFinalizer(() => unsubscribe)
+ const running = yield* execution.resume(id).pipe(Effect.exit, Effect.forkScoped)
+ yield* Effect.addFinalizer(() => execution.interrupt(id).pipe(Effect.andThen(TestClock.adjust("5 minutes"))))
+ const form = yield* Deferred.await(formReady)
+ const context = yield* map.contextEffect(ref).pipe(Effect.scoped)
+ const forms = Context.get(context, Form.Service)
+ yield* forms.cancel(form.id)
+ expect((yield* Fiber.join(running))._tag).toBe("Success")
+ yield* execution.awaitIdle(id)
+ expect(yield* forms.state(form.id)).toEqual({ status: "cancelled" })
+
+ yield* TestClock.adjust("62 minutes")
+ expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([])
+ }),
+ )
})
|
Issue for this PR
Fixes #51343 (inactivity cancellation of active sessions; no configurable timeout is added).
Type of change
What does this PR do?
A running session can spend an hour waiting for a tool, model, or human input without emitting durable events. The inactivity sweep currently interrupts it anyway.
Retain the location while
SessionExecution.activeowns a session there, including interruption cleanup. This keeps requests routed to the original services; genuinely idle locations still expire after one hour. Explicit Stop is unchanged.This deliberately changes the policy selected in #47629: silent active runs are no longer stopped by location cleanup. Hung work must be stopped explicitly or reach its own provider/tool timeout. Unlike #50499 and #51583, this is an ownership-only fix, without new activity tracking. Detached terminal-only workloads are out of scope.
How did you verify your code works?
packages/core,bun run test: 5,601 passed, 41 skipped, 0 failed.bun run check: all 36 typecheck tasks passed; lint has no errors (existing warnings remain). Changed-file Prettier andgit diff --checkpassed.Screenshots / recordings
Not applicable; no UI changes.
Checklist