Compare commits

...
Author SHA1 Message Date
jlongster 3a7de260ba fix(core): keep background shell locations alive 2026-10-01 20:48:40 +00:00
6 changed files with 194 additions and 4 deletions

No files matched your search

+19
View File
@@ -4,6 +4,7 @@ import { Array, Cause, Clock, Context, Deferred, Effect, Exit, Layer, Schema, Sc
import { makeGlobalNode } from "@opencode/util/effect/app-node"
import { Identifier } from "./id/id.js"
import { KV } from "./kv.js"
import { Location } from "./location.js"
import { SessionMessage } from "./session/message.js"
import { SessionSchema } from "./session/schema.js"
@@ -53,6 +54,7 @@ export type Info = {
type Active = {
info: Info
location?: Location.Ref
done: Deferred.Deferred<Info>
backgrounded: Deferred.Deferred<Info>
scope: Scope.Closeable
@@ -98,6 +100,7 @@ export type StartInput = {
type: string
title?: string
metadata?: Record<string, unknown>
location?: Location.Ref
recovery?: Recovery
notificationID?: SessionMessage.ID
run: Effect.Effect<string, unknown>
@@ -133,6 +136,8 @@ export interface Interface {
readonly background: (id: string) => Effect.Effect<Info | undefined>
readonly backgroundAll: (input: BackgroundAllInput) => Effect.Effect<Info[]>
readonly cancel: (id: string) => Effect.Effect<Info | undefined>
/** Locations owning process-local background shells that are still running. */
readonly runningBackgroundShellLocations: Effect.Effect<readonly Location.Ref[]>
readonly pendingBackground: Effect.Effect<readonly Background[]>
readonly completeBackground: (notificationID: SessionMessage.ID) => Effect.Effect<void>
}
@@ -275,6 +280,7 @@ export const make = Effect.gen(function* () {
metadata: input.metadata,
...(input.notificationID ? { notificationID: input.notificationID } : {}),
},
location: input.location,
done,
backgrounded,
scope,
@@ -451,6 +457,18 @@ export const make = Effect.gen(function* () {
return recovered
}).pipe(Effect.withSpan("Job.pendingBackground"))
const runningBackgroundShellLocations: Interface["runningBackgroundShellLocations"] = SynchronizedRef.get(
state.jobs,
).pipe(
Effect.map((jobs) =>
[...jobs.values()].flatMap((job) =>
job.info.status === "running" && job.isBackgrounded && job.recovery?.kind === "shell" && job.location
? [job.location]
: [],
),
),
)
const completeBackground: Interface["completeBackground"] = Effect.fn("Job.completeBackground")((notificationID) =>
SynchronizedRef.updateEffect(state.jobs, (jobs) =>
Effect.gen(function* () {
@@ -472,6 +490,7 @@ export const make = Effect.gen(function* () {
background,
backgroundAll,
cancel,
runningBackgroundShellLocations,
pendingBackground,
completeBackground,
})
+16 -3
View File
@@ -2,6 +2,7 @@ export * as LocationActivity from "./location-activity.js"
import { Clock, Context, Duration, Effect, Layer, RcMap, Schema } from "effect"
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"
@@ -22,6 +23,7 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly
const locations = yield* LocationServiceMap.Service
const execution = yield* SessionExecution.Service
const sessions = yield* SessionStore.Service
const jobs = yield* Job.Service
const timeToLive = Duration.toMillis(options.timeToLive ?? "60 minutes")
const entries = new Map<string, { readonly ref: Location.Ref; expiresAt: number }>()
const key = (ref: Location.Ref) => `${LocationServiceMap.canonical(ref).directory}\0${ref.workspaceID ?? ""}`
@@ -29,6 +31,9 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly
Effect.sync(() => {
entries.set(key(ref), { ref, expiresAt: clock.currentTimeMillisUnsafe() + timeToLive })
})
const runningShellLocations = jobs.runningBackgroundShellLocations.pipe(
Effect.map((running) => new Set(running.map(key))),
)
const unsubscribe = yield* bus.listen((event) => {
if (!isSessionEvent(event)) return Effect.void
@@ -50,11 +55,16 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly
const now = clock.currentTimeMillisUnsafe()
const expired = Array.from(entries.values()).filter((entry) => entry.expiresAt <= now)
if (expired.length === 0) return
const shells = yield* runningShellLocations
const active = yield* Effect.forEach(yield* execution.active, (sessionID) => sessions.get(sessionID))
yield* Effect.forEach(
expired,
(entry) =>
Effect.gen(function* () {
if (shells.has(key(entry.ref))) {
yield* touch(entry.ref)
return
}
const owners = active.flatMap((session) =>
session && key(session.location) === key(entry.ref) ? [session] : [],
)
@@ -69,8 +79,11 @@ export function layer(options: { readonly timeToLive?: Duration.Input; readonly
},
)
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))) {
// New execution or a background shell admitted during cleanup may now own the cached graph.
if (
remaining.some((session) => session && key(session.location) === key(entry.ref)) ||
(yield* runningShellLocations).has(key(entry.ref))
) {
yield* touch(entry.ref)
return
}
@@ -92,5 +105,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],
})
+3
View File
@@ -18,6 +18,7 @@ import { Shell } from "../../shell.js"
import { ShellParse } from "../../shell/parse.js"
import { ShellSelect } from "../../shell/select.js"
import { ShellResult } from "../../shell/result.js"
import { Location } from "../../location.js"
export const name = "shell"
export const DEFAULT_TIMEOUT_MS = 2 * 60 * 1_000
@@ -108,6 +109,7 @@ export const Plugin = {
const environment = yield* Environment.Service
const access = yield* FileAccess.Service
const shell = yield* Shell.Service
const location = yield* Location.Service
const shellSelect = yield* ShellSelect.Service
const compatibleShell = shellSelect.resolve({ priority: "compat" })
const permission = yield* Permission.Service
@@ -239,6 +241,7 @@ export const Plugin = {
type: name,
title: info.command,
metadata: { sessionID: context.sessionID, shellID: info.id },
location: Location.Ref.make({ directory: location.directory, workspaceID: location.workspaceID }),
recovery: {
kind: "shell",
sessionID: context.sessionID,
+2 -1
View File
@@ -7,6 +7,7 @@ import { makeGlobalNode } from "@opencode/util/effect/app-node"
import { Global } from "@opencode/util/global"
import { Bus } from "@opencode/core/bus"
import { Database } from "@opencode/core/database/database"
import { Job } from "@opencode/core/job"
import { Location } from "@opencode/core/location"
import { LocationActivity } from "@opencode/core/location-activity"
import { LocationServiceMap } from "@opencode/core/location-services"
@@ -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],
}),
),
],
@@ -7,6 +7,7 @@ import { makeGlobalNode } from "@opencode/util/effect/app-node"
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"
@@ -83,6 +84,7 @@ const it = testEffect(
LayerNode.group([
Database.node,
Bus.node,
Job.node,
SessionStore.node,
LocationServiceMap.node,
SessionExecution.node,
@@ -101,6 +103,149 @@ const it = testEffect(
)
describe("LocationActivity eviction", () => {
it.effect("keeps a running background shell's location but evicts its idle parent location", () =>
Effect.gen(function* () {
const db = (yield* Database.Service).db
const map = yield* LocationServiceMap.Service
const jobs = yield* Job.Service
const parent = Session.ID.make("ses_background_parent")
const child = Session.ID.make("ses_background_child")
const shell = Session.ID.make("ses_background_shell")
const parentRef = LocationServiceMap.canonical({ directory: AbsolutePath.make("/parent") })
const shellRef = LocationServiceMap.canonical({ directory: AbsolutePath.make("/shell") })
yield* db
.insert(ProjectTable)
.values({ id: Project.ID.global, worktree: parentRef.directory, sandboxes: [] })
.run()
.pipe(Effect.orDie)
yield* db
.insert(SessionTable)
.values([
{
id: parent,
project_id: Project.ID.global,
slug: "parent",
directory: parentRef.directory,
title: "Parent",
version: "test",
},
{
id: child,
project_id: Project.ID.global,
slug: "child",
directory: AbsolutePath.make("/child"),
title: "Child",
version: "test",
parent_id: parent,
},
{
id: shell,
project_id: Project.ID.global,
slug: "shell",
directory: shellRef.directory,
title: "Shell",
version: "test",
},
])
.run()
.pipe(Effect.orDie)
yield* Location.Service.pipe(Effect.provide(map.get(parentRef)), Effect.scoped)
yield* Location.Service.pipe(Effect.provide(map.get(shellRef)), Effect.scoped)
const shellDone = yield* Deferred.make<string>()
yield* jobs.start({
id: "background-shell",
type: "shell",
location: shellRef,
recovery: { kind: "shell", sessionID: shell, shellID: "background-shell", command: "sleep" },
run: Deferred.await(shellDone),
})
yield* jobs.background("background-shell")
const childDone = yield* Deferred.make<string>()
yield* jobs.start({
id: child,
type: "subagent",
recovery: {
kind: "subagent",
parentSessionID: parent,
childSessionID: child,
agent: "test",
description: "Child",
},
run: Deferred.await(childDone),
})
yield* jobs.background(child)
yield* TestClock.adjust("1 minute")
yield* TestClock.adjust("62 minutes")
expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([shellRef])
expect((yield* jobs.get(child))?.status).toBe("running")
expect((yield* jobs.get("background-shell"))?.status).toBe("running")
yield* Deferred.succeed(shellDone, "done")
yield* jobs.wait({ id: "background-shell" })
yield* TestClock.adjust("62 minutes")
expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([])
yield* Deferred.succeed(childDone, "done")
}),
)
it.effect("keeps a location when a background shell starts during eviction 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 jobs = yield* Job.Service
const sessionID = Session.ID.make("ses_shell_during_cleanup")
const ref = LocationServiceMap.canonical({ directory: AbsolutePath.make("/cleanup") })
yield* db
.insert(ProjectTable)
.values({ id: Project.ID.global, worktree: ref.directory, sandboxes: [] })
.run()
.pipe(Effect.orDie)
yield* db
.insert(SessionTable)
.values({
id: sessionID,
project_id: Project.ID.global,
slug: "cleanup",
directory: ref.directory,
title: "Cleanup",
version: "test",
})
.run()
.pipe(Effect.orDie)
const created = yield* Deferred.make<void>()
const unsubscribe = yield* bus.listen((event) =>
event.type === Form.Event.Created.type ? Deferred.succeed(created, undefined).pipe(Effect.asVoid) : Effect.void,
)
yield* Effect.addFinalizer(() => unsubscribe)
const running = yield* execution.resume(sessionID).pipe(Effect.exit, Effect.forkScoped)
yield* Deferred.await(created)
yield* TestClock.adjust("1 minute")
yield* TestClock.adjust("62 minutes")
expect(Array.from(yield* execution.active)).toEqual([sessionID])
const done = yield* Deferred.make<string>()
yield* jobs.start({
id: "shell-during-cleanup",
type: "shell",
location: ref,
recovery: { kind: "shell", sessionID, shellID: "shell-during-cleanup", command: "sleep" },
run: Deferred.await(done),
})
yield* jobs.background("shell-during-cleanup")
yield* TestClock.adjust("5 minutes")
expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([ref])
expect((yield* Fiber.join(running))._tag).toBe("Failure")
yield* Deferred.succeed(done, "done")
yield* jobs.wait({ id: "shell-during-cleanup" })
yield* TestClock.adjust("62 minutes")
expect(Array.from(yield* RcMap.keys(map.rcMap))).toEqual([])
}),
)
for (const [count, admission] of [
[1, "none"],
[2, "none"],
+9
View File
@@ -136,6 +136,7 @@ const shellPluginSupervisor = makeLocationNode({
Permission.node,
Session.node,
Job.node,
Location.node,
Shell.node,
ShellSelect.node,
Tool.node,
@@ -1504,6 +1505,12 @@ describe("ShellTool", () => {
expect(settled.metadata).toMatchObject({ truncated: false })
expect(shellID).toStartWith("sh_")
const jobs = yield* Job.Service
const location = yield* Location.Service
expect(yield* jobs.runningBackgroundShellLocations).toEqual([
Location.Ref.make({ directory: location.directory, workspaceID: location.workspaceID }),
])
const shell = yield* Shell.Service
if (!shellID) return
const id = ID.make(shellID)
@@ -1520,6 +1527,8 @@ describe("ShellTool", () => {
])
expect((yield* shell.list()).map((info) => info.id)).toContain(id)
expect((yield* shell.wait(id)).status).toBe("timeout")
yield* jobs.wait({ id: shellID })
expect(yield* jobs.runningBackgroundShellLocations).toEqual([])
expect((yield* Fiber.join(admitted)).valueOrUndefined?.data.item.payload).toMatchObject({
text: expect.stringContaining("Timed out before completion"),
description: idleCommand,