$DSH_HOME/config.yaml was an implicit composition layer: if the file existed, every launch applied an arbitrary Loader patch graph over the shipped tree, kept live by a dedicated HMR watcher. Three costs came from the implicitness, not the capability. A patch replaces its target row's whole config, so a file written months ago pins that row to the field set it knew and every default the shipped tree later adds silently stops applying. It competed with the typed settings namespaces llm-deepseek and llm-pi-ai already register, so which one wins was a function of layer order rather than meaning. And the explicit escape hatch it was supposedly redundant with did not exist on every surface: dsh -p, dsh meta, and dsh upgrade all rejected --config, so for them the implicit file was the only composition route at all. Complete the explicit layer first: --config and --config-replace now work on every booting surface. A headless --config-replace tree must still mount a webserver row, because that surface reaches its own agent over the same HTTP gateway the browser uses; AppCLIEntry names that contract in the failure instead of reporting a bare missing service. Then delete the implicit one. PERSONAL_CONFIG_FILENAME, loadPersonalPatches, watchPersonalPatches, and the config-only HMR row mounted for it are gone; a file left at that path is inert, and --dump-config no longer reads the Harness home. --config therefore stops *replacing* the personal overlay and simply *is* the user overlay. No migration: a user who wants the old behavior names the same file (dsh --config ~/.dsh/config.yaml), which a shell alias makes permanent.
128 lines
5.8 KiB
TypeScript
128 lines
5.8 KiB
TypeScript
/**
|
|
* `dsh -p "task"` — headless over the one shared composition: AppCLIEntry
|
|
* boots the same base plus Web overlay as `dsh web` (port 0, so parallel runs never
|
|
* collide), then in-process isomorphic injection (InProcessApiClient over
|
|
* toFetchHandler(ctx.apiProxy), so the full carrier chain — wire
|
|
* serialization, zod, SSE framing — really runs). The printed URL opens the
|
|
* live session in a browser while the task runs. Runs one task turn, prints
|
|
* the final assistant text, exits (completed → 0, else 1).
|
|
*/
|
|
|
|
import { fileURLToPath } from 'node:url'
|
|
import { resolveConfigPath } from '@deepseek-ai/dsh-app-boot'
|
|
import { InProcessApiClient, toFetchHandler } from '@deepseek-ai/dsh-host-apiproxy'
|
|
import type { MuxFrame } from '@deepseek-ai/dsh-host-apiproxy/api'
|
|
import type { RpcRequest, RpcResponse } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
|
|
import type { SessionId } from '@deepseek-ai/dsh-session'
|
|
import { AppCLIEntry } from './app-cli-entry.ts'
|
|
|
|
/** Outcome of one headless turn: aggregated final text plus the turn-end reason kind. */
|
|
interface TurnOutcome {
|
|
text: string
|
|
reason: string
|
|
}
|
|
|
|
/** Unwrap an RpcResponse or fail loud: business errors print and exit 1 (dispose first). */
|
|
async function unwrap<T>(response: RpcResponse<T>, dispose: () => Promise<void>): Promise<T> {
|
|
if (response.result.ok) return response.result.value
|
|
const { code, message } = response.result.error
|
|
process.stderr.write(`dsh: ${code}: ${message}\n`)
|
|
await dispose()
|
|
process.exit(1)
|
|
}
|
|
|
|
/**
|
|
* Consume mux frames until the task turn ends, per the cli-demo runOneShot
|
|
* correlation precedent: anchor on the first turn/start whose trigger kind is
|
|
* 'message' (startup-injected turns are skipped), aggregate text from that
|
|
* turn's assistant/message events (last one wins), finish on its turn/end.
|
|
*/
|
|
async function consumeUntilTurnEnd(frames: AsyncIterable<RpcRequest<MuxFrame>>, sessionId: SessionId): Promise<TurnOutcome> {
|
|
let targetTurn: number | undefined
|
|
let text = ''
|
|
try {
|
|
for await (const frame of frames) {
|
|
const payload = frame.payload
|
|
if (payload.type === 'stream/error') {
|
|
process.stderr.write(`dsh: stream error: ${payload.error.message}\n`)
|
|
return { text, reason: 'error' }
|
|
}
|
|
if (payload.type !== 'session/event' || payload.sessionId !== sessionId) continue
|
|
const event = payload.event
|
|
if (targetTurn === undefined) {
|
|
if (event.type === 'turn/start' && event.data.trigger.kind === 'message') targetTurn = event.data.turn
|
|
continue
|
|
}
|
|
if (event.type === 'assistant/message' && event.data.turn === targetTurn) {
|
|
const joined = event.data.message.content.filter(block => block.type === 'text').map(block => block.text).join('')
|
|
if (joined !== '') text = joined
|
|
}
|
|
if (event.type === 'turn/end' && event.data.turn === targetTurn) {
|
|
return { text, reason: event.data.reason.kind }
|
|
}
|
|
}
|
|
} catch (error: unknown) {
|
|
process.stderr.write(`dsh: event stream failed: ${String(error)}\n`)
|
|
}
|
|
return { text, reason: 'error' }
|
|
}
|
|
|
|
/**
|
|
* Run one headless turn for `task` and exit (completed → 0, else 1). The task
|
|
* is the non-empty prompt the argument adapter parsed from `-p`/`--prompt`
|
|
* (the adapter rejects an empty task, so no guard is needed here).
|
|
* @param task - the prompt text for the single turn.
|
|
* @param config - a `--config` overlay applied over the shipped composition, or `undefined`.
|
|
* @param configReplace - a `--config-replace` tree booted instead of the
|
|
* shipped composition, or `undefined`. It must mount a webserver row: this
|
|
* surface reaches its own agent over the same HTTP gateway the browser uses.
|
|
*/
|
|
export async function runHeadless(task: string, config?: string, configReplace?: string): Promise<void> {
|
|
// A missing DEEPSEEK_API_KEY throws here (plugin load is fail-loud, uncaught by design).
|
|
const entry = new AppCLIEntry({
|
|
configPath: fileURLToPath(new URL('../config/base.cordis.yml', import.meta.url)),
|
|
overlayPath: fileURLToPath(new URL('../config/web.cordis.yml', import.meta.url)),
|
|
...config !== undefined && { extraOverlayPath: resolveConfigPath(config, undefined) },
|
|
...configReplace !== undefined && { configReplacePath: resolveConfigPath(configReplace, undefined) },
|
|
dev: false,
|
|
port: 0,
|
|
})
|
|
const { ctx, port } = await entry.run()
|
|
const dispose = async (): Promise<void> => { await ctx.fiber.dispose() }
|
|
// Signal exits must still dispose the tree: the composition mounts
|
|
// exit-drained plugins (telemetry's queued tail and shutdown marker would
|
|
// otherwise be lost), and Node's default signal exit skips disposal.
|
|
let signalled = false
|
|
const disposeAndExit = (code: number): void => {
|
|
if (signalled) return
|
|
signalled = true
|
|
void dispose().finally(() => { process.exit(code) })
|
|
}
|
|
process.on('SIGTERM', () => { disposeAndExit(143) })
|
|
process.on('SIGINT', () => { disposeAndExit(130) })
|
|
// The headless session is web-observable while it runs (same composition).
|
|
process.stderr.write(`dsh: observing at http://127.0.0.1:${String(port)}\n`)
|
|
const api = new InProcessApiClient(toFetchHandler(ctx.apiProxy))
|
|
|
|
const created = await unwrap(await api.sessions.create({}), dispose)
|
|
|
|
// Open the stream before prompting so no frame is lost — kept in this order
|
|
// even though in-process delivery has no race, so the code survives a move
|
|
// to a remote HTTP carrier unchanged.
|
|
const abort = new AbortController()
|
|
const frames = api.events.mux({}, abort.signal)
|
|
const done = consumeUntilTurnEnd(frames, created.sessionId)
|
|
|
|
await unwrap(await api.sessions.prompt({
|
|
sessionId: created.sessionId,
|
|
mode: 'queue',
|
|
content: [{ type: 'text', text: task }],
|
|
}), dispose)
|
|
|
|
const outcome = await done
|
|
process.stdout.write(outcome.text + '\n')
|
|
abort.abort()
|
|
await dispose()
|
|
process.exit(outcome.reason === 'completed' ? 0 : 1)
|
|
}
|