mirror of
https://github.com/deepseek-ai/deepseek-harness
synced 2026-08-15 21:04:50 +00:00
Six mutation RPCs (create/edit/pause/resume/complete/clear) move into
dsh-host-apiproxy (the PR's host/runtime carrier is gone): goalService()
resolves ctx.get('goals') with a loud absence error, mutateGoal() resolves
the session's agent (agentFor, the command.* implicit-resume precedent) and
acknowledges with the new CAS ref only. GoalError codes ride err.details.
goal.get and the wire GoalView/goalViewSchema are gone: the read side is the
'goal' session projection (whole values on the history tail page and
session/projection frames), so responses never feed client state — the rule
whose absence forced the original PR's write-revision fences.
252 lines
12 KiB
TypeScript
252 lines
12 KiB
TypeScript
// Test-local programmable IApiClient fake (NOT the fixture: fixture is a demo
|
|
// data source on a real clock; behavior tests need per-case responses and
|
|
// deferred-controlled timing). Streams are hand pumps: pushMux/pushHost.
|
|
import type { CommandId } from '@deepseek-ai/dsh-commands/brand'
|
|
import type {
|
|
ClientResponse, CommandDescriptor, HostFrame, IApiClient, ModelTarget, MuxFrame,
|
|
RpcError, RpcReceipt, RpcRequest, RpcResponse, SessionId, SessionModels, SkillEntry,
|
|
WorkspaceId, WorkspaceView,
|
|
} from '@deepseek-ai/dsh-client-connection/client'
|
|
import { RpcId } from '@deepseek-ai/dsh-client-connection/client'
|
|
|
|
/** Programmable-default workspace row (branded id, ISO-ish times). */
|
|
function fakeWorkspace(id: string, over: Partial<WorkspaceView> = {}): WorkspaceView {
|
|
return {
|
|
workspaceId: id as WorkspaceId,
|
|
path: '/f/ws',
|
|
title: 'ws',
|
|
sessionIds: [],
|
|
createdAt: '2026-01-01T00:00:00.000Z',
|
|
updatedAt: '2026-01-01T00:00:00.000Z',
|
|
...over,
|
|
}
|
|
}
|
|
|
|
export interface Deferred<T> {
|
|
promise: Promise<T>
|
|
resolve(value: T): void
|
|
reject(error: unknown): void
|
|
}
|
|
|
|
/** Test-held settlement: the case decides when an RPC lands (history-pending injections etc.). */
|
|
export function deferred<T>(): Deferred<T> {
|
|
let resolve!: (value: T) => void
|
|
let reject!: (error: unknown) => void
|
|
const promise = new Promise<T>((res, rej) => {
|
|
resolve = res
|
|
reject = rej
|
|
})
|
|
return { promise, resolve, reject }
|
|
}
|
|
|
|
let nextRpc = 0
|
|
|
|
export function ok<T>(value: T): RpcResponse<T> {
|
|
return { rpcId: RpcId(`fake-${nextRpc++}`), result: { ok: true, value } }
|
|
}
|
|
|
|
export function err<T>(error: RpcError): RpcResponse<T> {
|
|
return { rpcId: RpcId(`fake-${nextRpc++}`), result: { ok: false, error } }
|
|
}
|
|
|
|
type StreamItem<F> = { kind: 'frame'; envelope: RpcRequest<F> } | { kind: 'end' } | { kind: 'fail'; error: unknown }
|
|
|
|
interface StreamConn<F> {
|
|
feed(item: StreamItem<F>): void
|
|
}
|
|
|
|
export class FakeApiClient implements IApiClient {
|
|
/** Chronological call record: [method, payload]. */
|
|
readonly calls: { method: string; payload: unknown }[] = []
|
|
|
|
// Programmable slots (defaults answer OK-empty); reassign per case.
|
|
onList: (payload: unknown) => Promise<RpcResponse<{ items: never[] }>> = () => Promise.resolve(ok({ items: [] }))
|
|
onCreate: (payload: unknown) => Promise<RpcResponse<{ sessionId: SessionId }>> = () => Promise.resolve(ok({ sessionId: 'fk-new' as SessionId }))
|
|
readonly defaultModel: ModelTarget = { provider: 'deepseek', model: 'deepseek-v4-flash' }
|
|
onHistory: (payload: { sessionId: SessionId; beforeSeq?: number; maxMessages?: number })
|
|
=> Promise<RpcResponse<{ events: never[]; hasMore: boolean }>> =
|
|
() => Promise.resolve(ok({ events: [], hasMore: false }))
|
|
|
|
onModels: (payload: unknown) => Promise<RpcResponse<SessionModels>> = () => Promise.resolve(ok({
|
|
current: this.defaultModel,
|
|
groups: [{
|
|
id: 'deepseek',
|
|
name: 'DeepSeek',
|
|
models: [{ id: 'deepseek-v4-flash', name: 'DeepSeek V4 Flash' }],
|
|
}],
|
|
failures: [],
|
|
}))
|
|
onSelectModel: (payload: { provider: string; model: string }) =>
|
|
Promise<RpcResponse<{ selected: ModelTarget }>> =
|
|
payload => Promise.resolve(ok({ selected: { provider: payload.provider, model: payload.model } }))
|
|
onPrompt: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
|
|
onCancel: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
|
|
onDescribe: (payload: unknown) => Promise<RpcResponse<{ version: string; cwd: string; attachedSessions: number }>> =
|
|
() => Promise.resolve(ok({ version: '0-fake', cwd: '/f', attachedSessions: 0 }))
|
|
onPickDirectory: (payload: unknown) => Promise<RpcResponse<{ path: string | null }>> =
|
|
() => Promise.resolve(ok({ path: null }))
|
|
onOpenPath: (payload: unknown) => Promise<RpcResponse<{ opened: true }>> =
|
|
() => Promise.resolve(ok({ opened: true as const }))
|
|
|
|
private readonly muxConns: StreamConn<MuxFrame>[] = []
|
|
private readonly hostConns: StreamConn<HostFrame>[] = []
|
|
|
|
// Parameters carry local structural annotations: the CI lint lane runs
|
|
// without built lib/, so IApiClient's indexed-access types collapse to any
|
|
// and inferred parameters would trip no-unsafe-argument.
|
|
readonly sessions: IApiClient['sessions'] = {
|
|
list: (payload: unknown) => this.record('session.list', payload, this.onList(payload)),
|
|
create: (payload: unknown) => this.record('session.create', payload, this.onCreate(payload)),
|
|
history: (payload: { sessionId: SessionId; beforeSeq?: number; maxMessages?: number }) =>
|
|
this.record('session.history', payload, this.onHistory(payload)),
|
|
models: (payload: unknown) => this.record('session.models', payload, this.onModels(payload)),
|
|
selectModel: (payload: { provider: string; model: string }) =>
|
|
this.record('session.selectModel', payload, this.onSelectModel(payload)),
|
|
prompt: (payload: unknown) => this.record('session.prompt', payload, this.onPrompt(payload)),
|
|
cancel: (payload: unknown) => this.record('session.cancel', payload, this.onCancel(payload)),
|
|
}
|
|
|
|
readonly host: IApiClient['host'] = {
|
|
describe: (payload: unknown) => this.record('host.describe', payload, this.onDescribe(payload)),
|
|
pickDirectory: (payload: unknown) => this.record('host.pickDirectory', payload, this.onPickDirectory(payload)),
|
|
openPath: (payload: unknown) => this.record('host.openPath', payload, this.onOpenPath(payload)),
|
|
}
|
|
|
|
onWorkspaceList: (payload: unknown) => Promise<RpcResponse<{ items: never[] }>> = () => Promise.resolve(ok({ items: [] }))
|
|
onWorkspaceCreate: (payload: unknown) => Promise<RpcResponse<{ workspace: WorkspaceView; created: boolean }>> =
|
|
() => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws'), created: true }))
|
|
|
|
onWorkspaceRename: (payload: unknown) => Promise<RpcResponse<{ workspace: WorkspaceView }>> =
|
|
() => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws') }))
|
|
|
|
onWorkspaceDelete: (payload: unknown) => Promise<RpcResponse<{ deleted: true }>> =
|
|
() => Promise.resolve(ok({ deleted: true }))
|
|
|
|
onWorkspaceInsertSessionBefore: (payload: unknown) => Promise<RpcResponse<{ workspace: WorkspaceView }>> =
|
|
() => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws') }))
|
|
|
|
readonly workspace: IApiClient['workspace'] = {
|
|
list: (payload: unknown) => this.record('workspace.list', payload, this.onWorkspaceList(payload)),
|
|
create: (payload: unknown) => this.record('workspace.create', payload, this.onWorkspaceCreate(payload)),
|
|
rename: (payload: unknown) => this.record('workspace.rename', payload, this.onWorkspaceRename(payload)),
|
|
delete: (payload: unknown) => this.record('workspace.delete', payload, this.onWorkspaceDelete(payload)),
|
|
insertSessionBefore: (payload: unknown) =>
|
|
this.record('workspace.insertSessionBefore', payload, this.onWorkspaceInsertSessionBefore(payload)),
|
|
}
|
|
|
|
// Payloads stay `unknown` (lint-lane note above); response rows are the real
|
|
// wire shapes so cases can program requires-bearing catalogs and dual-address
|
|
// skill lists without casts.
|
|
onCommandList: (payload: unknown) => Promise<RpcResponse<{ commands: CommandDescriptor[] }>>
|
|
= () => Promise.resolve(ok({ commands: [] }))
|
|
onCommandExecute: (payload: unknown) => Promise<RpcResponse<{ matched: boolean; commandId?: CommandId }>>
|
|
= () => Promise.resolve(ok({ matched: false }))
|
|
onSkillList: (payload: unknown) => Promise<RpcResponse<{ skills: SkillEntry[] }>>
|
|
= () => Promise.resolve(ok({ skills: [] }))
|
|
|
|
readonly commands: IApiClient['commands'] = {
|
|
list: (payload: unknown) => this.record('command.list', payload, this.onCommandList(payload)),
|
|
execute: (payload: unknown) => this.record('command.execute', payload, this.onCommandExecute(payload)),
|
|
}
|
|
|
|
readonly skills: IApiClient['skills'] = {
|
|
list: (payload: unknown) => this.record('skill.list', payload, this.onSkillList(payload)),
|
|
}
|
|
|
|
readonly goals: IApiClient['goals'] = {
|
|
create: payload => this.record('goal.create', payload, Promise.resolve(ok({ ref: { id: 'fake-goal' as never, revision: 1 } }))),
|
|
edit: payload => this.record('goal.edit', payload, Promise.resolve(ok({ ref: { id: 'fake-goal' as never, revision: 1 } }))),
|
|
pause: payload => this.record('goal.pause', payload, Promise.resolve(ok({ ref: { id: 'fake-goal' as never, revision: 1 } }))),
|
|
resume: payload => this.record('goal.resume', payload, Promise.resolve(ok({ ref: { id: 'fake-goal' as never, revision: 1 } }))),
|
|
complete: payload => this.record('goal.complete', payload, Promise.resolve(ok({ ref: { id: 'fake-goal' as never, revision: 1 } }))),
|
|
clear: payload => this.record('goal.clear', payload, Promise.resolve(ok({ cleared: true as const }))),
|
|
}
|
|
|
|
/** When true, streams never fire onOpen (misbehaving-carrier material for the handshake timeout guard). */
|
|
suppressStreamOpen = false
|
|
|
|
/** When true, onOpen callbacks are parked instead of fired; releaseStreamOpens() fires them.
|
|
* Lets a case hold the readiness handshake open (describe done, streams not yet "established"). */
|
|
holdStreamOpen = false
|
|
private heldOpens: (() => void)[] = []
|
|
|
|
releaseStreamOpens(): void {
|
|
const held = this.heldOpens
|
|
this.heldOpens = []
|
|
for (const fire of held) fire()
|
|
}
|
|
|
|
readonly events: IApiClient['events'] = {
|
|
mux: (_payload: unknown, signal: AbortSignal, onOpen?: () => void) => this.openStream(this.muxConns, signal, onOpen),
|
|
host: (_payload: unknown, signal: AbortSignal, onOpen?: () => void) => this.openStream(this.hostConns, signal, onOpen),
|
|
}
|
|
|
|
onRespond: (message: ClientResponse) => Promise<RpcReceipt> = () => Promise.resolve({ accepted: true })
|
|
|
|
respond(message: ClientResponse): Promise<RpcReceipt> {
|
|
return this.record('respond', message, this.onRespond(message))
|
|
}
|
|
|
|
/** Push one mux frame to every open mux stream (rpcId minted unless pinned by the case). */
|
|
pushMux(frame: MuxFrame, rpcId?: string): void {
|
|
for (const conn of [...this.muxConns]) conn.feed({ kind: 'frame', envelope: { rpcId: RpcId(rpcId ?? `push-${nextRpc++}`), payload: frame } })
|
|
}
|
|
|
|
pushHost(frame: HostFrame, rpcId?: string): void {
|
|
for (const conn of [...this.hostConns]) conn.feed({ kind: 'frame', envelope: { rpcId: RpcId(rpcId ?? `push-${nextRpc++}`), payload: frame } })
|
|
}
|
|
|
|
/** End (clean close) or fail (throw) every open stream — reconnect-path material. */
|
|
endStreams(): void {
|
|
for (const conn of [...this.muxConns, ...this.hostConns]) conn.feed({ kind: 'end' })
|
|
}
|
|
|
|
failStreams(error: unknown): void {
|
|
for (const conn of [...this.muxConns, ...this.hostConns]) conn.feed({ kind: 'fail', error })
|
|
}
|
|
|
|
get openMuxCount(): number {
|
|
return this.muxConns.length
|
|
}
|
|
|
|
callsOf(method: string): unknown[] {
|
|
return this.calls.filter(c => c.method === method).map(c => c.payload)
|
|
}
|
|
|
|
private record<T>(method: string, payload: unknown, response: Promise<T>): Promise<T> {
|
|
this.calls.push({ method, payload })
|
|
return response
|
|
}
|
|
|
|
private async *openStream<F>(registry: StreamConn<F>[], signal: AbortSignal, onOpen?: () => void): AsyncGenerator<RpcRequest<F>> {
|
|
const inbox: StreamItem<F>[] = []
|
|
let wake: (() => void) | null = null
|
|
const conn: StreamConn<F> = {
|
|
feed: (item) => {
|
|
inbox.push(item)
|
|
wake?.()
|
|
},
|
|
}
|
|
registry.push(conn)
|
|
if (this.holdStreamOpen && onOpen !== undefined) this.heldOpens.push(onOpen)
|
|
else if (!this.suppressStreamOpen) onOpen?.()
|
|
try {
|
|
while (!signal.aborted) {
|
|
while (inbox.length > 0) {
|
|
const item = inbox.shift() as StreamItem<F>
|
|
if (item.kind === 'end') return
|
|
if (item.kind === 'fail') throw item.error
|
|
yield item.envelope
|
|
}
|
|
await new Promise<void>((resolve) => {
|
|
wake = resolve
|
|
signal.addEventListener('abort', () => { resolve() }, { once: true })
|
|
})
|
|
wake = null
|
|
}
|
|
} finally {
|
|
registry.splice(registry.indexOf(conn), 1)
|
|
}
|
|
}
|
|
}
|