Skip to content

fix(core): preserve active sessions during idle cleanup - #53238

Open
possibilities wants to merge 1 commit into
anomalyco:devfrom
possibilities:session-inactivity
Open

possibilities wants to merge 1 commit into
anomalyco:devfrom
possibilities:session-inactivity

Conversation

@possibilities

Copy link
Copy Markdown

Issue for this PR

Fixes #51343 (inactivity cancellation of active sessions; no configurable timeout is added).

Type of change

  • Bug fix
  • New feature
  • Refactor / code improvement
  • Documentation

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.active owns 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?

  • Four regressions fail on the original code and pass after the fix. They cover forms answered after 125 minutes, multiple owners, silent work, independent idle workspaces, Stop cleanup, and eventual eviction.
  • From packages/core, bun run test: 5,601 passed, 41 skipped, 0 failed.
  • From the root, bun run check: all 36 typecheck tasks passed; lint has no errors (existing warnings remain). Changed-file Prettier and git diff --check passed.
  • Tests use real ownership/cache/forms/SQLite with a fixture runner, not live providers. An independent review checked settlement, movement, and explicit shutdown paths.

Screenshots / recordings

Not applicable; no UI changes.

Checklist

  • I have tested my changes locally
  • I have not included unrelated changes in this PR

@github-actions

github-actions Bot commented Oct 4, 2026

Copy link
Copy Markdown
Contributor

Thanks for your contribution!

This PR doesn't have a linked issue. All PRs must reference an existing issue.

Please:

  1. Open an issue describing the bug/feature (if one doesn't exist)
  2. Add Fixes #<number> or Closes #<number> to this PR description

See CONTRIBUTING.md for details.

@github-actions

github-actions Bot commented Oct 4, 2026

Copy link
Copy Markdown
Contributor

The following comment was made by an LLM, it may be inaccurate:

@possibilities

Copy link
Copy Markdown
Author

The linked bug is #51343, already referenced as Fixes #51343 under “Issue for this PR”. The PR targets v2 as required by CONTRIBUTING.md.

The standards run checked only closingIssuesReferences, which is empty for this non-default-branch target. The v2 workflow includes a PR-body fallback for that case, but the script that ran did not include it.

Could a maintainer clear the false-positive needs:issue label and approve the fork’s check, test, and nix-eval CI runs? They currently report action_required. Local full core tests and repository checks passed; their results are in the description.

@tlogemann

Copy link
Copy Markdown

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.

@possibilities

Copy link
Copy Markdown
Author

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?

@tlogemann

Copy link
Copy Markdown

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.

@possibilities

Copy link
Copy Markdown
Author

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 :)

@tlogemann

Copy link
Copy Markdown

^^ 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 :/

@possibilities

Copy link
Copy Markdown
Author

^^ 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.

@Ghilteras

Copy link
Copy Markdown

Additional v2.0.23 regression evidence for this fix:

I tested a local active-retention variant against upstream v2.0.23 (0fd7e2829449b052abf0078666669302923d77af) using Bun 1.4.2. The regression contract failed on unmodified source (1 pass, 2 fail) and passed after the retention change (3 pass). Six focused suites passed 101 tests, and direct-consumer checks passed 25 tests.

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: execution.active is re-read inside each expired-location handler rather than relying on one snapshot taken before iterating the expired entries. That narrows the stale-snapshot window for new work admitted during a sweep. This is a source comparison, not a claim that this race has been reproduced in production.

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([])
+    }),
+  )
 })

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

core: 60m idle location eviction interrupts a running session and rejects pending questions

3 participants