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`.
|
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
|
||||||
|
|||||||
@@ -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(() =>
|
||||||
|
|||||||
@@ -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 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 })
|
||||||
}),
|
}),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
@@ -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())
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
@@ -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 })
|
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(() => {
|
||||||
|
|||||||
@@ -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: {
|
||||||
|
|||||||
@@ -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,
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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" },
|
||||||
])
|
])
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user