Compare commits

..

6 Commits

Author SHA1 Message Date
Kit Langton c61e93a1f3 fix(cli): harden election edge cases 2026-07-08 10:54:28 -04:00
Kit Langton 68c62774ac fix(cli): elect one managed daemon 2026-07-08 10:54:28 -04:00
Dax Raad 8b634e4a58 refactor(tui): simplify client data state 2026-07-08 10:37:48 -04:00
James Long c6156f171c feat(simulation): stream UI recordings (#35909) 2026-07-08 10:33:11 -04:00
Dax Raad 4d2b06f8cf docs: prohibit bypassing git hooks 2026-07-08 10:11:46 -04:00
𝓛𝓲𝓽𝓽𝓵𝓮 𝓕𝓻𝓪𝓷𝓴 bf15c97e4b Revert "fix(tui): paginate session history"
This reverts commit cc29f86cdc.
2026-07-08 13:45:20 +00:00
21 changed files with 645 additions and 565 deletions
+2
View File
@@ -19,6 +19,8 @@ Valid types are `feat`, `fix`, `docs`, `chore`, `refactor`, and `test`. Scopes a
Examples: `fix(tui): simplify thinking toggle styling`, `docs: update contributing guide`, `chore(sdk): regenerate types`.
Never bypass Git hooks. Do not use `--no-verify` or otherwise disable, skip, or circumvent commit or push hooks. If a hook fails, fix the failure or stop and report it to the user.
## Style Guide
### General Principles
+27 -5
View File
@@ -7,6 +7,7 @@ import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { Global } from "@opencode-ai/core/global"
import { InstallationVersion } from "@opencode-ai/core/installation/version"
import { AppProcess } from "@opencode-ai/core/process"
import { EffectFlock } from "@opencode-ai/core/util/effect-flock"
import { start } from "@opencode-ai/server/process"
import { randomBytes, randomUUID } from "node:crypto"
import path from "node:path"
@@ -24,10 +25,13 @@ export type Options = {
readonly port?: number
}
type ManagedServiceOptions = Service.Options & { readonly file: string }
export const run = Effect.fn("cli.server-process.run")((options: Options) =>
processEffect(options).pipe(
Effect.catchTag("ServiceAlreadyOwned", () => Effect.logInfo("another process owns the managed service, exiting")),
Effect.provide(Updater.layer),
Effect.provide(AppNodeBuilder.build(LayerNode.group([Global.node, AppProcess.node]))),
Effect.provide(AppNodeBuilder.build(LayerNode.group([Global.node, AppProcess.node, EffectFlock.node]))),
Effect.provide(NodeServices.layer),
),
)
@@ -43,6 +47,21 @@ const processEffect = Effect.fnUntraced(function* (options: Options) {
delete process.env.OPENCODE_SERVER_PASSWORD
}
const config = options.mode === "service" ? yield* ServiceConfig.read() : {}
const serviceOptions = options.mode === "service" ? yield* ServiceConfig.options() : undefined
if (serviceOptions) {
const flock = yield* EffectFlock.Service
yield* flock.tryAcquire(`opencode-service:${serviceOptions.file}`).pipe(
Effect.filterOrFail(
(acquired) => acquired,
() => ({ _tag: "ServiceAlreadyOwned" }) as const,
),
// Retry only lost elections: a displaced owner may take a moment to release.
Effect.retry({
while: (error) => error._tag === "ServiceAlreadyOwned",
schedule: Schedule.spaced("100 millis").pipe(Schedule.both(Schedule.recurs(20))),
}),
)
}
const password =
options.mode === "service"
? yield* ServiceConfig.password()
@@ -55,7 +74,7 @@ const processEffect = Effect.fnUntraced(function* (options: Options) {
port: Option.fromNullishOr(options.port ?? config.port),
password,
}).pipe(Effect.provide(Logger.layer([], { mergeWithExisting: false })))
if (options.mode === "service") yield* register(address, password)
if (serviceOptions) yield* register(address, password, serviceOptions)
const url = HttpServer.formatAddress(address)
console.log(options.mode === "stdio" ? JSON.stringify({ url }) : `server listening on ${url}`)
if (options.mode === "default" && !environmentPassword) console.log(`server password ${password}`)
@@ -72,9 +91,12 @@ const infoJson = Schema.fromJsonString(Service.Info)
const encodeInfo = Schema.encodeEffect(infoJson)
const decodeInfo = Schema.decodeUnknownEffect(infoJson)
const register = Effect.fnUntraced(function* (address: HttpServer.Address, password: string) {
const register = Effect.fnUntraced(function* (
address: HttpServer.Address,
password: string,
options: ManagedServiceOptions,
) {
const fs = yield* FileSystem.FileSystem
const options = yield* ServiceConfig.options()
const id = randomUUID()
const temp = options.file + "." + id + ".tmp"
yield* fs.makeDirectory(path.dirname(options.file), { recursive: true })
@@ -98,7 +120,7 @@ const register = Effect.fnUntraced(function* (address: HttpServer.Address, passw
? Effect.void
: Effect.try({ try: () => process.kill(process.pid, "SIGTERM"), catch: (cause) => cause }).pipe(Effect.ignore),
),
Effect.repeat(Schedule.spaced("10 seconds")),
Effect.repeat(Schedule.spaced("1 second")),
Effect.forkScoped,
)
yield* Effect.addFinalizer(() =>
+61
View File
@@ -0,0 +1,61 @@
import { expect, test } from "bun:test"
import { spawn } from "node:child_process"
import fs from "node:fs/promises"
import os from "node:os"
import path from "node:path"
test("concurrent service candidates elect one owner before serving", async () => {
const root = await fs.mkdtemp(path.join(os.tmpdir(), "opencode-service-election-"))
const env = {
...process.env,
XDG_CACHE_HOME: path.join(root, "cache"),
XDG_CONFIG_HOME: path.join(root, "config"),
XDG_DATA_HOME: path.join(root, "data"),
XDG_STATE_HOME: path.join(root, "state"),
OPENCODE_DISABLE_AUTOUPDATE: "1",
}
const entry = path.join(import.meta.dir, "../src/index.ts")
const candidates = [
spawn(process.execPath, [entry, "serve", "--service"], { cwd: path.join(import.meta.dir, ".."), env }),
spawn(process.execPath, [entry, "serve", "--service"], { cwd: path.join(import.meta.dir, ".."), env }),
]
try {
const registration = path.join(root, "state", "opencode", "service-local.json")
await waitForFile(registration)
const info = await Bun.file(registration).json()
await waitFor(() => candidates.some((candidate) => candidate.exitCode !== null))
expect(candidates.filter((candidate) => candidate.exitCode === null)).toHaveLength(1)
expect(candidates.find((candidate) => candidate.exitCode !== null)?.exitCode).toBe(0)
expect(candidates.find((candidate) => candidate.exitCode === null)?.pid).toBe(info.pid)
expect(
await fetch(new URL("/api/health", info.url), {
headers: { authorization: "Basic " + btoa(`opencode:${info.password}`) },
}).then((response) => response.ok),
).toBe(true)
} finally {
await Promise.all(candidates.map(stop))
await fs.rm(root, { recursive: true, force: true })
}
}, 20_000)
async function waitForFile(file: string) {
await waitFor(() => Bun.file(file).exists())
}
async function waitFor(check: () => boolean | Promise<boolean>) {
const timeout = Date.now() + 10_000
while (Date.now() < timeout) {
if (await check()) return
await Bun.sleep(20)
}
throw new Error("Timed out waiting for service election")
}
async function stop(process: ReturnType<typeof spawn>) {
if (process.exitCode !== null || process.signalCode !== null) return
const closed = new Promise<void>((resolve) => process.once("close", () => resolve()))
process.kill("SIGTERM")
await closed
}
+53 -12
View File
@@ -66,6 +66,7 @@ export namespace EffectFlock {
)
const decodeMeta = Schema.decodeUnknownSync(LockMetaJson)
const decodeMetaOption = Schema.decodeUnknownOption(LockMetaJson)
const encodeMeta = Schema.encodeSync(LockMetaJson)
// ---------------------------------------------------------------------------
@@ -74,6 +75,7 @@ export namespace EffectFlock {
export interface Interface {
readonly acquire: (key: string, dir?: string) => Effect.Effect<void, LockError, Scope.Scope>
readonly tryAcquire: (key: string, dir?: string) => Effect.Effect<boolean, LockError, Scope.Scope>
readonly withLock: {
(key: string, dir?: string): <A, E, R>(body: Effect.Effect<A, E, R>) => Effect.Effect<A, E | LockError, R>
<A, E, R>(body: Effect.Effect<A, E, R>, key: string, dir?: string): Effect.Effect<A, E | LockError, R>
@@ -150,6 +152,27 @@ export namespace EffectFlock {
const isStale = Effect.fnUntraced(function* (lockDir: string, heartbeatPath: string, metaPath: string) {
const now = wall()
const raw = yield* fs.readFileString(metaPath).pipe(
Effect.map(Option.some),
Effect.catchIf(isPathGone, () => Effect.succeed(Option.none())),
Effect.orDie,
)
const owner = Option.isSome(raw) ? Option.getOrUndefined(decodeMetaOption(raw.value)) : undefined
if (owner?.hostname === hostname) {
const alive = yield* Effect.try({
try: () => {
process.kill(owner.pid, 0)
return true
},
catch: (cause) => (cause && typeof cause === "object" && "code" in cause ? cause.code : undefined),
}).pipe(
// Only ESRCH proves the pid is gone; other errors (e.g. EPERM) mean a
// process exists, so never break a possibly live owner's lock.
Effect.catch((code) => Effect.succeed(code !== "ESRCH")),
)
if (!alive) return true
}
const hb = yield* safeStat(heartbeatPath)
if (hb) return now - mtimeMs(hb) > STALE_MS
@@ -243,28 +266,46 @@ export namespace EffectFlock {
catch: (cause) => new ReleaseError({ detail: "metadata invalid", cause }),
}).pipe(Effect.orDie)
if (parsed.token !== handle.token) return yield* Effect.die(new ReleaseError({ detail: "token mismatch" }))
if (parsed.token !== handle.token) yield* Effect.die(new ReleaseError({ detail: "token mismatch" }))
yield* forceRemove(handle.lockDir)
})
// -- build service --
const acquire = Effect.fn("EffectFlock.acquire")(function* (key: string, dir?: string) {
const lockDir = dir ?? lockRoot
yield* ensureDir(lockDir)
const lockfile = path.join(lockDir, Hash.fast(key) + ".lock")
// acquireRelease: acquire is uninterruptible, release is guaranteed
const handle = yield* Effect.acquireRelease(acquireHandle(lockfile, key), (handle) => release(handle))
// Heartbeat fiber — scoped, so it's interrupted before release runs
const heartbeat = Effect.fnUntraced(function* (handle: Handle) {
yield* fs
.utimes(handle.heartbeatPath, new Date(), new Date())
.pipe(Effect.ignore, Effect.repeat(Schedule.spaced(HEARTBEAT_MS)), Effect.forkScoped)
})
const acquire = Effect.fn("EffectFlock.acquire")(function* (key: string, dir?: string) {
const lockDir = dir ?? lockRoot
yield* ensureDir(lockDir)
const handle = yield* Effect.acquireRelease(
acquireHandle(path.join(lockDir, Hash.fast(key) + ".lock"), key),
(handle) => release(handle),
)
yield* heartbeat(handle)
})
const tryAcquire = Effect.fn("EffectFlock.tryAcquire")(function* (key: string, dir?: string) {
return yield* Effect.uninterruptibleMask(() =>
Effect.gen(function* () {
const lockDir = dir ?? lockRoot
yield* ensureDir(lockDir)
const handle = yield* tryAcquireLockDir(path.join(lockDir, Hash.fast(key) + ".lock"), key).pipe(
Effect.map(Option.some),
Effect.catchTag("NotAcquired", () => Effect.succeed(Option.none())),
)
if (Option.isNone(handle)) return false
yield* Effect.addFinalizer(() => release(handle.value))
yield* heartbeat(handle.value)
return true
}),
)
})
const withLock: Interface["withLock"] = Function.dual(
(args) => Effect.isEffect(args[0]),
<A, E, R>(body: Effect.Effect<A, E, R>, key: string, dir?: string): Effect.Effect<A, E | LockError, R> =>
@@ -276,7 +317,7 @@ export namespace EffectFlock {
),
)
return Service.of({ acquire, withLock })
return Service.of({ acquire, tryAcquire, withLock })
}),
)
+19 -9
View File
@@ -6,7 +6,6 @@ import os from "os"
import { Cause, Effect, Exit } from "effect"
import { testEffect } from "../lib/effect"
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { EffectFlock } from "@opencode-ai/core/util/effect-flock"
import { Global } from "@opencode-ai/core/global"
import { Hash } from "@opencode-ai/core/util/hash"
@@ -134,6 +133,24 @@ describe("util.effect-flock", () => {
}),
)
it.live(
"tries once without waiting for an active owner",
Effect.gen(function* () {
const flock = yield* EffectFlock.Service
const tmp = yield* Effect.promise(() => fs.mkdtemp(path.join(os.tmpdir(), "eflock-test-")))
const dir = path.join(tmp, "locks")
yield* Effect.scoped(
Effect.gen(function* () {
expect(yield* flock.tryAcquire("eflock:try", dir)).toBe(true)
expect(yield* Effect.scoped(flock.tryAcquire("eflock:try", dir))).toBe(false)
}),
)
expect(yield* Effect.scoped(flock.tryAcquire("eflock:try", dir))).toBe(true)
yield* Effect.promise(() => fs.rm(tmp, { recursive: true, force: true }))
}),
)
it.live(
"withLock data-first",
Effect.gen(function* () {
@@ -358,7 +375,7 @@ describe("util.effect-flock", () => {
)
it.live(
"recovers after a crashed lock owner",
"immediately recovers after a local lock owner crashes",
() =>
Effect.promise(async () => {
const tmp = await fs.mkdtemp(path.join(os.tmpdir(), "eflock-crash-"))
@@ -371,13 +388,6 @@ describe("util.effect-flock", () => {
await waitForFile(ready, 5_000)
await stopWorker(proc)
// Backdate lock files so they're past STALE_MS (60s)
const lockDir = lock(dir, "eflock:crash")
const old = new Date(Date.now() - 120_000)
await fs.utimes(lockDir, old, old).catch(() => {})
await fs.utimes(path.join(lockDir, "heartbeat"), old, old).catch(() => {})
await fs.utimes(path.join(lockDir, "meta.json"), old, old).catch(() => {})
const done = path.join(tmp, "done.log")
const result = await run({ key: "eflock:crash", dir, done, holdMs: 10 })
expect(result.code).toBe(0)
+2 -1
View File
@@ -10,7 +10,8 @@
"./backend/*": "./src/backend/*.ts",
"./frontend": "./src/frontend/simulation.ts",
"./frontend/*": "./src/frontend/*.ts",
"./protocol": "./src/protocol/index.ts"
"./protocol": "./src/protocol/index.ts",
"./recording": "./src/recording.ts"
},
"scripts": {
"typecheck": "tsgo --noEmit"
+9 -59
View File
@@ -1,7 +1,7 @@
import { mkdtemp } from "node:fs/promises"
import { mkdir, mkdtemp } from "node:fs/promises"
import { tmpdir } from "node:os"
import { join } from "node:path"
import type { CapturedFrame, CliRenderer, Renderable } from "@opentui/core"
import { dirname, join } from "node:path"
import type { CliRenderer, Renderable } from "@opentui/core"
import { createMockKeys, createMockMouse, type MockInput, type MockMouse } from "@opentui/core/testing"
import type { SimulationProtocol } from "../protocol"
import { SimulationRenderer } from "./renderer"
@@ -104,67 +104,15 @@ export function state(harness: Harness) {
}
}
export async function screenshot(harness: Harness) {
export async function screenshot(harness: Harness, output?: string) {
await harness.renderOnce()
const image = SimulationPng.screenshot(harness.renderer)
const path = join(await mkdtemp(join(tmpdir(), "opencode-drive-")), "screenshot.png")
const path = output ?? join(await mkdtemp(join(tmpdir(), "opencode-drive-")), "screenshot.png")
if (output) await mkdir(dirname(output), { recursive: true })
await Bun.write(path, image.data)
return path
}
export function frame(harness: Harness): CapturedFrame {
const buffer = harness.renderer.currentRenderBuffer
return {
cols: buffer.width,
rows: buffer.height,
cursor: [0, 0],
lines: buffer.getSpanLines().map((line) => ({
spans: line.spans.map((span) => ({
text: span.text,
fg: span.fg,
bg: span.bg,
attributes: span.attributes,
width: span.width,
})),
})),
}
}
export async function video(frames: CapturedFrame[]) {
const directory = await mkdtemp(join(tmpdir(), "opencode-drive-recording-"))
await Promise.all(
frames.map((frame, index) =>
Bun.write(
join(directory, `frame-${index.toString().padStart(6, "0")}.png`),
SimulationPng.screenshotFrame(frame).data,
),
),
)
const path = join(directory, "recording.mp4")
const process = Bun.spawn(
[
"ffmpeg",
"-loglevel",
"error",
"-framerate",
"10",
"-i",
join(directory, "frame-%06d.png"),
"-c:v",
"libx264",
"-pix_fmt",
"yuv420p",
"-movflags",
"+faststart",
"-y",
path,
],
{ stderr: "pipe" },
)
if ((await process.exited) !== 0) throw new Error(`ffmpeg failed: ${await new Response(process.stderr).text()}`)
return path
}
export async function execute(harness: Harness, action: Action) {
switch (action.type) {
case "ui.type":
@@ -180,7 +128,9 @@ export async function execute(harness: Harness, action: Action) {
harness.mockInput.pressArrow(action.direction)
break
case "ui.focus":
all(harness.renderer.root).find((item) => item.num === action.target)?.focus()
all(harness.renderer.root)
.find((item) => item.num === action.target)
?.focus()
break
case "ui.click":
await harness.mockMouse.click(action.x, action.y)
+33 -5
View File
@@ -1,21 +1,43 @@
import type { CliRenderer, CliRendererConfig } from "@opentui/core"
import { createTestRenderer, type TestRendererSetup } from "@opentui/core/testing"
import { Timeline } from "../recording"
const setups = new WeakMap<CliRenderer, TestRendererSetup>()
const recordings = new WeakMap<CliRenderer, Timeline>()
/**
* Creates the headless simulation renderer: a real CliRenderer backed by an
* in-memory screen buffer instead of a terminal. The TestRendererSetup is
* kept module-side (keyed by renderer) so the harness can use the supported
* testing APIs without app code carrying it around.
* Creates a headless renderer with optional recording: a real CliRenderer
* backed by an in-memory screen buffer. The TestRendererSetup is kept
* module-side so the harness can use supported testing APIs without app
* code carrying it around.
*/
export async function create(options: CliRendererConfig): Promise<CliRenderer> {
export async function create(options: CliRendererConfig, path?: string): Promise<CliRenderer> {
if (!path) {
const setup = await createTestRenderer({
...options,
width: 100,
height: 40,
})
setups.set(setup.renderer, setup)
return setup.renderer
}
const recording = await Timeline.create(path, 100, 40)
const setup = await createTestRenderer({
...options,
width: 100,
height: 40,
stdout: recording as unknown as NodeJS.WriteStream,
bufferedOutput: "stdout",
onDestroy: () => {
void recording.finish().catch((error) => process.stderr.write(`Failed to finish UI recording: ${error}\n`))
options.onDestroy?.()
},
}).catch(async (error) => {
await recording.finish().catch(() => undefined)
throw error
})
setups.set(setup.renderer, setup)
recordings.set(setup.renderer, recording)
return setup.renderer
}
@@ -23,4 +45,10 @@ export function setupFor(renderer: CliRenderer): TestRendererSetup | undefined {
return setups.get(renderer)
}
export function finish(renderer: CliRenderer) {
const recording = recordings.get(renderer)
if (!recording) throw new Error("UI recording is not available")
return recording.finish()
}
export * as SimulationRenderer from "./renderer"
+9 -42
View File
@@ -1,4 +1,3 @@
import type { CapturedFrame } from "@opentui/core"
import { SimulationProtocol } from "../protocol"
import { SimulationActions, type Harness } from "./actions"
@@ -11,50 +10,20 @@ function parseRequest(input: string | Buffer) {
return SimulationProtocol.Frontend.decodeRequest(JSON.parse(typeof input === "string" ? input : input.toString()))
}
interface Recording {
readonly frames: CapturedFrame[]
readonly timer: ReturnType<typeof setInterval>
pending: Promise<void>
}
async function handle(
harness: Harness,
request: SimulationProtocol.Frontend.Request,
recording: { current?: Recording },
headless: boolean,
finishRecording?: () => Promise<string>,
) {
switch (request.method) {
case "ui.screenshot":
return SimulationActions.screenshot(harness)
return SimulationActions.screenshot(harness, request.params?.path)
case "ui.state": {
return SimulationActions.state(harness)
}
case "ui.start-record": {
if (recording.current) throw new Error("UI recording is already active")
const frames = [SimulationActions.frame(harness)]
const current: Recording = {
frames,
timer: setInterval(() => {
current.pending = current.pending.then(async () => {
if (headless) await harness.renderOnce()
frames.push(SimulationActions.frame(harness))
})
}, 100),
pending: Promise.resolve(),
}
recording.current = current
return { recording: true }
}
case "ui.end-record": {
if (!recording.current) throw new Error("UI recording is not active")
const current = recording.current
clearInterval(current.timer)
await current.pending
if (headless) await harness.renderOnce()
current.frames.push(SimulationActions.frame(harness))
recording.current = undefined
return SimulationActions.video(current.frames)
}
case "ui.recording.finish":
if (!finishRecording) throw new Error("UI recording is not available")
return finishRecording()
case "ui.type":
return SimulationActions.execute(harness, { type: "ui.type", text: request.params.text })
case "ui.enter":
@@ -79,14 +48,13 @@ async function handle(
}
}
export function start(harness: Harness, endpoint: string, headless: boolean): Server {
export function start(harness: Harness, endpoint: string, finishRecording?: () => Promise<string>): Server {
const url = new URL(endpoint)
const recording: { current?: Recording } = {}
const server = Bun.serve<{ readonly drive: true; readonly headless: boolean }>({
const server = Bun.serve<{ readonly drive: true }>({
hostname: url.hostname,
port: Number(url.port),
fetch(request, server) {
if (server.upgrade(request, { data: { drive: true, headless } })) return undefined
if (server.upgrade(request, { data: { drive: true } })) return undefined
return new Response("opencode drive ui websocket", { status: 426 })
},
websocket: {
@@ -94,7 +62,7 @@ export function start(harness: Harness, endpoint: string, headless: boolean): Se
let request: SimulationProtocol.Frontend.Request | undefined
try {
request = parseRequest(message)
const result = await handle(harness, request, recording, headless)
const result = await handle(harness, request, finishRecording)
const next = SimulationProtocol.JsonRpc.success(request.id, result)
if (next) socket.send(JSON.stringify(next))
} catch (error) {
@@ -106,7 +74,6 @@ export function start(harness: Harness, endpoint: string, headless: boolean): Se
return {
url: endpoint,
stop: () => {
if (recording.current) clearInterval(recording.current.timer)
server.stop(true)
},
}
@@ -14,11 +14,14 @@ import { SimulationServer } from "./server"
*/
export async function create(options: CliRendererConfig): Promise<CliRenderer> {
const headless = process.env.OPENCODE_DRIVE_RENDERER === "headless"
const renderer = headless ? await SimulationRenderer.create(options) : await createCliRenderer(options)
const manifest = DriveManifest.resolve()
const renderer = headless
? await SimulationRenderer.create(options, manifest.recording?.timeline)
: await createCliRenderer(options)
const server = SimulationServer.start(
SimulationActions.createHarness(renderer),
DriveManifest.resolve().endpoints.ui,
headless,
manifest.endpoints.ui,
headless && manifest.recording ? () => SimulationRenderer.finish(renderer) : undefined,
)
process.stderr.write(`opencode drive ui websocket: ${server.url}\n`)
renderer.once("destroy", () => server.stop())
+10 -9
View File
@@ -1,12 +1,15 @@
import { existsSync, readFileSync } from "node:fs"
import { homedir } from "node:os"
import { join } from "node:path"
import { isAbsolute, join } from "node:path"
export interface Manifest {
readonly endpoints: {
readonly ui: string
readonly backend: string
}
readonly recording?: {
readonly timeline: string
}
}
export const defaults: Manifest = {
@@ -32,18 +35,16 @@ export function resolve() {
if (!isManifest(manifest)) throw new Error(`Invalid drive manifest: ${file}`)
validateEndpoint(manifest.endpoints.ui, "ui")
validateEndpoint(manifest.endpoints.backend, "backend")
if (manifest.recording && !isAbsolute(manifest.recording.timeline)) {
throw new Error(`Invalid drive recording timeline path: ${manifest.recording.timeline}`)
}
return manifest
}
function isManifest(value: unknown): value is Manifest {
if (typeof value !== "object" || value === null) return false
if (!("endpoints" in value) || typeof value.endpoints !== "object" || value.endpoints === null) return false
return (
"ui" in value.endpoints &&
typeof value.endpoints.ui === "string" &&
"backend" in value.endpoints &&
typeof value.endpoints.backend === "string"
)
if (typeof value !== "object" || value === null || !("endpoints" in value)) return false
if (typeof value.endpoints !== "object" || value.endpoints === null) return false
return "ui" in value.endpoints && "backend" in value.endpoints
}
function validateEndpoint(value: string, name: string) {
+10 -13
View File
@@ -94,11 +94,11 @@ export namespace Frontend {
export const Screenshot = Schema.String
export type Screenshot = Schema.Schema.Type<typeof Screenshot>
export const StartRecord = Schema.Struct({ recording: Schema.Literal(true) })
export interface StartRecord extends Schema.Schema.Type<typeof StartRecord> {}
export const RecordingFinish = Schema.String
export type RecordingFinish = Schema.Schema.Type<typeof RecordingFinish>
export const EndRecord = Schema.String
export type EndRecord = Schema.Schema.Type<typeof EndRecord>
export const ArtifactParams = Schema.Struct({ path: Schema.optional(Schema.String) })
export interface ArtifactParams extends Schema.Schema.Type<typeof ArtifactParams> {}
export const TypeParams = Schema.Struct({ text: Schema.String })
export interface TypeParams extends Schema.Schema.Type<typeof TypeParams> {}
@@ -123,18 +123,16 @@ export namespace Frontend {
Schema.Struct({ ...JsonRpc.RequestFields, method: Schema.Literal("ui.click"), params: ClickParams }),
Schema.Struct({
...JsonRpc.RequestFields,
method: Schema.Literals([
"ui.enter",
"ui.screenshot",
"ui.state",
"ui.start-record",
"ui.end-record",
]),
method: Schema.Literal("ui.screenshot"),
params: Schema.optional(ArtifactParams),
}),
Schema.Struct({
...JsonRpc.RequestFields,
method: Schema.Literals(["ui.enter", "ui.state", "ui.recording.finish"]),
}),
])
export type Request = Schema.Schema.Type<typeof Request>
export const decodeRequest = Schema.decodeUnknownSync(Request)
}
export namespace Backend {
@@ -189,7 +187,6 @@ export namespace Backend {
matched: Schema.Boolean,
})
export interface NetworkLogEntry extends Schema.Schema.Type<typeof NetworkLogEntry> {}
}
export * as SimulationProtocol from "./index"
+115
View File
@@ -0,0 +1,115 @@
import { createWriteStream, type WriteStream } from "node:fs"
import { mkdir } from "node:fs/promises"
import { dirname } from "node:path"
import { Writable } from "node:stream"
import { finished } from "node:stream/promises"
import { Schema } from "effect"
export const Header = Schema.Struct({
type: Schema.Literal("header"),
version: Schema.Literal(1),
cols: Schema.Number,
rows: Schema.Number,
encoding: Schema.Literal("base64"),
})
export interface Header extends Schema.Schema.Type<typeof Header> {}
export const Output = Schema.Struct({
type: Schema.Literal("output"),
at_ms: Schema.Number,
data: Schema.String,
})
export interface Output extends Schema.Schema.Type<typeof Output> {}
export const Event = Schema.Union([Header, Output])
export type Event = Schema.Schema.Type<typeof Event>
export class Timeline extends Writable {
readonly isTTY = true
readonly path: string
readonly columns: number
readonly rows: number
private readonly output: WriteStream
private readonly started = performance.now()
private readonly timestamps: number[] = []
private done?: Promise<string>
private constructor(path: string, cols: number, rows: number, output: WriteStream) {
super()
this.path = path
this.columns = cols
this.rows = rows
this.output = output
// finish() reports stream failures; keep Writable from also throwing them process-wide.
this.on("error", () => {})
output.on("error", (error) => this.destroy(error))
}
static async create(path: string, cols: number, rows: number) {
await mkdir(dirname(path), { recursive: true })
const output = createWriteStream(path)
const timeline = new Timeline(path, cols, rows, output)
await new Promise<void>((resolve, reject) => {
output.write(
`${JSON.stringify({ type: "header", version: 1, cols, rows, encoding: "base64" } satisfies Header)}\n`,
(error) => (error ? reject(error) : resolve()),
)
})
return timeline
}
getColorDepth() {
return 24
}
override write(chunk: unknown, callback?: (error?: Error | null) => void): boolean
override write(chunk: unknown, encoding: BufferEncoding, callback?: (error?: Error | null) => void): boolean
override write(
chunk: unknown,
encoding?: BufferEncoding | ((error?: Error | null) => void),
callback?: (error?: Error | null) => void,
) {
if (!this.writableEnded) {
this.timestamps.push(this.elapsed())
if (typeof encoding === "function") return super.write(chunk, encoding)
if (encoding === undefined) return super.write(chunk, callback)
return super.write(chunk, encoding, callback)
}
const done = typeof encoding === "function" ? encoding : callback
queueMicrotask(() => done?.(null))
return true
}
override _write(chunk: Buffer, _encoding: BufferEncoding, callback: (error?: Error | null) => void) {
this.writeOutput(chunk, this.timestamps.shift() ?? this.elapsed(), callback)
}
override _final(callback: (error?: Error | null) => void) {
this.writeOutput(Buffer.alloc(0), this.elapsed(), (error) => {
if (error) return callback(error)
this.output.end(callback)
})
}
finish() {
if (this.done) return this.done
this.end()
this.done = finished(this).then(() => this.path)
return this.done
}
private elapsed() {
return Math.max(0, Math.round(performance.now() - this.started))
}
private writeOutput(data: Buffer, at_ms: number, callback: (error?: Error | null) => void) {
const event = {
type: "output",
at_ms,
data: data.toString("base64"),
} satisfies Output
this.output.write(`${JSON.stringify(event)}\n`, callback)
}
}
export * as SimulationRecording from "./recording"
@@ -0,0 +1,55 @@
import { expect, test } from "bun:test"
import { mkdtemp, rm } from "node:fs/promises"
import { tmpdir } from "node:os"
import { join } from "node:path"
import { SimulationRenderer } from "../src/frontend/renderer"
import { Timeline, type Event } from "../src/recording"
test("streams ANSI chunks into a versioned JSONL timeline", async () => {
const directory = await mkdtemp(join(tmpdir(), "simulation-recording-"))
const path = join(directory, "nested", "timeline.jsonl")
try {
const timeline = await Timeline.create(path, 80, 24)
await new Promise<void>((resolve, reject) => {
timeline.write(Buffer.from("\u001b[2Jhello"), (error) => (error ? reject(error) : resolve()))
})
expect(await timeline.finish()).toBe(path)
await new Promise<void>((resolve) => timeline.write(Buffer.from("ignored"), () => resolve()))
const events = (await Bun.file(path).text())
.trim()
.split("\n")
.map((line) => JSON.parse(line) as Event)
expect(events[0]).toEqual({ type: "header", version: 1, cols: 80, rows: 24, encoding: "base64" })
expect(events[1]?.type).toBe("output")
if (events[1]?.type !== "output") throw new Error("Missing output event")
expect(Buffer.from(events[1].data, "base64").toString()).toBe("\u001b[2Jhello")
expect(events[1].at_ms).toBeGreaterThanOrEqual(0)
expect(events.at(-1)).toMatchObject({ type: "output", data: "" })
} finally {
await rm(directory, { recursive: true, force: true })
}
})
test("captures native renderer output and finishes on destroy", async () => {
const directory = await mkdtemp(join(tmpdir(), "simulation-renderer-recording-"))
const path = join(directory, "timeline.jsonl")
const renderer = await SimulationRenderer.create({}, path)
try {
await SimulationRenderer.setupFor(renderer)?.renderOnce()
renderer.destroy()
expect(await SimulationRenderer.finish(renderer)).toBe(path)
const events = (await Bun.file(path).text())
.trim()
.split("\n")
.map((line) => JSON.parse(line) as Event)
expect(events.some((event) => event.type === "output")).toBe(true)
} finally {
if (!renderer.isDestroyed) renderer.destroy()
await SimulationRenderer.finish(renderer)
await rm(directory, { recursive: true, force: true })
}
})
+5 -7
View File
@@ -1146,11 +1146,8 @@ function App(props: { onSnapshot?: () => Promise<string[]>; pluginHost: TuiPlugi
return render({ params: route.data.data })
})
// Suppress the full-screen reconnecting overlay for transient disconnects (initial startup, host
// reload, sub-second event-stream blips). After the first successful connect, show it only once the
// connection has been lost for a full second; before the first connect give a longer grace period so
// startup never flashes it, but a server that dies before ever connecting still surfaces instead of
// leaving a silent empty app. Hide it immediately the moment status leaves "connecting".
// Suppress the full-screen overlay for transient startup and event-stream retry states.
// Initial connection gets a longer grace period; retries surface more quickly.
const [showReconnecting, setShowReconnecting] = createSignal(false)
let reconnectTimer: ReturnType<typeof setTimeout> | undefined
createEffect(() => {
@@ -1158,7 +1155,8 @@ function App(props: { onSnapshot?: () => Promise<string[]>; pluginHost: TuiPlugi
clearTimeout(reconnectTimer)
reconnectTimer = undefined
}
if (sdk.connection.status() !== "connecting") {
const status = sdk.connection.status()
if (status === "connected") {
setShowReconnecting(false)
return
}
@@ -1167,7 +1165,7 @@ function App(props: { onSnapshot?: () => Promise<string[]>; pluginHost: TuiPlugi
reconnectTimer = undefined
setShowReconnecting(true)
},
sdk.connection.connectedOnce() ? 1000 : 5000,
status === "reconnecting" ? 1000 : 5000,
).unref()
})
onCleanup(() => {
+75 -151
View File
@@ -4,74 +4,62 @@
// Reconnect may re-bootstrap; that is enough. UI and the server own ordering concerns.
import type {
AgentV2Info,
CommandV2Info,
AgentInfo,
CommandInfo,
FormFormInfo,
FormUrlInfo,
IntegrationInfo,
LocationRef,
McpServer,
ModelV2Info,
ModelInfo,
PermissionSavedInfo,
PermissionV2Request,
ProviderV2Info,
ReferenceInfo,
SessionMessage,
SessionMessageInfo,
SessionMessageAssistant,
SessionMessageAssistantReasoning,
SessionMessageAssistantText,
SessionMessageAssistantTool,
SessionV2Info,
SessionInfo,
Shell,
SkillV2Info,
SkillInfo,
V2Event,
} from "@opencode-ai/sdk/v2"
import { createStore, produce, reconcile } from "solid-js/store"
import { createSimpleContext } from "./helper"
import { useSDK } from "./sdk"
import { batch, createSignal, onCleanup } from "solid-js"
import { createSignal, onCleanup } from "solid-js"
export type DataSessionStatus = "idle" | "running"
const messageIDFromEvent = (eventID: string) => eventID.replace(/^evt_/, "msg_")
const MESSAGE_PAGE_SIZE = 25
export type FormInfo = FormFormInfo | FormUrlInfo
// Per-session message timeline plus older-history paging. `items` is ascending;
// `cursor` is the opaque server cursor for the next older page after a desc first load.
type SessionMessages = {
items: SessionMessage[]
cursor?: string
complete: boolean
loading: boolean
}
type LocationData = {
agent?: AgentV2Info[]
command?: CommandV2Info[]
agent?: AgentInfo[]
command?: CommandInfo[]
integration?: IntegrationInfo[]
mcp?: McpServer[]
model?: ModelV2Info[]
model?: ModelInfo[]
provider?: ProviderV2Info[]
reference?: ReferenceInfo[]
// Currently running shell commands for this location, keyed by shell id. Entries are removed
// once the command exits or is deleted, so this only ever holds in-flight shells.
shell?: Record<string, Shell>
skill?: SkillV2Info[]
skill?: SkillInfo[]
}
type Data = {
session: {
info: Record<string, SessionV2Info>
info: Record<string, SessionInfo>
// Family index keyed by a family's root (or furthest-known-ancestor when the
// true root is not yet loaded). The value is a flat deduplicated list of every
// session ID in that family, including the key itself once its info arrives.
family: Record<string, string[]>
status: Record<string, DataSessionStatus>
compaction: Partial<Record<string, string>>
compactionReason: Partial<Record<string, "auto" | "manual">>
message: Record<string, SessionMessages>
message: Record<string, SessionMessageInfo[]>
input: Record<string, string[]>
permission: Record<string, PermissionV2Request[]>
// Pending forms keyed by session ID.
@@ -83,10 +71,6 @@ type Data = {
location: Record<string, LocationData>
}
function emptyMessages(): SessionMessages {
return { items: [], complete: false, loading: false }
}
function locationKey(location: LocationRef) {
return JSON.stringify([location.directory, location.workspaceID])
}
@@ -111,8 +95,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
info: {},
family: {},
status: {},
compaction: {},
compactionReason: {},
message: {},
input: {},
permission: {},
@@ -136,37 +118,35 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
}
const message = {
update(sessionID: string, fn: (messages: SessionMessage[], index: Map<string, number>) => void) {
update(sessionID: string, fn: (messages: SessionMessageInfo[], index: Map<string, number>) => void) {
setStore(
"session",
"message",
produce((draft) => {
fn((draft[sessionID] ??= emptyMessages()).items, index(sessionID))
fn((draft[sessionID] ??= []), index(sessionID))
}),
)
},
append(messages: SessionMessage[], index: Map<string, number>, item: SessionMessage) {
append(messages: SessionMessageInfo[], index: Map<string, number>, item: SessionMessageInfo) {
if (index.has(item.id)) return
index.set(item.id, messages.length)
messages.push(item)
},
activeAssistant(messages: SessionMessage[]) {
activeAssistant(messages: SessionMessageInfo[]) {
const item = messages.findLast((item) => item.type === "assistant" && !item.time.completed)
return item?.type === "assistant" ? item : undefined
},
assistant(messages: SessionMessage[], index: Map<string, number>, messageID: string) {
assistant(messages: SessionMessageInfo[], index: Map<string, number>, messageID: string) {
const position = index.get(messageID)
const item = position === undefined ? undefined : messages[position]
return item?.type === "assistant" ? item : undefined
},
shell(messages: SessionMessage[], shellID: string) {
const item = messages.findLast((item) => item.type === "shell" && item.shell.id === shellID)
shell(messages: SessionMessageInfo[], shellID: string) {
const item = messages.findLast((item) => item.type === "shell" && item.shellID === shellID)
return item?.type === "shell" ? item : undefined
},
compaction(messages: SessionMessage[]) {
const item = messages.findLast(
(item) => item.type === "compaction" && (item.status === "queued" || item.status === "running"),
)
compaction(messages: SessionMessageInfo[]) {
const item = messages.findLast((item) => item.type === "compaction" && item.status === "running")
return item?.type === "compaction" ? item : undefined
},
latestTool(assistant: SessionMessageAssistant | undefined, callID?: string) {
@@ -244,7 +224,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
produce((draft) => {
delete draft.info[sessionID]
delete draft.status[sessionID]
delete draft.compaction[sessionID]
delete draft.message[sessionID]
delete draft.input[sessionID]
delete draft.permission[sessionID]
@@ -378,6 +357,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
id: messageIDFromEvent(event.id),
type: "system",
text: event.data.text,
metadata: event.metadata,
time: { created: event.created },
})
})
@@ -387,9 +367,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
message.append(draft, index, {
id: messageIDFromEvent(event.id),
type: "synthetic",
sessionID: event.data.sessionID,
text: event.data.text,
description: event.data.description,
metadata: event.data.metadata,
time: { created: event.created },
})
})
@@ -399,7 +379,11 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
message.append(draft, index, {
id: messageIDFromEvent(event.id),
type: "shell",
shell: event.data.shell,
shellID: event.data.shell.id,
command: event.data.shell.command,
status: event.data.shell.status,
exit: event.data.shell.exit,
metadata: event.metadata,
time: { created: event.created },
})
})
@@ -408,7 +392,8 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
message.update(event.data.sessionID, (draft) => {
const match = message.shell(draft, event.data.shell.id)
if (!match) return
match.shell = event.data.shell
match.status = event.data.shell.status
match.exit = event.data.shell.exit
match.output = event.data.output
match.time.completed = event.created
})
@@ -437,6 +422,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
type: "assistant",
agent: event.data.agent,
model: event.data.model,
metadata: event.metadata,
content: [],
snapshot: event.data.snapshot ? { start: event.data.snapshot } : undefined,
time: { created: event.created },
@@ -623,29 +609,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
setSessionStatus(event.data.sessionID, "running")
break
case "session.compaction.admitted":
message.update(event.data.sessionID, (draft, index) => {
if (message.compaction(draft)) return
message.append(draft, index, {
id: event.data.inputID,
type: "compaction",
status: "queued",
reason: "manual",
summary: "",
recent: "",
time: { created: event.created },
})
})
break
case "session.compaction.started":
setStore("session", "compaction", event.data.sessionID, "")
setStore("session", "compactionReason", event.data.sessionID, event.data.reason)
message.update(event.data.sessionID, (draft, index) => {
const current = message.compaction(draft)
if (current) {
current.status = "running"
current.reason = event.data.reason
return
}
message.append(draft, index, {
id: event.data.inputID ?? messageIDFromEvent(event.id),
type: "compaction",
@@ -661,10 +627,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
case "session.execution.failed":
case "session.execution.interrupted":
setSessionStatus(event.data.sessionID, "idle")
if (store.session.compaction[event.data.sessionID] !== undefined)
setStore("session", "compaction", event.data.sessionID, undefined)
if (store.session.compactionReason[event.data.sessionID] !== undefined)
setStore("session", "compactionReason", event.data.sessionID, undefined)
message.update(event.data.sessionID, (draft) => {
const currentAssistant = message.activeAssistant(draft)
if (currentAssistant) currentAssistant.retry = undefined
@@ -695,22 +657,26 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
})
break
case "session.compaction.delta":
setStore("session", "compaction", event.data.sessionID, (text) => (text ?? "") + event.data.text)
message.update(event.data.sessionID, (draft) => {
const current = message.compaction(draft)
if (current) current.summary += event.data.text
if (current?.status === "running") current.summary += event.data.text
})
break
case "session.compaction.ended":
setStore("session", "compaction", event.data.sessionID, undefined)
setStore("session", "compactionReason", event.data.sessionID, undefined)
message.update(event.data.sessionID, (draft, index) => {
const current = message.compaction(draft)
if (current) {
current.status = "completed"
current.reason = event.data.reason
current.summary = event.data.text
current.recent = event.data.recent
const position = draft.findLastIndex((item) => item.type === "compaction" && item.status === "running")
const current = draft[position]
if (current?.type === "compaction") {
draft[position] = {
id: current.id,
type: "compaction",
status: "completed",
reason: event.data.reason,
summary: event.data.text,
recent: event.data.recent,
metadata: current.metadata,
time: current.time,
}
return
}
message.append(draft, index, {
@@ -725,11 +691,26 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
})
break
case "session.compaction.failed":
setStore("session", "compaction", event.data.sessionID, undefined)
setStore("session", "compactionReason", event.data.sessionID, undefined)
message.update(event.data.sessionID, (draft) => {
const current = message.compaction(draft)
if (current) current.status = "failed"
message.update(event.data.sessionID, (draft, index) => {
const position = draft.findLastIndex((item) => item.type === "compaction" && item.status === "running")
const current = draft[position]
const failed: Extract<SessionMessageInfo, { type: "compaction"; status: "failed" }> = {
id: current?.id ?? event.data.inputID ?? messageIDFromEvent(event.id),
type: "compaction",
status: "failed",
reason: event.data.reason ?? "manual",
error: event.data.error ?? {
type: "compaction.failed",
message: "Compaction failed before recording an error",
},
metadata: current?.type === "compaction" ? current.metadata : event.metadata,
time: current?.type === "compaction" ? current.time : { created: event.created },
}
if (current?.type === "compaction") {
draft[position] = failed
return
}
message.append(draft, index, failed)
})
break
case "permission.v2.asked":
@@ -809,20 +790,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
const result = {
on: sdk.event.on,
listen: sdk.event.listen,
connection: {
status() {
return sdk.connection.status()
},
attempt() {
return sdk.connection.attempt()
},
error() {
return sdk.connection.error()
},
connectedOnce() {
return sdk.connection.connectedOnce()
},
},
session: {
list() {
return Object.values(store.session.info).toSorted((a, b) => b.time.updated - a.time.updated)
@@ -847,73 +814,30 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
return store.session.input[sessionID]?.includes(inputID) ?? false
},
},
compaction(sessionID: string) {
return store.session.compaction[sessionID]
},
async refresh(sessionID: string) {
setStore("session", "info", sessionID, mutable(await sdk.api.session.get({ sessionID })))
registerSession(sessionID)
},
message: {
ids(sessionID: string) {
return (store.session.message[sessionID]?.items ?? []).map((message) => message.id)
return (store.session.message[sessionID] ?? []).map((message) => message.id)
},
list(sessionID: string) {
return store.session.message[sessionID]?.items ?? []
return store.session.message[sessionID] ?? []
},
get(sessionID: string, messageID: string) {
const messages = store.session.message[sessionID]?.items
const messages = store.session.message[sessionID]
const position = messageIndex.get(sessionID)?.get(messageID)
return position === undefined ? undefined : messages?.[position]
},
cursor(sessionID: string) {
return store.session.message[sessionID]?.cursor
},
complete(sessionID: string) {
return store.session.message[sessionID]?.complete ?? false
},
loading(sessionID: string) {
return store.session.message[sessionID]?.loading ?? false
},
async refresh(sessionID: string) {
setStore("session", "message", sessionID, { ...emptyMessages(), loading: true })
setStore("session", "message", sessionID, [])
messageIndex.set(sessionID, new Map())
const response = await sdk.api.message.list({ sessionID, limit: MESSAGE_PAGE_SIZE, order: "desc" })
const items = mutable(response.data).toReversed()
messageIndex.set(sessionID, new Map(items.map((message, index) => [message.id, index])))
setStore("session", "message", sessionID, {
items,
cursor: response.cursor.next ?? undefined,
complete: response.data.length < MESSAGE_PAGE_SIZE,
loading: false,
})
const running = items.find((message) => message.type === "compaction" && message.status === "running")
setStore("session", "compaction", sessionID, running?.type === "compaction" ? running.summary : undefined)
setStore(
"session",
"compactionReason",
sessionID,
running?.type === "compaction" ? running.reason : undefined,
)
},
async more(sessionID: string) {
const current = store.session.message[sessionID]
if (!current || current.loading || current.complete || !current.cursor) return
const cursor = current.cursor
setStore("session", "message", sessionID, "loading", true)
const response = await sdk.api.message.list({ sessionID, limit: MESSAGE_PAGE_SIZE, cursor })
const older = mutable(response.data).toReversed()
const prepend = older.filter((item) => !messageIndex.get(sessionID)?.has(item.id))
const items = [...prepend, ...current.items]
messageIndex.set(sessionID, new Map(items.map((item, position) => [item.id, position])))
batch(() => {
setStore("session", "message", sessionID, "items", items)
setStore("session", "message", sessionID, {
cursor: response.cursor.next ?? undefined,
complete: response.data.length < MESSAGE_PAGE_SIZE,
loading: false,
})
})
const messages = mutable(
(await sdk.api.message.list({ sessionID, limit: 200, order: "desc" })).data,
).toReversed()
messageIndex.set(sessionID, new Map(messages.map((message, index) => [message.id, index])))
setStore("session", "message", sessionID, messages)
},
},
permission: {
+3 -8
View File
@@ -5,7 +5,7 @@ import { onCleanup, onMount } from "solid-js"
import { createStore } from "solid-js/store"
import { createSimpleContext } from "./helper"
export type SDKConnectionStatus = "connected" | "connecting"
export type SDKConnectionStatus = "connected" | "connecting" | "reconnecting"
type SDKEventMap = { [Type in V2Event["type"]]: Extract<V2Event, { type: Type }> }
const connectTimeout = 2_000
@@ -27,11 +27,9 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
status: SDKConnectionStatus
attempt: number
error?: string
connectedOnce: boolean
}>({
status: "connecting",
attempt: 0,
connectedOnce: false,
})
let stream: AbortController | undefined
@@ -70,7 +68,7 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
clearTimeout(timeout)
attempt = 0
events.emit(first.value.type, first.value)
setConnection({ status: "connected", attempt: 0, error: undefined, connectedOnce: true })
setConnection({ status: "connected", attempt: 0, error: undefined })
connected()
while (!abort.signal.aborted && !controller.signal.aborted) {
const event = await iterator.next()
@@ -98,7 +96,7 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
}
}
setConnection({
status: "connecting",
status: "reconnecting",
attempt,
error: error instanceof Error ? error.message : String(error),
})
@@ -136,9 +134,6 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
error() {
return connection.error
},
connectedOnce() {
return connection.connectedOnce
},
},
reload: props.reload,
}
+30 -81
View File
@@ -198,16 +198,6 @@ export function Session() {
?.id
})
// Admitted inputs and in-flight manual compaction sit after history, not in rows.
const pendingMessages = createMemo(() => {
const boundary = session()?.revert?.messageID
return messages().filter((message) => {
if (boundary && message.id >= boundary) return false
if (data.session.input.has(route.sessionID, message.id)) return true
return message.type === "compaction" && (message.status === "queued" || message.status === "running")
})
})
const lastAssistant = createMemo(() => {
return messages().findLast((x) => x.type === "assistant")
})
@@ -251,44 +241,37 @@ export function Session() {
}),
)
createEffect(
on(
() => route.sessionID,
(sessionID) => {
void (async () => {
if (data.session.message.list(sessionID).length === 0) {
await Promise.all([
data.session.refresh(sessionID),
data.session.message.refresh(sessionID),
data.session.permission.refresh(sessionID),
data.session.form.refresh(sessionID),
])
}
const info = data.session.get(sessionID)
if (!info) {
toast.show({
message: `Session not found: ${sessionID}`,
variant: "error",
duration: 5000,
})
navigate({ type: "home" })
return
}
project.workspace.set(info.location.workspaceID)
editor.reconnect(info.location.directory)
if (route.sessionID === sessionID && scroll) scroll.scrollBy(100_000)
})().catch((error) => {
if (route.sessionID !== sessionID) return
toast.show({
message: errorMessage(error),
variant: "error",
duration: 5000,
})
navigate({ type: "home" })
createEffect(() => {
const sessionID = route.sessionID
void (async () => {
await Promise.all([
data.session.refresh(sessionID),
data.session.permission.refresh(sessionID),
data.session.form.refresh(sessionID),
])
const info = data.session.get(sessionID)
if (!info) {
toast.show({
message: `Session not found: ${sessionID}`,
variant: "error",
duration: 5000,
})
},
),
)
navigate({ type: "home" })
return
}
project.workspace.set(info.location.workspaceID)
editor.reconnect(info.location.directory)
if (route.sessionID === sessionID && scroll) scroll.scrollBy(100_000)
})().catch((error) => {
if (route.sessionID !== sessionID) return
toast.show({
message: errorMessage(error),
variant: "error",
duration: 5000,
})
navigate({ type: "home" })
})
})
let seeded = false
let scroll: ScrollBoxRenderable
@@ -344,14 +327,12 @@ export function Session() {
if (!targetID) {
scroll.scrollBy(direction === "next" ? scroll.height : -scroll.height)
if (direction === "prev") loadOlder()
dialog.clear()
return
}
const child = scroll.getChildren().find((c) => c.id === targetID)
if (child) scroll.scrollBy(child.y - scroll.y - 1)
if (direction === "prev") loadOlder()
dialog.clear()
}
@@ -362,24 +343,6 @@ export function Session() {
}, 50)
}
let loadingOlder = false
function loadOlder() {
if (loadingOlder || scroll.scrollTop > 2) return
loadingOlder = true
const before = scroll.scrollHeight
void data.session.message.more(route.sessionID).then(
() => {
setTimeout(() => {
if (!scroll.isDestroyed) scroll.scrollBy(scroll.scrollHeight - before)
loadingOlder = false
}, 50)
},
() => {
loadingOlder = false
},
)
}
const sessionCommandList = createMemo(() => [
{
title: "Share session",
@@ -585,7 +548,6 @@ export function Session() {
hidden: true,
run: () => {
scroll.scrollBy(-scroll.height / 2)
loadOlder()
dialog.clear()
},
},
@@ -606,7 +568,6 @@ export function Session() {
hidden: true,
run: () => {
scroll.scrollBy(-1)
loadOlder()
dialog.clear()
},
},
@@ -627,7 +588,6 @@ export function Session() {
hidden: true,
run: () => {
scroll.scrollBy(-scroll.height / 4)
loadOlder()
dialog.clear()
},
},
@@ -648,7 +608,6 @@ export function Session() {
hidden: true,
run: () => {
scroll.scrollTo(0)
loadOlder()
dialog.clear()
},
},
@@ -958,9 +917,6 @@ export function Session() {
stickyStart="bottom"
flexGrow={1}
scrollAcceleration={scrollAcceleration()}
onMouseScroll={(event) => {
if (event.scroll?.direction === "up") void loadOlder()
}}
>
<For each={rows}>
{(row) => (
@@ -970,13 +926,6 @@ export function Session() {
/>
)}
</For>
<For each={pendingMessages()}>
{(message) => (
<box marginTop={1} flexShrink={0}>
<SessionMessageView message={message} />
</box>
)}
</For>
<BackgroundToolHint messages={messages()} />
<Show when={session()?.revert?.messageID}>
<RevertMessage
+80 -46
View File
@@ -1,6 +1,6 @@
import type { SessionMessage, SessionMessageAssistant } from "@opencode-ai/sdk/v2"
import type { SessionMessageAssistant, SessionMessageInfo } from "@opencode-ai/sdk/v2"
import { createEffect, on, onCleanup, type Accessor } from "solid-js"
import { createStore, produce } from "solid-js/store"
import { createStore, produce, reconcile } from "solid-js/store"
import { useData } from "../../context/data"
export type PartRef = {
@@ -25,22 +25,11 @@ export function createSessionRows(sessionID: Accessor<string>) {
const [rows, setRows] = createStore<SessionRow[]>([])
const revertBoundary = () => data.session.get(sessionID())?.revert?.messageID
function pendingIDs() {
const inputs = data.session.input.list(sessionID())
const pending = new Set(inputs)
for (const message of data.session.message.list(sessionID())) {
if (message.type === "compaction" && (message.status === "queued" || message.status === "running"))
pending.add(message.id)
}
return pending
}
function reduce() {
const messages = data.session.message.list(sessionID())
const inputs = new Set(data.session.input.list(sessionID()))
const boundary = revertBoundary()
const visible = boundary ? messages.filter((message) => message.id < boundary) : messages
const pending = pendingIDs()
const rows = reduceSessionRows(visible.filter((message) => !pending.has(message.id)))
const rows = reduceSessionRows(boundary ? messages.filter((message) => message.id < boundary) : messages, inputs)
partitionPending(rows, pendingPermissions())
return rows
}
@@ -63,30 +52,48 @@ export function createSessionRows(sessionID: Accessor<string>) {
})
createEffect(
on(sessionID, () => {
setRows(reduce())
on(sessionID, (id) => {
setRows(reconcile(reduce()))
void data.session.message.refresh(id).then(
() => {
if (sessionID() !== id) return
setRows(reconcile(reduce()))
},
() => undefined,
)
}),
)
// Re-reduce when the revert boundary changes (stage/clear/commit).
createEffect(
on(revertBoundary, () => {
setRows(reduce())
setRows(reconcile(reduce()))
}),
)
// Pending inputs and compaction leaving the pending set change history membership.
createEffect(
on(
() => {
const messages = data.session.message.list(sessionID())
const pending = data.session.input.list(sessionID()).join("\0")
const compaction = messages
.filter((message) => message.type === "compaction")
.map((message) => `${message.id}:${message.status}`)
.join("\0")
return `${pending}\u0001${compaction}`
},
() => setRows(reduce()),
() =>
data.session.message.list(sessionID()).flatMap((message) =>
message.type === "user"
? [
{
id: message.id,
created: message.time.created,
input: data.session.input.has(sessionID(), message.id),
},
]
: message.type === "compaction"
? [
{
id: message.id,
created: message.time.created,
input: message.status === "running",
},
]
: [],
),
() => setRows(reconcile(reduce())),
),
)
@@ -94,9 +101,12 @@ export function createSessionRows(sessionID: Accessor<string>) {
setRows(
produce((draft) => {
if (draft.some((row) => row.type === "message" && row.messageID === messageID)) return
if (pendingIDs().has(messageID)) return
completePrevious(draft)
draft.push({ type: "message", messageID })
const pending = isPending(messageID)
const message = data.session.message.get(sessionID(), messageID)
const index =
message?.type === "compaction" && pending ? queuedStart(draft) : pending ? draft.length : queuedStart(draft)
if (!pending) completePrevious(draft, index)
draft.splice(index, 0, { type: "message", messageID })
}),
)
@@ -104,14 +114,15 @@ export function createSessionRows(sessionID: Accessor<string>) {
setRows(
produce((draft) => {
if (hasPart(draft, ref)) return
const index = queuedStart(draft)
if (name && exploration(name)) {
const previous = draft.at(-1)
const previous = draft[index - 1]
if (previous?.type === "group" && previous.kind === "exploration") {
previous.refs.push(ref)
return
}
completePrevious(draft)
draft.push({
completePrevious(draft, index)
draft.splice(index, 0, {
type: "group",
kind: "exploration",
refs: [ref],
@@ -120,8 +131,8 @@ export function createSessionRows(sessionID: Accessor<string>) {
})
return
}
completePrevious(draft)
draft.push({ type: "part", ref })
completePrevious(draft, index)
draft.splice(index, 0, { type: "part", ref })
}),
)
@@ -129,8 +140,9 @@ export function createSessionRows(sessionID: Accessor<string>) {
setRows(
produce((draft) => {
if (draft.some((row) => row.type === "assistant-footer" && row.messageID === messageID)) return
completePrevious(draft)
draft.push({ type: "assistant-footer", messageID })
const index = queuedStart(draft)
completePrevious(draft, index)
draft.splice(index, 0, { type: "assistant-footer", messageID })
}),
)
@@ -142,13 +154,26 @@ export function createSessionRows(sessionID: Accessor<string>) {
}),
)
const isPending = (messageID: string) => {
const message = data.session.message.get(sessionID(), messageID)
if (message?.type === "user") return data.session.input.has(sessionID(), messageID)
return message?.type === "compaction" && message.status === "running"
}
const queuedStart = (rows: SessionRow[]) => {
const index = rows.findIndex((row) => row.type === "message" && isPending(row.messageID))
return index === -1 ? rows.length : index
}
const message = (event: { id: string; data: { sessionID: string } }) => {
if (event.data.sessionID === sessionID()) appendMessage(event.id.replace(/^evt_/, "msg_"))
}
const input = (event: { data: { sessionID: string; inputID: string } }) => {
if (event.data.sessionID === sessionID()) appendMessage(event.data.inputID)
}
const subscriptions = [
data.on("session.prompt.promoted", (event) => {
if (event.data.sessionID === sessionID()) appendMessage(event.data.inputID)
}),
data.on("session.prompt.admitted", input),
data.on("session.compaction.started", message),
data.on("session.instructions.updated", message),
data.on("session.synthetic", (event) => {
if (event.data.sessionID === sessionID() && event.data.description?.trim())
@@ -157,7 +182,9 @@ export function createSessionRows(sessionID: Accessor<string>) {
data.on("session.shell.started", message),
data.on("session.agent.selected", message),
data.on("session.model.selected", message),
data.on("session.compaction.ended", (event) => {
if (event.data.reason !== "manual") message(event)
}),
data.on("session.text.delta", (event) => {
if (event.data.sessionID === sessionID())
appendPart({ messageID: event.data.assistantMessageID, partID: `text:${event.data.ordinal}` })
@@ -197,11 +224,18 @@ export function createSessionRows(sessionID: Accessor<string>) {
return rows
}
export function reduceSessionRows(messages: SessionMessage[]) {
return messages.reduce<SessionRow[]>((rows, message) => {
export function reduceSessionRows(messages: SessionMessageInfo[], inputs = new Set<string>()) {
const isInput = (message: SessionMessageInfo) => inputs.has(message.id)
const pendingCompactions = messages.filter((message) => message.type === "compaction" && message.status === "running")
const pending = new Set([...pendingCompactions.map((message) => message.id), ...inputs])
return [
...messages.filter((message) => !pending.has(message.id)),
...pendingCompactions,
...messages.filter(isInput),
].reduce<SessionRow[]>((rows, message) => {
if (message.type !== "assistant") {
if (message.type === "synthetic" && !message.description?.trim()) return rows
completePrevious(rows)
if (!pending.has(message.id)) completePrevious(rows)
rows.push({ type: "message", messageID: message.id })
return rows
}
+27 -106
View File
@@ -6,7 +6,7 @@ import { SessionMessage } from "@opencode-ai/core/session/message"
import { EventV2 } from "@opencode-ai/core/event"
import { onMount } from "solid-js"
import { ProjectProvider } from "../../../src/context/project"
import { SDKProvider } from "../../../src/context/sdk"
import { SDKProvider, useSDK } from "../../../src/context/sdk"
import { DataProvider, useData } from "../../../src/context/data"
import { createSessionRows, type SessionRow } from "../../../src/routes/session/rows"
import { createApi, createClient, createEventStream, createFetch, directory, json } from "../../fixture/tui-sdk"
@@ -114,76 +114,6 @@ test("refreshes resources into reactive getters", async () => {
}
})
test("pages older messages through nested message state", async () => {
const events = createEventStream()
const sessionID = "ses_message_page"
const pages: Array<{ limit?: string | null; order?: string | null; cursor?: string | null }> = []
// Full first page (desc) so complete stays false until a short older page arrives.
const first = Array.from({ length: 50 }, (_, index) => {
const n = 51 - index
return { id: `msg_${n}`, type: "user" as const, text: String(n), time: { created: n } }
})
const calls = createFetch((url) => {
if (url.pathname !== `/api/session/${sessionID}/message`) return
pages.push({
limit: url.searchParams.get("limit"),
order: url.searchParams.get("order"),
cursor: url.searchParams.get("cursor"),
})
if (!url.searchParams.get("cursor"))
return json({
data: first,
cursor: { next: "cursor-older" },
})
return json({
data: [{ id: "msg_1", type: "user", text: "one", time: { created: 1 } }],
cursor: {},
})
}, events)
let data!: ReturnType<typeof useData>
function Probe() {
data = useData()
return <box />
}
const app = await testRender(() => (
<TestTuiContexts>
<SDKProvider client={createClient(calls.fetch)} api={createApi(calls.fetch)}>
<ProjectProvider>
<DataProvider>
<Probe />
</DataProvider>
</ProjectProvider>
</SDKProvider>
</TestTuiContexts>
))
try {
await data.session.message.refresh(sessionID)
expect(pages).toEqual([{ limit: "50", order: "desc", cursor: null }])
expect(data.session.message.ids(sessionID)).toEqual(first.toReversed().map((message) => message.id))
expect(data.session.message.cursor(sessionID)).toBe("cursor-older")
expect(data.session.message.complete(sessionID)).toBe(false)
expect(data.session.message.loading(sessionID)).toBe(false)
await data.session.message.more(sessionID)
expect(pages).toEqual([
{ limit: "50", order: "desc", cursor: null },
{ limit: "50", order: null, cursor: "cursor-older" },
])
expect(data.session.message.ids(sessionID)[0]).toBe("msg_1")
expect(data.session.message.ids(sessionID)).toHaveLength(51)
expect(data.session.message.cursor(sessionID)).toBeUndefined()
expect(data.session.message.complete(sessionID)).toBe(true)
await data.session.message.more(sessionID)
expect(pages).toHaveLength(2)
} finally {
app.renderer.destroy()
}
})
test("applies absolute usage events to session info", async () => {
const events = createEventStream()
const sessionID = "ses_usage_refresh"
@@ -564,9 +494,11 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
})
}, events)
let data!: ReturnType<typeof useData>
let sdk!: ReturnType<typeof useSDK>
function Probe() {
data = useData()
sdk = useSDK()
return <box />
}
@@ -585,41 +517,39 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
try {
await wait(() => data.location.model.list()?.[0]?.id === "model-1")
await wait(() => data.session.status("session-stale") === "running")
expect(data.connection.status()).toBe("connected")
expect(data.connection.attempt()).toBe(0)
expect(sdk.connection.status()).toBe("connected")
expect(sdk.connection.attempt()).toBe(0)
events.disconnect()
await wait(() => data.connection.status() === "connecting")
expect(data.connection.attempt()).toBe(1)
expect(data.connection.error()).toBe("Event stream disconnected")
await wait(() => sdk.connection.status() === "reconnecting")
expect(sdk.connection.attempt()).toBe(1)
expect(sdk.connection.error()).toBe("Event stream disconnected")
await wait(() => requests.active === 2 && data.connection.status() === "connected", 4000)
await wait(() => requests.active === 2 && sdk.connection.status() === "connected", 4000)
resolveActive(json({ data: { "session-new": { type: "running" } } }))
await wait(() => data.location.model.list()?.[0]?.id === "model-2", 4000)
await wait(() => data.session.status("session-stale") === "idle")
expect(data.session.status("session-new")).toBe("running")
expect(requests.event).toBe(2)
expect(data.connection.status()).toBe("connected")
expect(data.connection.attempt()).toBe(0)
expect(data.connection.error()).toBeUndefined()
expect(sdk.connection.status()).toBe("connected")
expect(sdk.connection.attempt()).toBe(0)
expect(sdk.connection.error()).toBeUndefined()
} finally {
app.renderer.destroy()
}
})
test("keeps pending prompts out of history rows until promoted", async () => {
test("completes exploration when a queued prompt is promoted", async () => {
const events = createEventStream()
const sessionID = "session-promotion"
const calls = createFetch((url) => {
if (url.pathname === `/api/session/${sessionID}/message`) return json({ data: [], cursor: {} })
}, events)
let rows!: ReturnType<typeof createSessionRows>
let data!: ReturnType<typeof useData>
function Probe() {
rows = createSessionRows(() => sessionID)
data = useData()
return <box />
}
@@ -674,8 +604,7 @@ test("keeps pending prompts out of history rows until promoted", async () => {
delivery: "steer",
},
})
await wait(() => data.session.input.has(sessionID, "message-user"))
expect(rows.some((row) => row.type === "message" && row.messageID === "message-user")).toBe(false)
await wait(() => rows.at(-1)?.type === "message")
expect(rows.find((row) => row.type === "group")?.completed).toBe(false)
emitEvent(events, {
@@ -685,8 +614,7 @@ test("keeps pending prompts out of history rows until promoted", async () => {
durable: durable(sessionID, 3),
data: { sessionID, inputID: "message-user" },
})
await wait(() => rows.some((row) => row.type === "message" && row.messageID === "message-user"))
expect(data.session.input.has(sessionID, "message-user")).toBe(false)
await wait(() => rows.find((row) => row.type === "group")?.completed === true)
expect(rows.at(-1)).toEqual({ type: "message", messageID: "message-user" })
} finally {
app.renderer.destroy()
@@ -747,7 +675,7 @@ test("removes committed revert messages from local state", async () => {
}
})
test("connectedOnce is false until first connect and persists across disconnect", async () => {
test("distinguishes initial connection from reconnection", async () => {
const encoder = new TextEncoder()
let stream: ReadableStreamDefaultController<Uint8Array> | undefined
const eventResponse = () =>
@@ -773,10 +701,10 @@ test("connectedOnce is false until first connect and persists across disconnect"
const calls = createFetch((url) => {
if (url.pathname === "/api/event") return eventResponse()
})
let data!: ReturnType<typeof useData>
let sdk!: ReturnType<typeof useSDK>
function Probe() {
data = useData()
sdk = useSDK()
return <box />
}
@@ -794,16 +722,13 @@ test("connectedOnce is false until first connect and persists across disconnect"
try {
await wait(() => stream !== undefined)
expect(data.connection.status()).toBe("connecting")
expect(data.connection.connectedOnce()).toBe(false)
expect(sdk.connection.status()).toBe("connecting")
connect()
await wait(() => data.connection.status() === "connected")
expect(data.connection.connectedOnce()).toBe(true)
await wait(() => sdk.connection.status() === "connected")
disconnect()
await wait(() => data.connection.status() === "connecting")
expect(data.connection.connectedOnce()).toBe(true)
await wait(() => sdk.connection.status() === "reconnecting")
} finally {
app.renderer.destroy()
}
@@ -1443,9 +1368,11 @@ test("adds and dismisses permission requests from live events", async () => {
const events = createEventStream()
const calls = createFetch(undefined, events)
let data!: ReturnType<typeof useData>
let sdk!: ReturnType<typeof useSDK>
function Probe() {
data = useData()
sdk = useSDK()
return <box />
}
@@ -1462,7 +1389,7 @@ test("adds and dismisses permission requests from live events", async () => {
))
try {
await wait(() => data.connection.status() === "connected")
await wait(() => sdk.connection.status() === "connected")
emitEvent(events, {
id: "evt_permission_asked_1",
created: 0,
@@ -1561,9 +1488,11 @@ test("adds, dismisses, and refreshes form requests", async () => {
return json({ data: [{ id: "frm_remote", sessionID: "ses_1", mode: "form", fields: [] }] })
}, events)
let data!: ReturnType<typeof useData>
let sdk!: ReturnType<typeof useSDK>
function Probe() {
data = useData()
sdk = useSDK()
return <box />
}
@@ -1580,7 +1509,7 @@ test("adds, dismisses, and refreshes form requests", async () => {
))
try {
await wait(() => data.connection.status() === "connected")
await wait(() => sdk.connection.status() === "connected")
emitEvent(events, {
id: "evt_form_created_1",
created: 0,
@@ -1751,14 +1680,6 @@ test("settles pending tools when a live failure arrives", async () => {
name: "bash",
},
})
await wait(() => {
const assistant = sync.session.message.get("session-1", "msg_explicit_assistant_9")
return (
assistant?.type === "assistant" &&
assistant.content[0]?.type === "tool" &&
assistant.content[0].state.status === "streaming"
)
})
emitEvent(events, {
id: "evt_called_1",
created: 0,
+14 -8
View File
@@ -207,25 +207,31 @@ test("renders a footer for a pre-output retry assistant after replay", () => {
expect(reduceSessionRows([message])).toEqual([{ type: "assistant-footer", messageID: "assistant-retry" }])
})
test("history reduce keeps chronological order without pending reordering", () => {
test("places a running compaction barrier before every queued user message", () => {
const queued = (id: string, text: string, created: number): SessionMessageInfo => ({
type: "user",
id,
text,
time: { created },
})
const messages: SessionMessageInfo[] = [
{ type: "user", id: "user-1", text: "Before", time: { created: 1 } },
queued("user-before", "Before", 1),
{
type: "compaction",
id: "compaction",
status: "completed",
status: "running",
reason: "manual",
summary: "done",
summary: "",
recent: "",
time: { created: 2 },
},
{ type: "user", id: "user-2", text: "After", time: { created: 3 } },
queued("user-after", "After", 3),
]
expect(reduceSessionRows(messages)).toEqual([
{ type: "message", messageID: "user-1" },
expect(reduceSessionRows(messages, new Set(["user-before", "user-after"]))).toEqual([
{ type: "message", messageID: "compaction" },
{ type: "message", messageID: "user-2" },
{ type: "message", messageID: "user-before" },
{ type: "message", messageID: "user-after" },
])
})