Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
92 changes: 92 additions & 0 deletions packages/core/src/inactivity.ts
Original file line number Diff line number Diff line change
@@ -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<void>
/** Owned work prevents expiration until its scope closes, then starts a fresh idle window. */
readonly hold: (owner: Owner) => Effect.Effect<void, never, Scope.Scope>
/** Find expired live executions and cached graphs; prune observations that no longer have owners. */
readonly expired: (input: {
readonly sessions: ReadonlySet<SessionSchema.ID>
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<SessionSchema.ID, Activity>()
const locations = new Map<string, Activity>()
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: [] })
12 changes: 10 additions & 2 deletions packages/core/src/job.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -177,6 +178,7 @@ function decrementSession(input: Map<SessionSchema.ID, number>, 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,
Expand Down Expand Up @@ -285,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,
Expand Down Expand Up @@ -477,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] })
133 changes: 72 additions & 61 deletions packages/core/src/location-activity.ts
Original file line number Diff line number Diff line change
@@ -1,96 +1,107 @@
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 { 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"
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)

/** Observe Session progress and expire process-local execution separately from shared Location resources. */
export class Service extends Context.Service<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 sessions = yield* SessionStore.Service
const inactivity = yield* Inactivity.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 ?? ""}`
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 === SessionEvent.Viewed.type) return
if (
isSessionEvent(event) &&
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
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 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
}
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],
deps: [Bus.node, LocationServiceMap.node, SessionExecution.node, SessionStore.node, Inactivity.node],
})
11 changes: 7 additions & 4 deletions packages/core/src/shell.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -255,11 +257,9 @@ const layer = () =>
input: CreateInput,
before?: (input: ShellCreateBefore) => Effect.Effect<void, E, R>,
) {
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,
Expand Down Expand Up @@ -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, {
Expand Down Expand Up @@ -443,6 +445,7 @@ export const node = makeLocationNode({
layer: layer(),
deps: [
Bus.node,
Inactivity.node,
Location.node,
Global.node,
ShellSelect.node,
Expand Down
3 changes: 2 additions & 1 deletion packages/core/test/browser-idle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 { 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"
Expand All @@ -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, Inactivity.node],
}),
),
],
Expand Down
3 changes: 2 additions & 1 deletion packages/core/test/job.test.ts
Original file line number Diff line number Diff line change
@@ -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"
Expand All @@ -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", () =>
Expand Down
Loading
Loading