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`. 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 ## Style Guide
### General Principles ### 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 { Global } from "@opencode-ai/core/global"
import { InstallationVersion } from "@opencode-ai/core/installation/version" import { InstallationVersion } from "@opencode-ai/core/installation/version"
import { AppProcess } from "@opencode-ai/core/process" import { AppProcess } from "@opencode-ai/core/process"
import { EffectFlock } from "@opencode-ai/core/util/effect-flock"
import { start } from "@opencode-ai/server/process" import { start } from "@opencode-ai/server/process"
import { randomBytes, randomUUID } from "node:crypto" import { randomBytes, randomUUID } from "node:crypto"
import path from "node:path" import path from "node:path"
@@ -24,10 +25,13 @@ export type Options = {
readonly port?: number readonly port?: number
} }
type ManagedServiceOptions = Service.Options & { readonly file: string }
export const run = Effect.fn("cli.server-process.run")((options: Options) => export const run = Effect.fn("cli.server-process.run")((options: Options) =>
processEffect(options).pipe( processEffect(options).pipe(
Effect.catchTag("ServiceAlreadyOwned", () => Effect.logInfo("another process owns the managed service, exiting")),
Effect.provide(Updater.layer), 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), Effect.provide(NodeServices.layer),
), ),
) )
@@ -43,6 +47,21 @@ const processEffect = Effect.fnUntraced(function* (options: Options) {
delete process.env.OPENCODE_SERVER_PASSWORD delete process.env.OPENCODE_SERVER_PASSWORD
} }
const config = options.mode === "service" ? yield* ServiceConfig.read() : {} 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 = const password =
options.mode === "service" options.mode === "service"
? yield* ServiceConfig.password() ? yield* ServiceConfig.password()
@@ -55,7 +74,7 @@ const processEffect = Effect.fnUntraced(function* (options: Options) {
port: Option.fromNullishOr(options.port ?? config.port), port: Option.fromNullishOr(options.port ?? config.port),
password, password,
}).pipe(Effect.provide(Logger.layer([], { mergeWithExisting: false }))) }).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) const url = HttpServer.formatAddress(address)
console.log(options.mode === "stdio" ? JSON.stringify({ url }) : `server listening on ${url}`) console.log(options.mode === "stdio" ? JSON.stringify({ url }) : `server listening on ${url}`)
if (options.mode === "default" && !environmentPassword) console.log(`server password ${password}`) 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 encodeInfo = Schema.encodeEffect(infoJson)
const decodeInfo = Schema.decodeUnknownEffect(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 fs = yield* FileSystem.FileSystem
const options = yield* ServiceConfig.options()
const id = randomUUID() const id = randomUUID()
const temp = options.file + "." + id + ".tmp" const temp = options.file + "." + id + ".tmp"
yield* fs.makeDirectory(path.dirname(options.file), { recursive: true }) 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.void
: Effect.try({ try: () => process.kill(process.pid, "SIGTERM"), catch: (cause) => cause }).pipe(Effect.ignore), : 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, Effect.forkScoped,
) )
yield* Effect.addFinalizer(() => 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 decodeMeta = Schema.decodeUnknownSync(LockMetaJson)
const decodeMetaOption = Schema.decodeUnknownOption(LockMetaJson)
const encodeMeta = Schema.encodeSync(LockMetaJson) const encodeMeta = Schema.encodeSync(LockMetaJson)
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@@ -74,6 +75,7 @@ export namespace EffectFlock {
export interface Interface { export interface Interface {
readonly acquire: (key: string, dir?: string) => Effect.Effect<void, LockError, Scope.Scope> 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: { readonly withLock: {
(key: string, dir?: string): <A, E, R>(body: Effect.Effect<A, E, R>) => Effect.Effect<A, E | LockError, R> (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> <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 isStale = Effect.fnUntraced(function* (lockDir: string, heartbeatPath: string, metaPath: string) {
const now = wall() 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) const hb = yield* safeStat(heartbeatPath)
if (hb) return now - mtimeMs(hb) > STALE_MS if (hb) return now - mtimeMs(hb) > STALE_MS
@@ -243,28 +266,46 @@ export namespace EffectFlock {
catch: (cause) => new ReleaseError({ detail: "metadata invalid", cause }), catch: (cause) => new ReleaseError({ detail: "metadata invalid", cause }),
}).pipe(Effect.orDie) }).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) yield* forceRemove(handle.lockDir)
}) })
// -- build service -- // -- build service --
const acquire = Effect.fn("EffectFlock.acquire")(function* (key: string, dir?: string) { const heartbeat = Effect.fnUntraced(function* (handle: Handle) {
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
yield* fs yield* fs
.utimes(handle.heartbeatPath, new Date(), new Date()) .utimes(handle.heartbeatPath, new Date(), new Date())
.pipe(Effect.ignore, Effect.repeat(Schedule.spaced(HEARTBEAT_MS)), Effect.forkScoped) .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( const withLock: Interface["withLock"] = Function.dual(
(args) => Effect.isEffect(args[0]), (args) => Effect.isEffect(args[0]),
<A, E, R>(body: Effect.Effect<A, E, R>, key: string, dir?: string): 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> =>
@@ -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 { Cause, Effect, Exit } from "effect"
import { testEffect } from "../lib/effect" import { testEffect } from "../lib/effect"
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" 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 { EffectFlock } from "@opencode-ai/core/util/effect-flock"
import { Global } from "@opencode-ai/core/global" import { Global } from "@opencode-ai/core/global"
import { Hash } from "@opencode-ai/core/util/hash" 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( it.live(
"withLock data-first", "withLock data-first",
Effect.gen(function* () { Effect.gen(function* () {
@@ -358,7 +375,7 @@ describe("util.effect-flock", () => {
) )
it.live( it.live(
"recovers after a crashed lock owner", "immediately recovers after a local lock owner crashes",
() => () =>
Effect.promise(async () => { Effect.promise(async () => {
const tmp = await fs.mkdtemp(path.join(os.tmpdir(), "eflock-crash-")) 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 waitForFile(ready, 5_000)
await stopWorker(proc) 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 done = path.join(tmp, "done.log")
const result = await run({ key: "eflock:crash", dir, done, holdMs: 10 }) const result = await run({ key: "eflock:crash", dir, done, holdMs: 10 })
expect(result.code).toBe(0) expect(result.code).toBe(0)
+2 -1
View File
@@ -10,7 +10,8 @@
"./backend/*": "./src/backend/*.ts", "./backend/*": "./src/backend/*.ts",
"./frontend": "./src/frontend/simulation.ts", "./frontend": "./src/frontend/simulation.ts",
"./frontend/*": "./src/frontend/*.ts", "./frontend/*": "./src/frontend/*.ts",
"./protocol": "./src/protocol/index.ts" "./protocol": "./src/protocol/index.ts",
"./recording": "./src/recording.ts"
}, },
"scripts": { "scripts": {
"typecheck": "tsgo --noEmit" "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 { tmpdir } from "node:os"
import { join } from "node:path" import { dirname, join } from "node:path"
import type { CapturedFrame, CliRenderer, Renderable } from "@opentui/core" import type { CliRenderer, Renderable } from "@opentui/core"
import { createMockKeys, createMockMouse, type MockInput, type MockMouse } from "@opentui/core/testing" import { createMockKeys, createMockMouse, type MockInput, type MockMouse } from "@opentui/core/testing"
import type { SimulationProtocol } from "../protocol" import type { SimulationProtocol } from "../protocol"
import { SimulationRenderer } from "./renderer" 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() await harness.renderOnce()
const image = SimulationPng.screenshot(harness.renderer) 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) await Bun.write(path, image.data)
return path 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) { export async function execute(harness: Harness, action: Action) {
switch (action.type) { switch (action.type) {
case "ui.type": case "ui.type":
@@ -180,7 +128,9 @@ export async function execute(harness: Harness, action: Action) {
harness.mockInput.pressArrow(action.direction) harness.mockInput.pressArrow(action.direction)
break break
case "ui.focus": 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 break
case "ui.click": case "ui.click":
await harness.mockMouse.click(action.x, action.y) await harness.mockMouse.click(action.x, action.y)
+33 -5
View File
@@ -1,21 +1,43 @@
import type { CliRenderer, CliRendererConfig } from "@opentui/core" import type { CliRenderer, CliRendererConfig } from "@opentui/core"
import { createTestRenderer, type TestRendererSetup } from "@opentui/core/testing" import { createTestRenderer, type TestRendererSetup } from "@opentui/core/testing"
import { Timeline } from "../recording"
const setups = new WeakMap<CliRenderer, TestRendererSetup>() const setups = new WeakMap<CliRenderer, TestRendererSetup>()
const recordings = new WeakMap<CliRenderer, Timeline>()
/** /**
* Creates the headless simulation renderer: a real CliRenderer backed by an * Creates a headless renderer with optional recording: a real CliRenderer
* in-memory screen buffer instead of a terminal. The TestRendererSetup is * backed by an in-memory screen buffer. The TestRendererSetup is kept
* kept module-side (keyed by renderer) so the harness can use the supported * module-side so the harness can use supported testing APIs without app
* testing APIs without app code carrying it around. * 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({ const setup = await createTestRenderer({
...options, ...options,
width: 100, width: 100,
height: 40, 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) setups.set(setup.renderer, setup)
recordings.set(setup.renderer, recording)
return setup.renderer return setup.renderer
} }
@@ -23,4 +45,10 @@ export function setupFor(renderer: CliRenderer): TestRendererSetup | undefined {
return setups.get(renderer) 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" export * as SimulationRenderer from "./renderer"
+9 -42
View File
@@ -1,4 +1,3 @@
import type { CapturedFrame } from "@opentui/core"
import { SimulationProtocol } from "../protocol" import { SimulationProtocol } from "../protocol"
import { SimulationActions, type Harness } from "./actions" 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())) 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( async function handle(
harness: Harness, harness: Harness,
request: SimulationProtocol.Frontend.Request, request: SimulationProtocol.Frontend.Request,
recording: { current?: Recording }, finishRecording?: () => Promise<string>,
headless: boolean,
) { ) {
switch (request.method) { switch (request.method) {
case "ui.screenshot": case "ui.screenshot":
return SimulationActions.screenshot(harness) return SimulationActions.screenshot(harness, request.params?.path)
case "ui.state": { case "ui.state": {
return SimulationActions.state(harness) return SimulationActions.state(harness)
} }
case "ui.start-record": { case "ui.recording.finish":
if (recording.current) throw new Error("UI recording is already active") if (!finishRecording) throw new Error("UI recording is not available")
const frames = [SimulationActions.frame(harness)] return finishRecording()
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.type": case "ui.type":
return SimulationActions.execute(harness, { type: "ui.type", text: request.params.text }) return SimulationActions.execute(harness, { type: "ui.type", text: request.params.text })
case "ui.enter": 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 url = new URL(endpoint)
const recording: { current?: Recording } = {} const server = Bun.serve<{ readonly drive: true }>({
const server = Bun.serve<{ readonly drive: true; readonly headless: boolean }>({
hostname: url.hostname, hostname: url.hostname,
port: Number(url.port), port: Number(url.port),
fetch(request, server) { 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 }) return new Response("opencode drive ui websocket", { status: 426 })
}, },
websocket: { websocket: {
@@ -94,7 +62,7 @@ export function start(harness: Harness, endpoint: string, headless: boolean): Se
let request: SimulationProtocol.Frontend.Request | undefined let request: SimulationProtocol.Frontend.Request | undefined
try { try {
request = parseRequest(message) 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) const next = SimulationProtocol.JsonRpc.success(request.id, result)
if (next) socket.send(JSON.stringify(next)) if (next) socket.send(JSON.stringify(next))
} catch (error) { } catch (error) {
@@ -106,7 +74,6 @@ export function start(harness: Harness, endpoint: string, headless: boolean): Se
return { return {
url: endpoint, url: endpoint,
stop: () => { stop: () => {
if (recording.current) clearInterval(recording.current.timer)
server.stop(true) server.stop(true)
}, },
} }
@@ -14,11 +14,14 @@ import { SimulationServer } from "./server"
*/ */
export async function create(options: CliRendererConfig): Promise<CliRenderer> { export async function create(options: CliRendererConfig): Promise<CliRenderer> {
const headless = process.env.OPENCODE_DRIVE_RENDERER === "headless" 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( const server = SimulationServer.start(
SimulationActions.createHarness(renderer), SimulationActions.createHarness(renderer),
DriveManifest.resolve().endpoints.ui, manifest.endpoints.ui,
headless, headless && manifest.recording ? () => SimulationRenderer.finish(renderer) : undefined,
) )
process.stderr.write(`opencode drive ui websocket: ${server.url}\n`) process.stderr.write(`opencode drive ui websocket: ${server.url}\n`)
renderer.once("destroy", () => server.stop()) renderer.once("destroy", () => server.stop())
+10 -9
View File
@@ -1,12 +1,15 @@
import { existsSync, readFileSync } from "node:fs" import { existsSync, readFileSync } from "node:fs"
import { homedir } from "node:os" import { homedir } from "node:os"
import { join } from "node:path" import { isAbsolute, join } from "node:path"
export interface Manifest { export interface Manifest {
readonly endpoints: { readonly endpoints: {
readonly ui: string readonly ui: string
readonly backend: string readonly backend: string
} }
readonly recording?: {
readonly timeline: string
}
} }
export const defaults: Manifest = { export const defaults: Manifest = {
@@ -32,18 +35,16 @@ export function resolve() {
if (!isManifest(manifest)) throw new Error(`Invalid drive manifest: ${file}`) if (!isManifest(manifest)) throw new Error(`Invalid drive manifest: ${file}`)
validateEndpoint(manifest.endpoints.ui, "ui") validateEndpoint(manifest.endpoints.ui, "ui")
validateEndpoint(manifest.endpoints.backend, "backend") 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 return manifest
} }
function isManifest(value: unknown): value is Manifest { function isManifest(value: unknown): value is Manifest {
if (typeof value !== "object" || value === null) return false if (typeof value !== "object" || value === null || !("endpoints" in value)) return false
if (!("endpoints" in value) || typeof value.endpoints !== "object" || value.endpoints === null) return false if (typeof value.endpoints !== "object" || value.endpoints === null) return false
return ( return "ui" in value.endpoints && "backend" in value.endpoints
"ui" in value.endpoints &&
typeof value.endpoints.ui === "string" &&
"backend" in value.endpoints &&
typeof value.endpoints.backend === "string"
)
} }
function validateEndpoint(value: string, name: string) { function validateEndpoint(value: string, name: string) {
+10 -13
View File
@@ -94,11 +94,11 @@ export namespace Frontend {
export const Screenshot = Schema.String export const Screenshot = Schema.String
export type Screenshot = Schema.Schema.Type<typeof Screenshot> export type Screenshot = Schema.Schema.Type<typeof Screenshot>
export const StartRecord = Schema.Struct({ recording: Schema.Literal(true) }) export const RecordingFinish = Schema.String
export interface StartRecord extends Schema.Schema.Type<typeof StartRecord> {} export type RecordingFinish = Schema.Schema.Type<typeof RecordingFinish>
export const EndRecord = Schema.String export const ArtifactParams = Schema.Struct({ path: Schema.optional(Schema.String) })
export type EndRecord = Schema.Schema.Type<typeof EndRecord> export interface ArtifactParams extends Schema.Schema.Type<typeof ArtifactParams> {}
export const TypeParams = Schema.Struct({ text: Schema.String }) export const TypeParams = Schema.Struct({ text: Schema.String })
export interface TypeParams extends Schema.Schema.Type<typeof TypeParams> {} 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.Literal("ui.click"), params: ClickParams }),
Schema.Struct({ Schema.Struct({
...JsonRpc.RequestFields, ...JsonRpc.RequestFields,
method: Schema.Literals([ method: Schema.Literal("ui.screenshot"),
"ui.enter", params: Schema.optional(ArtifactParams),
"ui.screenshot", }),
"ui.state", Schema.Struct({
"ui.start-record", ...JsonRpc.RequestFields,
"ui.end-record", method: Schema.Literals(["ui.enter", "ui.state", "ui.recording.finish"]),
]),
}), }),
]) ])
export type Request = Schema.Schema.Type<typeof Request> export type Request = Schema.Schema.Type<typeof Request>
export const decodeRequest = Schema.decodeUnknownSync(Request) export const decodeRequest = Schema.decodeUnknownSync(Request)
} }
export namespace Backend { export namespace Backend {
@@ -189,7 +187,6 @@ export namespace Backend {
matched: Schema.Boolean, matched: Schema.Boolean,
}) })
export interface NetworkLogEntry extends Schema.Schema.Type<typeof NetworkLogEntry> {} export interface NetworkLogEntry extends Schema.Schema.Type<typeof NetworkLogEntry> {}
} }
export * as SimulationProtocol from "./index" 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 }) return render({ params: route.data.data })
}) })
// Suppress the full-screen reconnecting overlay for transient disconnects (initial startup, host // Suppress the full-screen overlay for transient startup and event-stream retry states.
// reload, sub-second event-stream blips). After the first successful connect, show it only once the // Initial connection gets a longer grace period; retries surface more quickly.
// 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".
const [showReconnecting, setShowReconnecting] = createSignal(false) const [showReconnecting, setShowReconnecting] = createSignal(false)
let reconnectTimer: ReturnType<typeof setTimeout> | undefined let reconnectTimer: ReturnType<typeof setTimeout> | undefined
createEffect(() => { createEffect(() => {
@@ -1158,7 +1155,8 @@ function App(props: { onSnapshot?: () => Promise<string[]>; pluginHost: TuiPlugi
clearTimeout(reconnectTimer) clearTimeout(reconnectTimer)
reconnectTimer = undefined reconnectTimer = undefined
} }
if (sdk.connection.status() !== "connecting") { const status = sdk.connection.status()
if (status === "connected") {
setShowReconnecting(false) setShowReconnecting(false)
return return
} }
@@ -1167,7 +1165,7 @@ function App(props: { onSnapshot?: () => Promise<string[]>; pluginHost: TuiPlugi
reconnectTimer = undefined reconnectTimer = undefined
setShowReconnecting(true) setShowReconnecting(true)
}, },
sdk.connection.connectedOnce() ? 1000 : 5000, status === "reconnecting" ? 1000 : 5000,
).unref() ).unref()
}) })
onCleanup(() => { onCleanup(() => {
+75 -151
View File
@@ -4,74 +4,62 @@
// Reconnect may re-bootstrap; that is enough. UI and the server own ordering concerns. // Reconnect may re-bootstrap; that is enough. UI and the server own ordering concerns.
import type { import type {
AgentV2Info, AgentInfo,
CommandV2Info, CommandInfo,
FormFormInfo, FormFormInfo,
FormUrlInfo, FormUrlInfo,
IntegrationInfo, IntegrationInfo,
LocationRef, LocationRef,
McpServer, McpServer,
ModelV2Info, ModelInfo,
PermissionSavedInfo, PermissionSavedInfo,
PermissionV2Request, PermissionV2Request,
ProviderV2Info, ProviderV2Info,
ReferenceInfo, ReferenceInfo,
SessionMessage, SessionMessageInfo,
SessionMessageAssistant, SessionMessageAssistant,
SessionMessageAssistantReasoning, SessionMessageAssistantReasoning,
SessionMessageAssistantText, SessionMessageAssistantText,
SessionMessageAssistantTool, SessionMessageAssistantTool,
SessionV2Info, SessionInfo,
Shell, Shell,
SkillV2Info, SkillInfo,
V2Event, V2Event,
} from "@opencode-ai/sdk/v2" } from "@opencode-ai/sdk/v2"
import { createStore, produce, reconcile } from "solid-js/store" import { createStore, produce, reconcile } from "solid-js/store"
import { createSimpleContext } from "./helper" import { createSimpleContext } from "./helper"
import { useSDK } from "./sdk" import { useSDK } from "./sdk"
import { batch, createSignal, onCleanup } from "solid-js" import { createSignal, onCleanup } from "solid-js"
export type DataSessionStatus = "idle" | "running" export type DataSessionStatus = "idle" | "running"
const messageIDFromEvent = (eventID: string) => eventID.replace(/^evt_/, "msg_") const messageIDFromEvent = (eventID: string) => eventID.replace(/^evt_/, "msg_")
const MESSAGE_PAGE_SIZE = 25
export type FormInfo = FormFormInfo | FormUrlInfo 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 = { type LocationData = {
agent?: AgentV2Info[] agent?: AgentInfo[]
command?: CommandV2Info[] command?: CommandInfo[]
integration?: IntegrationInfo[] integration?: IntegrationInfo[]
mcp?: McpServer[] mcp?: McpServer[]
model?: ModelV2Info[] model?: ModelInfo[]
provider?: ProviderV2Info[] provider?: ProviderV2Info[]
reference?: ReferenceInfo[] reference?: ReferenceInfo[]
// Currently running shell commands for this location, keyed by shell id. Entries are removed // 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. // once the command exits or is deleted, so this only ever holds in-flight shells.
shell?: Record<string, Shell> shell?: Record<string, Shell>
skill?: SkillV2Info[] skill?: SkillInfo[]
} }
type Data = { type Data = {
session: { session: {
info: Record<string, SessionV2Info> info: Record<string, SessionInfo>
// Family index keyed by a family's root (or furthest-known-ancestor when the // 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 // 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. // session ID in that family, including the key itself once its info arrives.
family: Record<string, string[]> family: Record<string, string[]>
status: Record<string, DataSessionStatus> status: Record<string, DataSessionStatus>
compaction: Partial<Record<string, string>> message: Record<string, SessionMessageInfo[]>
compactionReason: Partial<Record<string, "auto" | "manual">>
message: Record<string, SessionMessages>
input: Record<string, string[]> input: Record<string, string[]>
permission: Record<string, PermissionV2Request[]> permission: Record<string, PermissionV2Request[]>
// Pending forms keyed by session ID. // Pending forms keyed by session ID.
@@ -83,10 +71,6 @@ type Data = {
location: Record<string, LocationData> location: Record<string, LocationData>
} }
function emptyMessages(): SessionMessages {
return { items: [], complete: false, loading: false }
}
function locationKey(location: LocationRef) { function locationKey(location: LocationRef) {
return JSON.stringify([location.directory, location.workspaceID]) return JSON.stringify([location.directory, location.workspaceID])
} }
@@ -111,8 +95,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
info: {}, info: {},
family: {}, family: {},
status: {}, status: {},
compaction: {},
compactionReason: {},
message: {}, message: {},
input: {}, input: {},
permission: {}, permission: {},
@@ -136,37 +118,35 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
} }
const message = { 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( setStore(
"session", "session",
"message", "message",
produce((draft) => { 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 if (index.has(item.id)) return
index.set(item.id, messages.length) index.set(item.id, messages.length)
messages.push(item) messages.push(item)
}, },
activeAssistant(messages: SessionMessage[]) { activeAssistant(messages: SessionMessageInfo[]) {
const item = messages.findLast((item) => item.type === "assistant" && !item.time.completed) const item = messages.findLast((item) => item.type === "assistant" && !item.time.completed)
return item?.type === "assistant" ? item : undefined 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 position = index.get(messageID)
const item = position === undefined ? undefined : messages[position] const item = position === undefined ? undefined : messages[position]
return item?.type === "assistant" ? item : undefined return item?.type === "assistant" ? item : undefined
}, },
shell(messages: SessionMessage[], shellID: string) { shell(messages: SessionMessageInfo[], shellID: string) {
const item = messages.findLast((item) => item.type === "shell" && item.shell.id === shellID) const item = messages.findLast((item) => item.type === "shell" && item.shellID === shellID)
return item?.type === "shell" ? item : undefined return item?.type === "shell" ? item : undefined
}, },
compaction(messages: SessionMessage[]) { compaction(messages: SessionMessageInfo[]) {
const item = messages.findLast( const item = messages.findLast((item) => item.type === "compaction" && item.status === "running")
(item) => item.type === "compaction" && (item.status === "queued" || item.status === "running"),
)
return item?.type === "compaction" ? item : undefined return item?.type === "compaction" ? item : undefined
}, },
latestTool(assistant: SessionMessageAssistant | undefined, callID?: string) { latestTool(assistant: SessionMessageAssistant | undefined, callID?: string) {
@@ -244,7 +224,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
produce((draft) => { produce((draft) => {
delete draft.info[sessionID] delete draft.info[sessionID]
delete draft.status[sessionID] delete draft.status[sessionID]
delete draft.compaction[sessionID]
delete draft.message[sessionID] delete draft.message[sessionID]
delete draft.input[sessionID] delete draft.input[sessionID]
delete draft.permission[sessionID] delete draft.permission[sessionID]
@@ -378,6 +357,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
id: messageIDFromEvent(event.id), id: messageIDFromEvent(event.id),
type: "system", type: "system",
text: event.data.text, text: event.data.text,
metadata: event.metadata,
time: { created: event.created }, time: { created: event.created },
}) })
}) })
@@ -387,9 +367,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
message.append(draft, index, { message.append(draft, index, {
id: messageIDFromEvent(event.id), id: messageIDFromEvent(event.id),
type: "synthetic", type: "synthetic",
sessionID: event.data.sessionID,
text: event.data.text, text: event.data.text,
description: event.data.description, description: event.data.description,
metadata: event.data.metadata,
time: { created: event.created }, time: { created: event.created },
}) })
}) })
@@ -399,7 +379,11 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
message.append(draft, index, { message.append(draft, index, {
id: messageIDFromEvent(event.id), id: messageIDFromEvent(event.id),
type: "shell", 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 }, time: { created: event.created },
}) })
}) })
@@ -408,7 +392,8 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
message.update(event.data.sessionID, (draft) => { message.update(event.data.sessionID, (draft) => {
const match = message.shell(draft, event.data.shell.id) const match = message.shell(draft, event.data.shell.id)
if (!match) return 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.output = event.data.output
match.time.completed = event.created match.time.completed = event.created
}) })
@@ -437,6 +422,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
type: "assistant", type: "assistant",
agent: event.data.agent, agent: event.data.agent,
model: event.data.model, model: event.data.model,
metadata: event.metadata,
content: [], content: [],
snapshot: event.data.snapshot ? { start: event.data.snapshot } : undefined, snapshot: event.data.snapshot ? { start: event.data.snapshot } : undefined,
time: { created: event.created }, time: { created: event.created },
@@ -623,29 +609,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
setSessionStatus(event.data.sessionID, "running") setSessionStatus(event.data.sessionID, "running")
break break
case "session.compaction.admitted": 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 break
case "session.compaction.started": 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) => { 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, { message.append(draft, index, {
id: event.data.inputID ?? messageIDFromEvent(event.id), id: event.data.inputID ?? messageIDFromEvent(event.id),
type: "compaction", type: "compaction",
@@ -661,10 +627,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
case "session.execution.failed": case "session.execution.failed":
case "session.execution.interrupted": case "session.execution.interrupted":
setSessionStatus(event.data.sessionID, "idle") 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) => { message.update(event.data.sessionID, (draft) => {
const currentAssistant = message.activeAssistant(draft) const currentAssistant = message.activeAssistant(draft)
if (currentAssistant) currentAssistant.retry = undefined if (currentAssistant) currentAssistant.retry = undefined
@@ -695,22 +657,26 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
}) })
break break
case "session.compaction.delta": case "session.compaction.delta":
setStore("session", "compaction", event.data.sessionID, (text) => (text ?? "") + event.data.text)
message.update(event.data.sessionID, (draft) => { message.update(event.data.sessionID, (draft) => {
const current = message.compaction(draft) const current = message.compaction(draft)
if (current) current.summary += event.data.text if (current?.status === "running") current.summary += event.data.text
}) })
break break
case "session.compaction.ended": 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) => { message.update(event.data.sessionID, (draft, index) => {
const current = message.compaction(draft) const position = draft.findLastIndex((item) => item.type === "compaction" && item.status === "running")
if (current) { const current = draft[position]
current.status = "completed" if (current?.type === "compaction") {
current.reason = event.data.reason draft[position] = {
current.summary = event.data.text id: current.id,
current.recent = event.data.recent type: "compaction",
status: "completed",
reason: event.data.reason,
summary: event.data.text,
recent: event.data.recent,
metadata: current.metadata,
time: current.time,
}
return return
} }
message.append(draft, index, { message.append(draft, index, {
@@ -725,11 +691,26 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
}) })
break break
case "session.compaction.failed": case "session.compaction.failed":
setStore("session", "compaction", event.data.sessionID, undefined) message.update(event.data.sessionID, (draft, index) => {
setStore("session", "compactionReason", event.data.sessionID, undefined) const position = draft.findLastIndex((item) => item.type === "compaction" && item.status === "running")
message.update(event.data.sessionID, (draft) => { const current = draft[position]
const current = message.compaction(draft) const failed: Extract<SessionMessageInfo, { type: "compaction"; status: "failed" }> = {
if (current) current.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 break
case "permission.v2.asked": case "permission.v2.asked":
@@ -809,20 +790,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
const result = { const result = {
on: sdk.event.on, on: sdk.event.on,
listen: sdk.event.listen, 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: { session: {
list() { list() {
return Object.values(store.session.info).toSorted((a, b) => b.time.updated - a.time.updated) 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 return store.session.input[sessionID]?.includes(inputID) ?? false
}, },
}, },
compaction(sessionID: string) {
return store.session.compaction[sessionID]
},
async refresh(sessionID: string) { async refresh(sessionID: string) {
setStore("session", "info", sessionID, mutable(await sdk.api.session.get({ sessionID }))) setStore("session", "info", sessionID, mutable(await sdk.api.session.get({ sessionID })))
registerSession(sessionID) registerSession(sessionID)
}, },
message: { message: {
ids(sessionID: string) { 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) { list(sessionID: string) {
return store.session.message[sessionID]?.items ?? [] return store.session.message[sessionID] ?? []
}, },
get(sessionID: string, messageID: string) { 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) const position = messageIndex.get(sessionID)?.get(messageID)
return position === undefined ? undefined : messages?.[position] 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) { async refresh(sessionID: string) {
setStore("session", "message", sessionID, { ...emptyMessages(), loading: true }) setStore("session", "message", sessionID, [])
messageIndex.set(sessionID, new Map()) messageIndex.set(sessionID, new Map())
const response = await sdk.api.message.list({ sessionID, limit: MESSAGE_PAGE_SIZE, order: "desc" }) const messages = mutable(
const items = mutable(response.data).toReversed() (await sdk.api.message.list({ sessionID, limit: 200, order: "desc" })).data,
messageIndex.set(sessionID, new Map(items.map((message, index) => [message.id, index]))) ).toReversed()
setStore("session", "message", sessionID, { messageIndex.set(sessionID, new Map(messages.map((message, index) => [message.id, index])))
items, setStore("session", "message", sessionID, messages)
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,
})
})
}, },
}, },
permission: { permission: {
+3 -8
View File
@@ -5,7 +5,7 @@ import { onCleanup, onMount } from "solid-js"
import { createStore } from "solid-js/store" import { createStore } from "solid-js/store"
import { createSimpleContext } from "./helper" 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 }> } type SDKEventMap = { [Type in V2Event["type"]]: Extract<V2Event, { type: Type }> }
const connectTimeout = 2_000 const connectTimeout = 2_000
@@ -27,11 +27,9 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
status: SDKConnectionStatus status: SDKConnectionStatus
attempt: number attempt: number
error?: string error?: string
connectedOnce: boolean
}>({ }>({
status: "connecting", status: "connecting",
attempt: 0, attempt: 0,
connectedOnce: false,
}) })
let stream: AbortController | undefined let stream: AbortController | undefined
@@ -70,7 +68,7 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
clearTimeout(timeout) clearTimeout(timeout)
attempt = 0 attempt = 0
events.emit(first.value.type, first.value) events.emit(first.value.type, first.value)
setConnection({ status: "connected", attempt: 0, error: undefined, connectedOnce: true }) setConnection({ status: "connected", attempt: 0, error: undefined })
connected() connected()
while (!abort.signal.aborted && !controller.signal.aborted) { while (!abort.signal.aborted && !controller.signal.aborted) {
const event = await iterator.next() const event = await iterator.next()
@@ -98,7 +96,7 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
} }
} }
setConnection({ setConnection({
status: "connecting", status: "reconnecting",
attempt, attempt,
error: error instanceof Error ? error.message : String(error), error: error instanceof Error ? error.message : String(error),
}) })
@@ -136,9 +134,6 @@ export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
error() { error() {
return connection.error return connection.error
}, },
connectedOnce() {
return connection.connectedOnce
},
}, },
reload: props.reload, reload: props.reload,
} }
+30 -81
View File
@@ -198,16 +198,6 @@ export function Session() {
?.id ?.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(() => { const lastAssistant = createMemo(() => {
return messages().findLast((x) => x.type === "assistant") return messages().findLast((x) => x.type === "assistant")
}) })
@@ -251,44 +241,37 @@ export function Session() {
}), }),
) )
createEffect( createEffect(() => {
on( const sessionID = route.sessionID
() => route.sessionID, void (async () => {
(sessionID) => { await Promise.all([
void (async () => { data.session.refresh(sessionID),
if (data.session.message.list(sessionID).length === 0) { data.session.permission.refresh(sessionID),
await Promise.all([ data.session.form.refresh(sessionID),
data.session.refresh(sessionID), ])
data.session.message.refresh(sessionID), const info = data.session.get(sessionID)
data.session.permission.refresh(sessionID), if (!info) {
data.session.form.refresh(sessionID), toast.show({
]) message: `Session not found: ${sessionID}`,
} variant: "error",
const info = data.session.get(sessionID) duration: 5000,
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" })
}) })
}, 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 seeded = false
let scroll: ScrollBoxRenderable let scroll: ScrollBoxRenderable
@@ -344,14 +327,12 @@ export function Session() {
if (!targetID) { if (!targetID) {
scroll.scrollBy(direction === "next" ? scroll.height : -scroll.height) scroll.scrollBy(direction === "next" ? scroll.height : -scroll.height)
if (direction === "prev") loadOlder()
dialog.clear() dialog.clear()
return return
} }
const child = scroll.getChildren().find((c) => c.id === targetID) const child = scroll.getChildren().find((c) => c.id === targetID)
if (child) scroll.scrollBy(child.y - scroll.y - 1) if (child) scroll.scrollBy(child.y - scroll.y - 1)
if (direction === "prev") loadOlder()
dialog.clear() dialog.clear()
} }
@@ -362,24 +343,6 @@ export function Session() {
}, 50) }, 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(() => [ const sessionCommandList = createMemo(() => [
{ {
title: "Share session", title: "Share session",
@@ -585,7 +548,6 @@ export function Session() {
hidden: true, hidden: true,
run: () => { run: () => {
scroll.scrollBy(-scroll.height / 2) scroll.scrollBy(-scroll.height / 2)
loadOlder()
dialog.clear() dialog.clear()
}, },
}, },
@@ -606,7 +568,6 @@ export function Session() {
hidden: true, hidden: true,
run: () => { run: () => {
scroll.scrollBy(-1) scroll.scrollBy(-1)
loadOlder()
dialog.clear() dialog.clear()
}, },
}, },
@@ -627,7 +588,6 @@ export function Session() {
hidden: true, hidden: true,
run: () => { run: () => {
scroll.scrollBy(-scroll.height / 4) scroll.scrollBy(-scroll.height / 4)
loadOlder()
dialog.clear() dialog.clear()
}, },
}, },
@@ -648,7 +608,6 @@ export function Session() {
hidden: true, hidden: true,
run: () => { run: () => {
scroll.scrollTo(0) scroll.scrollTo(0)
loadOlder()
dialog.clear() dialog.clear()
}, },
}, },
@@ -958,9 +917,6 @@ export function Session() {
stickyStart="bottom" stickyStart="bottom"
flexGrow={1} flexGrow={1}
scrollAcceleration={scrollAcceleration()} scrollAcceleration={scrollAcceleration()}
onMouseScroll={(event) => {
if (event.scroll?.direction === "up") void loadOlder()
}}
> >
<For each={rows}> <For each={rows}>
{(row) => ( {(row) => (
@@ -970,13 +926,6 @@ export function Session() {
/> />
)} )}
</For> </For>
<For each={pendingMessages()}>
{(message) => (
<box marginTop={1} flexShrink={0}>
<SessionMessageView message={message} />
</box>
)}
</For>
<BackgroundToolHint messages={messages()} /> <BackgroundToolHint messages={messages()} />
<Show when={session()?.revert?.messageID}> <Show when={session()?.revert?.messageID}>
<RevertMessage <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 { 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" import { useData } from "../../context/data"
export type PartRef = { export type PartRef = {
@@ -25,22 +25,11 @@ export function createSessionRows(sessionID: Accessor<string>) {
const [rows, setRows] = createStore<SessionRow[]>([]) const [rows, setRows] = createStore<SessionRow[]>([])
const revertBoundary = () => data.session.get(sessionID())?.revert?.messageID 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() { function reduce() {
const messages = data.session.message.list(sessionID()) const messages = data.session.message.list(sessionID())
const inputs = new Set(data.session.input.list(sessionID()))
const boundary = revertBoundary() const boundary = revertBoundary()
const visible = boundary ? messages.filter((message) => message.id < boundary) : messages const rows = reduceSessionRows(boundary ? messages.filter((message) => message.id < boundary) : messages, inputs)
const pending = pendingIDs()
const rows = reduceSessionRows(visible.filter((message) => !pending.has(message.id)))
partitionPending(rows, pendingPermissions()) partitionPending(rows, pendingPermissions())
return rows return rows
} }
@@ -63,30 +52,48 @@ export function createSessionRows(sessionID: Accessor<string>) {
}) })
createEffect( createEffect(
on(sessionID, () => { on(sessionID, (id) => {
setRows(reduce()) 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( createEffect(
on(revertBoundary, () => { on(revertBoundary, () => {
setRows(reduce()) setRows(reconcile(reduce()))
}), }),
) )
// Pending inputs and compaction leaving the pending set change history membership.
createEffect( createEffect(
on( on(
() => { () =>
const messages = data.session.message.list(sessionID()) data.session.message.list(sessionID()).flatMap((message) =>
const pending = data.session.input.list(sessionID()).join("\0") message.type === "user"
const compaction = messages ? [
.filter((message) => message.type === "compaction") {
.map((message) => `${message.id}:${message.status}`) id: message.id,
.join("\0") created: message.time.created,
return `${pending}\u0001${compaction}` input: data.session.input.has(sessionID(), message.id),
}, },
() => setRows(reduce()), ]
: 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( setRows(
produce((draft) => { produce((draft) => {
if (draft.some((row) => row.type === "message" && row.messageID === messageID)) return if (draft.some((row) => row.type === "message" && row.messageID === messageID)) return
if (pendingIDs().has(messageID)) return const pending = isPending(messageID)
completePrevious(draft) const message = data.session.message.get(sessionID(), messageID)
draft.push({ type: "message", 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( setRows(
produce((draft) => { produce((draft) => {
if (hasPart(draft, ref)) return if (hasPart(draft, ref)) return
const index = queuedStart(draft)
if (name && exploration(name)) { if (name && exploration(name)) {
const previous = draft.at(-1) const previous = draft[index - 1]
if (previous?.type === "group" && previous.kind === "exploration") { if (previous?.type === "group" && previous.kind === "exploration") {
previous.refs.push(ref) previous.refs.push(ref)
return return
} }
completePrevious(draft) completePrevious(draft, index)
draft.push({ draft.splice(index, 0, {
type: "group", type: "group",
kind: "exploration", kind: "exploration",
refs: [ref], refs: [ref],
@@ -120,8 +131,8 @@ export function createSessionRows(sessionID: Accessor<string>) {
}) })
return return
} }
completePrevious(draft) completePrevious(draft, index)
draft.push({ type: "part", ref }) draft.splice(index, 0, { type: "part", ref })
}), }),
) )
@@ -129,8 +140,9 @@ export function createSessionRows(sessionID: Accessor<string>) {
setRows( setRows(
produce((draft) => { produce((draft) => {
if (draft.some((row) => row.type === "assistant-footer" && row.messageID === messageID)) return if (draft.some((row) => row.type === "assistant-footer" && row.messageID === messageID)) return
completePrevious(draft) const index = queuedStart(draft)
draft.push({ type: "assistant-footer", messageID }) 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 } }) => { const message = (event: { id: string; data: { sessionID: string } }) => {
if (event.data.sessionID === sessionID()) appendMessage(event.id.replace(/^evt_/, "msg_")) 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 = [ const subscriptions = [
data.on("session.prompt.promoted", (event) => { data.on("session.prompt.admitted", input),
if (event.data.sessionID === sessionID()) appendMessage(event.data.inputID) data.on("session.compaction.started", message),
}),
data.on("session.instructions.updated", message), data.on("session.instructions.updated", message),
data.on("session.synthetic", (event) => { data.on("session.synthetic", (event) => {
if (event.data.sessionID === sessionID() && event.data.description?.trim()) 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.shell.started", message),
data.on("session.agent.selected", message), data.on("session.agent.selected", message),
data.on("session.model.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) => { data.on("session.text.delta", (event) => {
if (event.data.sessionID === sessionID()) if (event.data.sessionID === sessionID())
appendPart({ messageID: event.data.assistantMessageID, partID: `text:${event.data.ordinal}` }) appendPart({ messageID: event.data.assistantMessageID, partID: `text:${event.data.ordinal}` })
@@ -197,11 +224,18 @@ export function createSessionRows(sessionID: Accessor<string>) {
return rows return rows
} }
export function reduceSessionRows(messages: SessionMessage[]) { export function reduceSessionRows(messages: SessionMessageInfo[], inputs = new Set<string>()) {
return messages.reduce<SessionRow[]>((rows, message) => { 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 !== "assistant") {
if (message.type === "synthetic" && !message.description?.trim()) return rows 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 }) rows.push({ type: "message", messageID: message.id })
return rows 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 { EventV2 } from "@opencode-ai/core/event"
import { onMount } from "solid-js" import { onMount } from "solid-js"
import { ProjectProvider } from "../../../src/context/project" 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 { DataProvider, useData } from "../../../src/context/data"
import { createSessionRows, type SessionRow } from "../../../src/routes/session/rows" import { createSessionRows, type SessionRow } from "../../../src/routes/session/rows"
import { createApi, createClient, createEventStream, createFetch, directory, json } from "../../fixture/tui-sdk" 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 () => { test("applies absolute usage events to session info", async () => {
const events = createEventStream() const events = createEventStream()
const sessionID = "ses_usage_refresh" const sessionID = "ses_usage_refresh"
@@ -564,9 +494,11 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
}) })
}, events) }, events)
let data!: ReturnType<typeof useData> let data!: ReturnType<typeof useData>
let sdk!: ReturnType<typeof useSDK>
function Probe() { function Probe() {
data = useData() data = useData()
sdk = useSDK()
return <box /> return <box />
} }
@@ -585,41 +517,39 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
try { try {
await wait(() => data.location.model.list()?.[0]?.id === "model-1") await wait(() => data.location.model.list()?.[0]?.id === "model-1")
await wait(() => data.session.status("session-stale") === "running") await wait(() => data.session.status("session-stale") === "running")
expect(data.connection.status()).toBe("connected") expect(sdk.connection.status()).toBe("connected")
expect(data.connection.attempt()).toBe(0) expect(sdk.connection.attempt()).toBe(0)
events.disconnect() events.disconnect()
await wait(() => data.connection.status() === "connecting") await wait(() => sdk.connection.status() === "reconnecting")
expect(data.connection.attempt()).toBe(1) expect(sdk.connection.attempt()).toBe(1)
expect(data.connection.error()).toBe("Event stream disconnected") 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" } } })) resolveActive(json({ data: { "session-new": { type: "running" } } }))
await wait(() => data.location.model.list()?.[0]?.id === "model-2", 4000) await wait(() => data.location.model.list()?.[0]?.id === "model-2", 4000)
await wait(() => data.session.status("session-stale") === "idle") await wait(() => data.session.status("session-stale") === "idle")
expect(data.session.status("session-new")).toBe("running") expect(data.session.status("session-new")).toBe("running")
expect(requests.event).toBe(2) expect(requests.event).toBe(2)
expect(data.connection.status()).toBe("connected") expect(sdk.connection.status()).toBe("connected")
expect(data.connection.attempt()).toBe(0) expect(sdk.connection.attempt()).toBe(0)
expect(data.connection.error()).toBeUndefined() expect(sdk.connection.error()).toBeUndefined()
} finally { } finally {
app.renderer.destroy() 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 events = createEventStream()
const sessionID = "session-promotion" const sessionID = "session-promotion"
const calls = createFetch((url) => { const calls = createFetch((url) => {
if (url.pathname === `/api/session/${sessionID}/message`) return json({ data: [], cursor: {} }) if (url.pathname === `/api/session/${sessionID}/message`) return json({ data: [], cursor: {} })
}, events) }, events)
let rows!: ReturnType<typeof createSessionRows> let rows!: ReturnType<typeof createSessionRows>
let data!: ReturnType<typeof useData>
function Probe() { function Probe() {
rows = createSessionRows(() => sessionID) rows = createSessionRows(() => sessionID)
data = useData()
return <box /> return <box />
} }
@@ -674,8 +604,7 @@ test("keeps pending prompts out of history rows until promoted", async () => {
delivery: "steer", delivery: "steer",
}, },
}) })
await wait(() => data.session.input.has(sessionID, "message-user")) await wait(() => rows.at(-1)?.type === "message")
expect(rows.some((row) => row.type === "message" && row.messageID === "message-user")).toBe(false)
expect(rows.find((row) => row.type === "group")?.completed).toBe(false) expect(rows.find((row) => row.type === "group")?.completed).toBe(false)
emitEvent(events, { emitEvent(events, {
@@ -685,8 +614,7 @@ test("keeps pending prompts out of history rows until promoted", async () => {
durable: durable(sessionID, 3), durable: durable(sessionID, 3),
data: { sessionID, inputID: "message-user" }, data: { sessionID, inputID: "message-user" },
}) })
await wait(() => rows.some((row) => row.type === "message" && row.messageID === "message-user")) await wait(() => rows.find((row) => row.type === "group")?.completed === true)
expect(data.session.input.has(sessionID, "message-user")).toBe(false)
expect(rows.at(-1)).toEqual({ type: "message", messageID: "message-user" }) expect(rows.at(-1)).toEqual({ type: "message", messageID: "message-user" })
} finally { } finally {
app.renderer.destroy() 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() const encoder = new TextEncoder()
let stream: ReadableStreamDefaultController<Uint8Array> | undefined let stream: ReadableStreamDefaultController<Uint8Array> | undefined
const eventResponse = () => const eventResponse = () =>
@@ -773,10 +701,10 @@ test("connectedOnce is false until first connect and persists across disconnect"
const calls = createFetch((url) => { const calls = createFetch((url) => {
if (url.pathname === "/api/event") return eventResponse() if (url.pathname === "/api/event") return eventResponse()
}) })
let data!: ReturnType<typeof useData> let sdk!: ReturnType<typeof useSDK>
function Probe() { function Probe() {
data = useData() sdk = useSDK()
return <box /> return <box />
} }
@@ -794,16 +722,13 @@ test("connectedOnce is false until first connect and persists across disconnect"
try { try {
await wait(() => stream !== undefined) await wait(() => stream !== undefined)
expect(data.connection.status()).toBe("connecting") expect(sdk.connection.status()).toBe("connecting")
expect(data.connection.connectedOnce()).toBe(false)
connect() connect()
await wait(() => data.connection.status() === "connected") await wait(() => sdk.connection.status() === "connected")
expect(data.connection.connectedOnce()).toBe(true)
disconnect() disconnect()
await wait(() => data.connection.status() === "connecting") await wait(() => sdk.connection.status() === "reconnecting")
expect(data.connection.connectedOnce()).toBe(true)
} finally { } finally {
app.renderer.destroy() app.renderer.destroy()
} }
@@ -1443,9 +1368,11 @@ test("adds and dismisses permission requests from live events", async () => {
const events = createEventStream() const events = createEventStream()
const calls = createFetch(undefined, events) const calls = createFetch(undefined, events)
let data!: ReturnType<typeof useData> let data!: ReturnType<typeof useData>
let sdk!: ReturnType<typeof useSDK>
function Probe() { function Probe() {
data = useData() data = useData()
sdk = useSDK()
return <box /> return <box />
} }
@@ -1462,7 +1389,7 @@ test("adds and dismisses permission requests from live events", async () => {
)) ))
try { try {
await wait(() => data.connection.status() === "connected") await wait(() => sdk.connection.status() === "connected")
emitEvent(events, { emitEvent(events, {
id: "evt_permission_asked_1", id: "evt_permission_asked_1",
created: 0, 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: [] }] }) return json({ data: [{ id: "frm_remote", sessionID: "ses_1", mode: "form", fields: [] }] })
}, events) }, events)
let data!: ReturnType<typeof useData> let data!: ReturnType<typeof useData>
let sdk!: ReturnType<typeof useSDK>
function Probe() { function Probe() {
data = useData() data = useData()
sdk = useSDK()
return <box /> return <box />
} }
@@ -1580,7 +1509,7 @@ test("adds, dismisses, and refreshes form requests", async () => {
)) ))
try { try {
await wait(() => data.connection.status() === "connected") await wait(() => sdk.connection.status() === "connected")
emitEvent(events, { emitEvent(events, {
id: "evt_form_created_1", id: "evt_form_created_1",
created: 0, created: 0,
@@ -1751,14 +1680,6 @@ test("settles pending tools when a live failure arrives", async () => {
name: "bash", 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, { emitEvent(events, {
id: "evt_called_1", id: "evt_called_1",
created: 0, 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" }]) 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[] = [ const messages: SessionMessageInfo[] = [
{ type: "user", id: "user-1", text: "Before", time: { created: 1 } }, queued("user-before", "Before", 1),
{ {
type: "compaction", type: "compaction",
id: "compaction", id: "compaction",
status: "completed", status: "running",
reason: "manual", reason: "manual",
summary: "done", summary: "",
recent: "", recent: "",
time: { created: 2 }, time: { created: 2 },
}, },
{ type: "user", id: "user-2", text: "After", time: { created: 3 } }, queued("user-after", "After", 3),
] ]
expect(reduceSessionRows(messages)).toEqual([ expect(reduceSessionRows(messages, new Set(["user-before", "user-after"]))).toEqual([
{ type: "message", messageID: "user-1" },
{ type: "message", messageID: "compaction" }, { type: "message", messageID: "compaction" },
{ type: "message", messageID: "user-2" }, { type: "message", messageID: "user-before" },
{ type: "message", messageID: "user-after" },
]) ])
}) })