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
74 changes: 74 additions & 0 deletions packages/core/src/instance-activity.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
import { resolve } from "path"
import { FSUtil } from "./fs-util"

export interface Identity {
readonly directory: string
}

export interface Snapshot {
readonly identity: Identity
readonly generation: number
readonly runningPtys: number
}

interface State {
generation: number
runningPtys: number
}

const states = new Map<string, State>()
let generation = 0

/** Stable for the lifetime of an instance: absolute and lexical, without filesystem lookups. */
export function identify(directory: string): Identity {
return { directory: resolve(FSUtil.windowsPath(directory)) }
}

const change = (identity: Identity, update?: (state: State) => void) => {
const state = states.get(identity.directory) ?? { generation: 0, runningPtys: 0 }
update?.(state)
generation++
state.generation = generation
states.set(identity.directory, state)
}

export function touch(identity: Identity) {
change(identity)
}

export function ptyStarted(identity: Identity) {
change(identity, (state) => state.runningPtys++)
}

export function ptyStopped(identity: Identity) {
const state = states.get(identity.directory)
if (!state || state.runningPtys === 0) return
change(identity, (current) => current.runningPtys--)
}

export function hasRunningPty(identity: Identity) {
return (states.get(identity.directory)?.runningPtys ?? 0) > 0
}

export function snapshot(identity: Identity): Snapshot {
const state = states.get(identity.directory)
return {
identity,
generation: state?.generation ?? 0,
runningPtys: state?.runningPtys ?? 0,
}
}

/** Must be called in the same synchronous call stack as the guarded state transition. */
export function isCurrentAndIdle(value: Snapshot) {
const state = states.get(value.identity.directory)
return (state?.generation ?? 0) === value.generation && (state?.runningPtys ?? 0) === 0
}

export function forget(identity: Identity) {
const state = states.get(identity.directory)
if (!state || state.runningPtys > 0) return
states.delete(identity.directory)
}

export * as InstanceActivity from "./instance-activity"
13 changes: 13 additions & 0 deletions packages/core/src/pty.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import { Location } from "./location"
import { PtyID } from "./pty/schema"
import { Shell } from "./shell"
import { lazy } from "./util/lazy"
import { PtyActivity } from "./pty/activity"

const BUFFER_LIMIT = 1024 * 1024 * 2
// Exited sessions stay observable (status, exit code, retained output) until removed explicitly.
Expand All @@ -34,6 +35,7 @@ type Active = {
cursor: number
subscribers: Map<object, Subscriber>
listeners: Disp[]
counted: boolean
}

export const Info = Pty.Info
Expand Down Expand Up @@ -94,6 +96,7 @@ const layer = Layer.effect(
Effect.gen(function* () {
const events = yield* EventV2.Service
const location = yield* Location.Service
const activity = PtyActivity.identify(location.directory)
const config = yield* Config.Service
const context = yield* Effect.context()
const runFork = Effect.runForkWith(context)
Expand All @@ -117,6 +120,10 @@ const layer = Layer.effect(
for (const listener of session.listeners) listener.dispose()
session.listeners.length = 0
if (session.info.status === "running") {
if (session.counted) {
session.counted = false
PtyActivity.stopped(activity)
}
try {
session.process.kill()
} catch {}
Expand Down Expand Up @@ -198,8 +205,10 @@ const layer = Layer.effect(
cursor: 0,
subscribers: new Map(),
listeners: [],
counted: true,
}
sessions.set(id, session)
PtyActivity.started(activity)
session.listeners.push(
proc.onData((chunk) => {
session.cursor += chunk.length
Expand All @@ -222,6 +231,10 @@ const layer = Layer.effect(
}),
proc.onExit(({ exitCode }) => {
if (session.info.status === "exited") return
if (session.counted) {
session.counted = false
PtyActivity.stopped(activity)
}
session.info.status = "exited"
session.info.exitCode = exitCode
notifyEnd(session, { exitCode })
Expand Down
20 changes: 20 additions & 0 deletions packages/core/src/pty/activity.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
import { InstanceActivity } from "../instance-activity"

export function identify(directory: string) {
return InstanceActivity.identify(directory)
}

export function started(identity: InstanceActivity.Identity) {
InstanceActivity.ptyStarted(identity)
}

export function stopped(identity: InstanceActivity.Identity) {
InstanceActivity.ptyStopped(identity)
}

export function hasRunning(directory: string) {
const identity = identify(directory)
return InstanceActivity.hasRunningPty(identity)
}

export * as PtyActivity from "./activity"
64 changes: 64 additions & 0 deletions packages/core/test/pty/activity.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
import { afterEach, describe, expect, test } from "bun:test"
import { mkdir, mkdtemp, rename, rm, symlink } from "fs/promises"
import { tmpdir } from "os"
import { join } from "path"
import { InstanceActivity } from "../../src/instance-activity"
import { PtyActivity } from "../../src/pty/activity"

const canonical = "/tmp/opencode-pty-activity/b"
const alias = "/tmp/opencode-pty-activity/a/../b"
const canonicalIdentity = PtyActivity.identify(canonical)

afterEach(() => {
PtyActivity.stopped(canonicalIdentity)
InstanceActivity.forget(canonicalIdentity)
})

describe("PtyActivity", () => {
test("canonicalizes lexical directory aliases", () => {
const aliasIdentity = PtyActivity.identify(alias)
const before = InstanceActivity.snapshot(canonicalIdentity).generation
PtyActivity.started(aliasIdentity)
expect(PtyActivity.hasRunning(canonical)).toBe(true)
const running = InstanceActivity.snapshot(canonicalIdentity).generation
expect(running).toBeGreaterThan(before)

PtyActivity.stopped(aliasIdentity)
expect(PtyActivity.hasRunning(alias)).toBe(false)
expect(InstanceActivity.snapshot(canonicalIdentity).generation).toBeGreaterThan(running)
})

test("keeps the owning identity across rename and symlink replacement", async () => {
const root = await mkdtemp(join(tmpdir(), "opencode-pty-retarget-"))
const owner = join(root, "owner")
const moved = join(root, "moved")
const replacement = join(root, "replacement")
await mkdir(owner)
await mkdir(replacement)
const identity = PtyActivity.identify(owner)

try {
PtyActivity.started(identity)
const running = InstanceActivity.snapshot(identity)
expect(running.runningPtys).toBe(1)

await rename(owner, moved)
await symlink(replacement, owner, "dir")

expect(PtyActivity.hasRunning(owner)).toBe(true)
PtyActivity.stopped(identity)
expect(PtyActivity.hasRunning(owner)).toBe(false)
expect(InstanceActivity.snapshot(identity).runningPtys).toBe(0)

PtyActivity.started(identity)
PtyActivity.stopped(identity)
const balanced = InstanceActivity.snapshot(identity)
expect(balanced.runningPtys).toBe(0)
expect(balanced.generation).toBeGreaterThan(running.generation)
} finally {
PtyActivity.stopped(identity)
InstanceActivity.forget(identity)
await rm(root, { recursive: true, force: true })
}
})
})
48 changes: 24 additions & 24 deletions packages/opencode/src/acp/directory.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,5 @@
import { Agent } from "@/agent/agent"
import { Command } from "@/command"
import { InstanceRef } from "@/effect/instance-ref"
import { InstanceBootstrap } from "@/project/bootstrap"
import { InstanceStore } from "@/project/instance-store"
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { ProviderV2 } from "@opencode-ai/core/provider"
Expand Down Expand Up @@ -114,28 +112,30 @@ export const loaderLayer = Layer.effect(

return Loader.of({
load: Effect.fn("ACPDirectoryLoader.load")(function* (directory) {
const ctx = yield* store.load({ directory })
return yield* Effect.gen(function* () {
const providers = yield* provider.list()
const [agents, defaultAgent, commands, defaultModel] = yield* Effect.all(
[agent.list(), agent.defaultInfo(), command.list(), provider.defaultModel().pipe(Effect.option)],
{ concurrency: "unbounded" },
)
return build({
directory,
providers,
modes: agents
.filter((item) => item.mode !== "subagent" && item.hidden !== true)
.map((item) => ({
id: item.name,
name: item.name,
...(item.description ? { description: item.description } : {}),
})),
defaultModeID: defaultAgent.name,
commands: commands.toSorted((a, b) => a.name.localeCompare(b.name)),
...(defaultModel._tag === "Some" ? { defaultModel: defaultModel.value } : {}),
})
}).pipe(Effect.provideService(InstanceRef, ctx))
return yield* store.provide(
{ directory },
Effect.gen(function* () {
const providers = yield* provider.list()
const [agents, defaultAgent, commands, defaultModel] = yield* Effect.all(
[agent.list(), agent.defaultInfo(), command.list(), provider.defaultModel().pipe(Effect.option)],
{ concurrency: "unbounded" },
)
return build({
directory,
providers,
modes: agents
.filter((item) => item.mode !== "subagent" && item.hidden !== true)
.map((item) => ({
id: item.name,
name: item.name,
...(item.description ? { description: item.description } : {}),
})),
defaultModeID: defaultAgent.name,
commands: commands.toSorted((a, b) => a.name.localeCompare(b.name)),
...(defaultModel._tag === "Some" ? { defaultModel: defaultModel.value } : {}),
})
}),
)
}),
})
}),
Expand Down
7 changes: 1 addition & 6 deletions packages/opencode/src/acp/usage.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,5 @@
import type { AgentSideConnection, Usage } from "@agentclientprotocol/sdk"
import type { AssistantMessage as OpenCodeAssistantMessage, Message } from "@opencode-ai/sdk/v2"
import { InstanceRef } from "@/effect/instance-ref"
import { InstanceBootstrap } from "@/project/bootstrap"
import { InstanceStore } from "@/project/instance-store"
import { makeGlobalNode, Node } from "@opencode-ai/core/effect/app-node"
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
Expand Down Expand Up @@ -130,10 +128,7 @@ export const contextLimitLoaderLayer = Layer.effect(

return ContextLimitLoader.of({
providers: Effect.fn("ACPUsageContextLimitLoader.providers")(function* (directory) {
const ctx = yield* store.load({ directory })
return yield* Effect.gen(function* () {
return yield* provider.list()
}).pipe(Effect.provideService(InstanceRef, ctx))
return yield* store.provide({ directory }, provider.list())
}),
})
}),
Expand Down
24 changes: 22 additions & 2 deletions packages/opencode/src/background/job.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { BackgroundJob as CoreBackgroundJob } from "@opencode-ai/core/background-job"
import { InstanceState } from "@/effect/instance-state"
import { Effect, Layer } from "effect"
import { InstanceActivity } from "@opencode-ai/core/instance-activity"

export {
Service,
Expand All @@ -19,11 +20,30 @@ const layer = Layer.effect(
CoreBackgroundJob.Service,
Effect.gen(function* () {
const state = yield* InstanceState.make(() => CoreBackgroundJob.make)
const tracked = (activity: InstanceActivity.Identity, run: Effect.Effect<string, unknown>) =>
Effect.sync(() => InstanceActivity.touch(activity)).pipe(
Effect.andThen(run),
Effect.ensuring(Effect.sync(() => InstanceActivity.touch(activity))),
)
return CoreBackgroundJob.Service.of({
list: () => InstanceState.useEffect(state, (jobs) => jobs.list()),
get: (id) => InstanceState.useEffect(state, (jobs) => jobs.get(id)),
start: (input) => InstanceState.useEffect(state, (jobs) => jobs.start(input)),
extend: (input) => InstanceState.useEffect(state, (jobs) => jobs.extend(input)),
start: (input) =>
Effect.gen(function* () {
const activity = InstanceActivity.identify(yield* InstanceState.directory)
InstanceActivity.touch(activity)
return yield* InstanceState.useEffect(state, (jobs) =>
jobs.start({ ...input, run: tracked(activity, input.run) }),
)
}),
extend: (input) =>
Effect.gen(function* () {
const activity = InstanceActivity.identify(yield* InstanceState.directory)
InstanceActivity.touch(activity)
return yield* InstanceState.useEffect(state, (jobs) =>
jobs.extend({ ...input, run: tracked(activity, input.run) }),
)
}),
wait: (input) => InstanceState.useEffect(state, (jobs) => jobs.wait(input)),
waitForPromotion: (id) => InstanceState.useEffect(state, (jobs) => jobs.waitForPromotion(id)),
promote: (id) => InstanceState.useEffect(state, (jobs) => jobs.promote(id)),
Expand Down
14 changes: 14 additions & 0 deletions packages/opencode/src/effect/runtime-flags.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,18 @@ const positiveInteger = (name: string) =>
Config.map((value) => (Number.isInteger(value) && value > 0 ? value : undefined)),
Config.orElse(() => Config.succeed(undefined)),
)
const instanceIdleTimeout = Config.string("OPENCODE_EXPERIMENTAL_INSTANCE_IDLE_TIMEOUT_MS").pipe(
Config.option,
Config.map((input) => {
if (Option.isNone(input)) return { value: undefined, warning: undefined }
const value = Number(input.value)
if (Number.isInteger(value) && value >= 60 * 60 * 1000) return { value, warning: undefined }
return {
value: undefined,
warning: "ignoring OPENCODE_EXPERIMENTAL_INSTANCE_IDLE_TIMEOUT_MS: expected an integer of at least 3600000 ms",
}
}),
)
const experimental = bool("OPENCODE_EXPERIMENTAL")
const enabledByExperimental = (name: string) =>
Config.all({ experimental, enabled: Config.boolean(name).pipe(Config.option) }).pipe(
Expand Down Expand Up @@ -51,6 +63,8 @@ export class Service extends ConfigService.Service<Service>()("@opencode/Runtime
experimentalIconDiscovery: enabledByExperimental("OPENCODE_EXPERIMENTAL_ICON_DISCOVERY"),
outputTokenMax: positiveInteger("OPENCODE_EXPERIMENTAL_OUTPUT_TOKEN_MAX"),
bashDefaultTimeoutMs: positiveInteger("OPENCODE_EXPERIMENTAL_BASH_DEFAULT_TIMEOUT_MS"),
instanceIdleTimeoutMs: instanceIdleTimeout.pipe(Config.map((input) => input.value)),
instanceIdleTimeoutWarning: instanceIdleTimeout.pipe(Config.map((input) => input.warning)),
experimentalNativeLlm: bool("OPENCODE_EXPERIMENTAL_NATIVE_LLM"),
experimentalWebSockets: bool("OPENCODE_EXPERIMENTAL_WEBSOCKETS"),
client: Config.string("OPENCODE_CLIENT").pipe(Config.withDefault("cli")),
Expand Down
Loading
Loading