mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-15 17:08:21 -04:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c61e93a1f3 | |||
| 68c62774ac | |||
| 8b634e4a58 | |||
| c6156f171c | |||
| 4d2b06f8cf | |||
| bf15c97e4b |
@@ -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
|
||||
|
||||
@@ -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(() =>
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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 })
|
||||
}),
|
||||
)
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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 })
|
||||
}
|
||||
})
|
||||
@@ -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(() => {
|
||||
|
||||
@@ -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: {
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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" },
|
||||
])
|
||||
})
|
||||
|
||||
|
||||
Reference in New Issue
Block a user