Merge branch 'structured-output-subagent-seam' into worktree-dynamic-workflows
# Conflicts: # examples/acp-agent/tests/snapshots/fs-policy-reject/session.jsonl # examples/acp-agent/tests/snapshots/fs-policy-reject/stdout.golden.jsonl
This commit is contained in:
@@ -22,6 +22,10 @@ import { isTurnOpen, lastTurnNumber, runLoop } from './loop.ts'
|
||||
* the agent/* event taxonomy — plugins never need this class.
|
||||
*/
|
||||
export class ReactLoopAgent implements Agent {
|
||||
/**
|
||||
* The queued + steering FIFOs behind {@link send}/{@link steer}. Public so
|
||||
* the driver loop can drain it; {@link cancel} clears it wholesale.
|
||||
*/
|
||||
readonly inbox = new Inbox()
|
||||
|
||||
private _status: AgentStatus = 'idle'
|
||||
@@ -256,6 +260,8 @@ export class ReactLoopAgent implements Agent {
|
||||
* promise (unblocking the idle wait), releases any `whenIdle` waiters, and
|
||||
* aborts the current request if any. The returned `agent.done` promise
|
||||
* resolves once the loop exits.
|
||||
* @returns the disposer — idempotent and infallible (it runs inside the
|
||||
* fiber's LIFO disposal chain, where a throw would skip later disposers).
|
||||
*/
|
||||
start(): () => void {
|
||||
this.done = runLoop(this.ctx, this, {
|
||||
|
||||
@@ -24,30 +24,47 @@ export class Inbox {
|
||||
private steeringMessages: InboxMessage[] = []
|
||||
private wakeup: (() => void) | undefined
|
||||
|
||||
/** Resolves when a queued message arrives (used by the idle loop). */
|
||||
/** True while queued messages are pending — read by the idle wait's fast path and the loop's turn-start checks. */
|
||||
get hasQueued(): boolean {
|
||||
return this.queuedMessages.length > 0
|
||||
}
|
||||
|
||||
/** True while steering messages are pending — read by `cancel()`'s arm gate and the loop's stop-override check. */
|
||||
get hasSteering(): boolean {
|
||||
return this.steeringMessages.length > 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Add a message to the queued FIFO and wake a parked {@link waitForQueued}.
|
||||
* @param message - the message to queue for the next turn start.
|
||||
*/
|
||||
enqueue(message: InboxMessage): void {
|
||||
this.queuedMessages.push(message)
|
||||
this.wakeup?.()
|
||||
}
|
||||
|
||||
/**
|
||||
* Add a message to the steering FIFO. Deliberately no wakeup: steering is
|
||||
* drained between steps of a running turn, never by the idle wait —
|
||||
* `Agent.steer()` on an idle agent falls back to `send()` instead.
|
||||
* @param message - the message to inject between steps of the running turn.
|
||||
*/
|
||||
steer(message: InboxMessage): void {
|
||||
this.steeringMessages.push(message)
|
||||
}
|
||||
|
||||
/** Drain all queued messages (turn start). */
|
||||
/**
|
||||
* Drain all queued messages (turn start).
|
||||
* @returns the drained messages in arrival order; the queued FIFO is left empty.
|
||||
*/
|
||||
drainQueued(): InboxMessage[] {
|
||||
return this.queuedMessages.splice(0)
|
||||
}
|
||||
|
||||
/** Drain all steering messages (between steps). */
|
||||
/**
|
||||
* Drain all steering messages (between steps).
|
||||
* @returns the drained messages in arrival order; the steering FIFO is left empty.
|
||||
*/
|
||||
drainSteering(): InboxMessage[] {
|
||||
return this.steeringMessages.splice(0)
|
||||
}
|
||||
@@ -62,7 +79,12 @@ export class Inbox {
|
||||
this.steeringMessages.length = 0
|
||||
}
|
||||
|
||||
/** Wait until a queued message arrives or `cancel` resolves. */
|
||||
/**
|
||||
* Wait until a queued message arrives or `cancel` resolves.
|
||||
* @param cancel - a promise whose resolution abandons the wait without a
|
||||
* message (the driver loop passes the agent's disposed promise so a parked
|
||||
* loop can exit).
|
||||
*/
|
||||
waitForQueued(cancel: Promise<void>): Promise<void> {
|
||||
if (this.hasQueued) return Promise.resolve()
|
||||
const { promise, resolve } = Promise.withResolvers<void>()
|
||||
|
||||
@@ -29,6 +29,10 @@ declare module 'cordis' {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Plugin config: the agents to create — or resume, via `resumeSessionId` —
|
||||
* declaratively at startup, so a cordis.yml deployment needs no code.
|
||||
*/
|
||||
export interface Config {
|
||||
/** Agents created from configuration at startup. */
|
||||
agents: (AgentOptions & {
|
||||
|
||||
@@ -185,6 +185,9 @@ export interface LoopHandle {
|
||||
* re-enqueue leftover steering as queued ⟵ steering is never stranded
|
||||
* idle (emit agent/status) unless more queued
|
||||
* ```
|
||||
* @param ctx - the plugin context the loop reaches events (agent/…, session/flush) and services (systemPrompt, llm, tools) through.
|
||||
* @param agent - the agent this invocation drives for its whole lifetime (its inbox, session, and options).
|
||||
* @param handle - the bridge to the agent's mutable state: status/abort setters plus the disposal and cancel-marker reads.
|
||||
*/
|
||||
export async function runLoop(ctx: Context, agent: ReactLoopAgent, handle: LoopHandle): Promise<void> {
|
||||
// Per-instance transmission bookkeeping: whether THIS loop instance has
|
||||
@@ -875,7 +878,11 @@ function withoutToolCalls(message: Message): Message {
|
||||
return { ...message, content: message.content.filter(block => block.type !== 'tool-call') }
|
||||
}
|
||||
|
||||
/** The last turn number in a (possibly seeded) session log, or 0. */
|
||||
/**
|
||||
* The last turn number in a (possibly seeded) session log, or 0.
|
||||
* @param session - the session whose log is scanned for the latest `turn/start`.
|
||||
* @returns the latest `turn/start`'s turn number, or 0 when the log has none (the next turn is this plus one).
|
||||
*/
|
||||
export function lastTurnNumber(session: Session): number {
|
||||
const lastStart = session.events.findLast(event => event.type === 'turn/start')
|
||||
return lastStart?.data.turn ?? 0
|
||||
@@ -889,6 +896,8 @@ export function lastTurnNumber(session: Session): number {
|
||||
* returns to idle), so status is not a reliable open-turn signal. Used by
|
||||
* `inject()` to choose between appending into an open turn vs. wrapping the
|
||||
* injection in its own one-shot turn (the turn-enclosure RFC).
|
||||
* @param session - the session whose log is inspected.
|
||||
* @returns true when the log's last turn boundary is a `turn/start` with no matching `turn/end` yet.
|
||||
*/
|
||||
export function isTurnOpen(session: Session): boolean {
|
||||
const last = session.events.findLast(e => e.type === 'turn/start' || e.type === 'turn/end')
|
||||
|
||||
@@ -19,7 +19,10 @@ export interface TransmissionLog {
|
||||
loggedHeader: boolean
|
||||
}
|
||||
|
||||
/** Fresh bookkeeping for a newly-started loop instance. */
|
||||
/**
|
||||
* Fresh bookkeeping for a newly-started loop instance.
|
||||
* @returns state with `loggedHeader` false, so the instance's first request appends an anchoring snapshot.
|
||||
*/
|
||||
export function createTransmissionLog(): TransmissionLog {
|
||||
return { loggedHeader: false }
|
||||
}
|
||||
|
||||
@@ -50,7 +50,11 @@ import type {} from '@deepseek-ai/dsh-system-prompt'
|
||||
/** Identifies one live agent in the registry. */
|
||||
export type AgentId = Branded<'AgentId'>
|
||||
|
||||
/** Brand a string as an {@link AgentId}. */
|
||||
/**
|
||||
* Brand a string as an {@link AgentId}.
|
||||
* @param id - the raw agent id string.
|
||||
* @returns the same string, branded (a compile-time cast — no runtime cost).
|
||||
*/
|
||||
export function AgentId(id: string): AgentId {
|
||||
return id as AgentId
|
||||
}
|
||||
@@ -80,10 +84,22 @@ export interface AgentOptions {
|
||||
model?: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Options for {@link Agent.send}/{@link Agent.steer}/{@link Agent.inject}. An
|
||||
* absent `source` resolves to `{ kind: 'user' }`, so a plugin supplying content
|
||||
* must label itself here or its message is recorded as a user prompt (see
|
||||
* {@link HookContext} on why that label is load-bearing).
|
||||
*/
|
||||
export interface SendOptions {
|
||||
source?: MessageSource
|
||||
}
|
||||
|
||||
/**
|
||||
* An agent's lifecycle state, emitted on every transition as `agent/status`:
|
||||
* `idle` (parked, waiting for queued work), `running` (a turn is in progress),
|
||||
* `disposed` (terminal — no transition leaves it, and `send`/`steer`/`inject`
|
||||
* throw).
|
||||
*/
|
||||
export type AgentStatus = 'idle' | 'running' | 'disposed'
|
||||
|
||||
/**
|
||||
|
||||
555
packages/core/agent/tests/verify-export-jsdoc.spec.ts
Normal file
555
packages/core/agent/tests/verify-export-jsdoc.spec.ts
Normal file
@@ -0,0 +1,555 @@
|
||||
/**
|
||||
* Negative-path tests for the export-surface JSDoc gate
|
||||
* (`scripts/verify-export-jsdoc.ts`).
|
||||
*
|
||||
* The gate's positive half runs against the real tree in CI (`pnpm run
|
||||
* verify-export-jsdoc`, part of doc-sync). What that run cannot prove is that
|
||||
* the walk REJECTS an undocumented surface the way it promises to — and that
|
||||
* every deliberate exemption (heritage members, plugin-protocol slots,
|
||||
* constructors, overload implementations, augmentation bodies, re-exports)
|
||||
* actually holds. These tests drive `collectExportJsdocViolations()` against
|
||||
* synthetic fixture packages, mirroring the gen-cordis-catalog negative
|
||||
* tests.
|
||||
*/
|
||||
|
||||
import { mkdirSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { dirname, join } from 'node:path'
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
import { collectExportJsdocViolations } from '../../../../scripts/verify-export-jsdoc.ts'
|
||||
|
||||
const roots: string[] = []
|
||||
|
||||
afterEach(() => {
|
||||
while (roots.length) rmSync(roots.pop()!, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
/** Write fixture files under `packages/group/fix/src/` and return the scan root. */
|
||||
function fixture(files: Record<string, string>): string {
|
||||
const root = mkdtempSync(join(tmpdir(), 'export-jsdoc-'))
|
||||
roots.push(root)
|
||||
for (const [rel, content] of Object.entries(files)) {
|
||||
const abs = join(root, 'packages', 'group', 'fix', 'src', rel)
|
||||
mkdirSync(dirname(abs), { recursive: true })
|
||||
writeFileSync(abs, content)
|
||||
}
|
||||
return root
|
||||
}
|
||||
|
||||
/** Single-file fixture shorthand: the content becomes `src/index.ts`. */
|
||||
const make = (content: string): string => fixture({ 'index.ts': content })
|
||||
|
||||
describe('verify-export-jsdoc functions and consts', () => {
|
||||
it('accepts a fully documented surface', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/**
|
||||
* Add one to a count.
|
||||
* @param n - the count to bump.
|
||||
* @returns the count plus one.
|
||||
*/
|
||||
export function bump(n: number): number { return n + 1 }
|
||||
|
||||
/**
|
||||
* Fire-and-forget (void needs no @returns).
|
||||
* @param flag - whether to arm.
|
||||
*/
|
||||
export function poke(flag: boolean): void { void flag }
|
||||
|
||||
/** The default retry budget. */
|
||||
export const RETRIES = 3
|
||||
|
||||
/**
|
||||
* Halve a count.
|
||||
* @param n - the count to halve.
|
||||
* @returns the count halved.
|
||||
*/
|
||||
export const halve = (n: number): number => n / 2
|
||||
`))).toEqual([])
|
||||
})
|
||||
|
||||
it('flags an exported function with no JSDoc at all', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'export function bare(): void {}\n',
|
||||
))).toEqual([expect.stringMatching(/exported function 'bare' .* has no JSDoc\./)])
|
||||
})
|
||||
|
||||
it('flags a missing @param and a missing @returns', () => {
|
||||
const violations = collectExportJsdocViolations(make(
|
||||
'/** Docs without tags. */\nexport function f(x: number): number { return x }\n',
|
||||
))
|
||||
expect(violations).toEqual([
|
||||
expect.stringMatching(/exported function 'f' .* is missing @param x\./),
|
||||
expect.stringMatching(/exported function 'f' .* is missing @returns \(return type: number\)\./),
|
||||
])
|
||||
})
|
||||
|
||||
it('flags an unannotated (inferred) return type', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/**\n * Docs.\n * @param x - value.\n */\nexport function f(x: number) { return x }\n',
|
||||
))).toEqual([expect.stringMatching(/no return type annotation/)])
|
||||
})
|
||||
|
||||
it('flags tags-only JSDoc with no description prose', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/**\n * @param x - value.\n */\nexport function f(x: number): void {}\n',
|
||||
))).toEqual([expect.stringMatching(/no description prose above its block tags/)])
|
||||
})
|
||||
|
||||
it('flags a stale @param and a binding-pattern parameter', () => {
|
||||
const violations = collectExportJsdocViolations(make(
|
||||
'/**\n * Docs.\n * @param ghost - not real.\n */\nexport function f({ a }: { a: number }): void {}\n',
|
||||
))
|
||||
expect(violations).toEqual([
|
||||
expect.stringMatching(/parameter '\{ a \}' is a binding pattern; the export surface needs simple identifier parameters/),
|
||||
expect.stringMatching(/@param ghost does not match any parameter \(stale tag\?\)/),
|
||||
])
|
||||
})
|
||||
|
||||
it('exempts a `this` receiver annotation from @param', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/**\n * Docs.\n * @param x - value.\n */\nexport function f(this: object, x: number): void {}\n',
|
||||
))).toEqual([])
|
||||
})
|
||||
|
||||
it('waives @returns for a declarator-annotated const but not an unannotated one', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
type Fn = (x: number) => number
|
||||
/**
|
||||
* Uses the named signature.
|
||||
* @param x - value.
|
||||
*/
|
||||
export const good: Fn = x => x
|
||||
/**
|
||||
* No signature anywhere.
|
||||
* @param x - value.
|
||||
*/
|
||||
export const bad = (x: number) => x
|
||||
`))).toEqual([expect.stringMatching(/exported const 'bad' .* has no return type annotation/)])
|
||||
})
|
||||
|
||||
it('requires description prose on a non-function const', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'export const LIMIT = 10\n',
|
||||
))).toEqual([expect.stringMatching(/exported const 'LIMIT' .* has no JSDoc\./)])
|
||||
})
|
||||
})
|
||||
|
||||
describe('verify-export-jsdoc type-level exports', () => {
|
||||
it('requires description prose on interfaces, type aliases, and enums', () => {
|
||||
const violations = collectExportJsdocViolations(make(
|
||||
'export interface I { a: number }\nexport type T = number\nexport enum E { A }\n',
|
||||
))
|
||||
expect(violations).toEqual([
|
||||
expect.stringMatching(/exported interface 'I' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported type 'T' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported enum 'E' .* has no JSDoc\./),
|
||||
])
|
||||
})
|
||||
|
||||
it('skips `declare module` augmentation bodies (the cordis gate owns them)', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
"declare module 'cordis' {\n interface Events {\n 'fix/x'(): void\n }\n}\nexport {}\n",
|
||||
))).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
describe('verify-export-jsdoc export forms', () => {
|
||||
it('resolves an `export { … }` list to the local declaration', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'function f(): void {}\nexport { f }\n',
|
||||
))).toEqual([expect.stringMatching(/exported function 'f' .* has no JSDoc\./)])
|
||||
})
|
||||
|
||||
it('does not treat a never-exported sibling declarator as surface (review round 2)', () => {
|
||||
// `export { publicValue }` resolves to the whole variable statement; only
|
||||
// the named declarator is surface — the gate must not demand JSDoc for
|
||||
// the private sibling sharing the statement.
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/** The public knob. */\nconst publicValue = 1, privateHelper = 2\nexport { publicValue }\nvoid privateHelper\n',
|
||||
))).toEqual([])
|
||||
})
|
||||
|
||||
it('unions declarators across multiple export lists over one statement (review round 2)', () => {
|
||||
// Two lists each name one declarator of the same undocumented statement:
|
||||
// both are surface (deduplicating on first resolution would drop `b`),
|
||||
// while the never-exported `c` stays out.
|
||||
const violations = collectExportJsdocViolations(make(
|
||||
'const a = 1, b = 2, c = 3\nexport { a }\nexport { b }\nvoid c\n',
|
||||
))
|
||||
expect(violations).toEqual([
|
||||
expect.stringMatching(/exported const 'a' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported const 'b' .* has no JSDoc\./),
|
||||
])
|
||||
})
|
||||
|
||||
it('scopes a default-export identifier to its own declarator (review round 2)', () => {
|
||||
// `export default` of an identifier reaches the statement through the
|
||||
// same name lookup as an export list; the sibling stays private.
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/** The app entry. */\nconst app = 1, scratch = 2\nexport default app\nvoid scratch\n',
|
||||
))).toEqual([])
|
||||
})
|
||||
|
||||
it('reports a re-exported module once, at its defining file', () => {
|
||||
const violations = collectExportJsdocViolations(fixture({
|
||||
'index.ts': "export * from './other.ts'\n",
|
||||
'other.ts': 'export function f(): void {}\n',
|
||||
}))
|
||||
expect(violations).toEqual([expect.stringMatching(/other\.ts:1\) has no JSDoc\./)])
|
||||
})
|
||||
|
||||
it('exempts overload implementations when the signatures are documented', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/**
|
||||
* From a number.
|
||||
* @param x - the number.
|
||||
* @returns its text.
|
||||
*/
|
||||
export function f(x: number): string
|
||||
/**
|
||||
* From a flag.
|
||||
* @param x - the flag.
|
||||
* @returns its text.
|
||||
*/
|
||||
export function f(x: boolean): string
|
||||
export function f(x: number | boolean): string { return String(x) }
|
||||
`))).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
describe('verify-export-jsdoc classes', () => {
|
||||
it('flags an undocumented class, method, property, and accessor', () => {
|
||||
const violations = collectExportJsdocViolations(make(`
|
||||
export class C {
|
||||
state = 1
|
||||
get view(): number { return this.state }
|
||||
run(x: number): number { return x }
|
||||
}
|
||||
`))
|
||||
expect(violations).toEqual([
|
||||
expect.stringMatching(/exported class 'C' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported class property 'C.state' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported class accessor 'C.view' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported class method 'C.run' .* has no JSDoc\./),
|
||||
])
|
||||
})
|
||||
|
||||
it('exempts members declared by an extends/implements heritage type', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/** Seam. */
|
||||
export abstract class Base {
|
||||
/**
|
||||
* Do it.
|
||||
* @param x - input.
|
||||
* @returns output.
|
||||
*/
|
||||
abstract run(x: number): number
|
||||
}
|
||||
/** Iface. */
|
||||
export interface Sized {
|
||||
/** Byte size. */
|
||||
size: number
|
||||
}
|
||||
/** Impl. */
|
||||
export class Impl extends Base implements Sized {
|
||||
size = 0
|
||||
run(x: number): number { return x }
|
||||
}
|
||||
`))).toEqual([])
|
||||
})
|
||||
|
||||
it('skips private/protected/#private members and constructors', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/** Documented. */
|
||||
export class C {
|
||||
#secret = 1
|
||||
private hidden(): void {}
|
||||
protected hook(): void {}
|
||||
constructor(x: number) { void x }
|
||||
}
|
||||
`))).toEqual([])
|
||||
})
|
||||
|
||||
it('exempts plugin-protocol statics but checks other statics', () => {
|
||||
const violations = collectExportJsdocViolations(make(`
|
||||
/** Plugin. */
|
||||
export class C {
|
||||
static Config = { a: 1 }
|
||||
static inject = ['bash']
|
||||
static reusable = true
|
||||
static other = 1
|
||||
}
|
||||
`))
|
||||
expect(violations).toEqual([expect.stringMatching(/exported class property 'C.other' .* has no JSDoc\./)])
|
||||
})
|
||||
|
||||
it("covers a set accessor by the getter's doc", () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/** Documented. */
|
||||
export class C {
|
||||
/** The current width. */
|
||||
get width(): number { return 1 }
|
||||
set width(_v: number) {}
|
||||
}
|
||||
`))).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
describe('verify-export-jsdoc plugin protocol and namespaces', () => {
|
||||
it('exempts top-level plugin-protocol exports', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
export const name = 'fix'
|
||||
export const inject = ['bash']
|
||||
export const reusable = true
|
||||
export const Config = { parse: true }
|
||||
export function apply(): void {}
|
||||
`))).toEqual([])
|
||||
})
|
||||
|
||||
it('recurses into namespaces with qualified names and honors the merge idiom', () => {
|
||||
const violations = collectExportJsdocViolations(make(`
|
||||
/** The plugin class. */
|
||||
export class Fix {}
|
||||
export namespace Fix {
|
||||
export interface Config { a: number }
|
||||
}
|
||||
export namespace Loose {
|
||||
export const x = 1
|
||||
}
|
||||
`))
|
||||
expect(violations).toEqual([
|
||||
expect.stringMatching(/exported interface 'Fix.Config' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported namespace 'Loose' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported const 'Loose.x' .* has no JSDoc\./),
|
||||
])
|
||||
})
|
||||
})
|
||||
|
||||
describe('verify-export-jsdoc fail-closed forms (review round 1)', () => {
|
||||
it('checks the function contract on a non-identifier default export', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/** Doubles. */\nexport default (x: number): number => x * 2\n',
|
||||
))).toEqual([
|
||||
expect.stringMatching(/default export .* is missing @param x\./),
|
||||
expect.stringMatching(/default export .* is missing @returns \(return type: number\)\./),
|
||||
])
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/**\n * Doubles.\n * @param x - the input.\n * @returns twice the input.\n */\nexport default (x: number): number => x * 2\n',
|
||||
))).toEqual([])
|
||||
})
|
||||
|
||||
it('treats an inline function-type annotation as the surface signature', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/** Maps a number. */\nexport declare const f: (x: number) => number\n',
|
||||
))).toEqual([
|
||||
expect.stringMatching(/exported const 'f' .* is missing @param x\./),
|
||||
expect.stringMatching(/exported const 'f' .* is missing @returns \(return type: number\)\./),
|
||||
])
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/**\n * Maps a number.\n * @param x - the input.\n * @returns the mapped value.\n */\nexport const f: (x: number) => number = v => v\n',
|
||||
))).toEqual([])
|
||||
})
|
||||
|
||||
it('recurses into an ambient declare namespace where members export implicitly', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'export declare namespace N {\n function f(x: number): number\n}\n',
|
||||
))).toEqual([
|
||||
expect.stringMatching(/exported namespace 'N' .* has no JSDoc\./),
|
||||
expect.stringMatching(/exported function 'N.f' .* has no JSDoc\./),
|
||||
])
|
||||
})
|
||||
|
||||
it('requires an export-import alias to document itself (its target may be unwalked)', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/** Holder. */\nexport namespace N {\n /** The value. */\n export const x = 1\n}\nexport import y = N.x\n',
|
||||
))).toEqual([expect.stringMatching(/exported alias 'y' .* has no JSDoc\./)])
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'namespace N {\n export const x = 1\n}\n/** Alias surfacing the internal counter. */\nexport import y = N.x\n',
|
||||
))).toEqual([])
|
||||
})
|
||||
|
||||
it('refuses an export-import alias to a callable, class, or namespace target', () => {
|
||||
const refusal = /exported alias 'g' .* aliases a callable, class, or namespace target/
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'namespace N {\n export function f(x: number): number { return x }\n}\n/** Alias. */\nexport import g = N.f\n',
|
||||
))).toEqual([expect.stringMatching(refusal)])
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'namespace N {\n export class C {\n run(x: number): number { return x }\n }\n}\n/** Alias. */\nexport import g = N.C\n',
|
||||
))).toEqual([expect.stringMatching(refusal)])
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'namespace N {\n export namespace Sub {\n export function f(x: number): number { return x }\n }\n}\n/** Alias. */\nexport import g = N.Sub\n',
|
||||
))).toEqual([expect.stringMatching(refusal)])
|
||||
})
|
||||
|
||||
it('classifies wrapped function initializers and default exports (parens, satisfies)', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'type Fn = (x: number) => number\n/** Wrapped. */\nexport const f = (((x: number): number => x)) satisfies Fn\n',
|
||||
))).toEqual([
|
||||
expect.stringMatching(/exported const 'f' .* is missing @param x\./),
|
||||
expect.stringMatching(/exported const 'f' .* is missing @returns \(return type: number\)\./),
|
||||
])
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'type Fn = (x: number) => number\n/** Wrapped. */\nexport default (((x: number): number => x * 2) satisfies Fn)\n',
|
||||
))).toEqual([
|
||||
expect.stringMatching(/default export .* is missing @param x\./),
|
||||
expect.stringMatching(/default export .* is missing @returns \(return type: number\)\./),
|
||||
])
|
||||
})
|
||||
|
||||
it('treats a single-call-signature type literal as the surface signature', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/** Maps. */\nexport declare const f: { (x: number): number }\n',
|
||||
))).toEqual([
|
||||
expect.stringMatching(/exported const 'f' .* is missing @param x\./),
|
||||
expect.stringMatching(/exported const 'f' .* is missing @returns \(return type: number\)\./),
|
||||
])
|
||||
})
|
||||
|
||||
it('refuses a hybrid callable type literal instead of narrowing the check', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'/** Hybrid. */\nexport declare const f: { (x: number): number; flush: () => void }\n',
|
||||
))).toEqual([expect.stringMatching(/exported const 'f'.*callable type literal is not gate-classifiable; extract a named type/)])
|
||||
})
|
||||
|
||||
it('refuses an export-equals assignment instead of failing open', () => {
|
||||
expect(collectExportJsdocViolations(make(
|
||||
'const x = 1\nexport = x\n',
|
||||
))).toEqual([expect.stringMatching(/export-equals assignment .* is not a gate-supported export form/)])
|
||||
})
|
||||
})
|
||||
|
||||
describe('verify-export-jsdoc heritage refinement (review round 1)', () => {
|
||||
it('requires @param for parameters the base member never names', () => {
|
||||
const violations = collectExportJsdocViolations(make(`
|
||||
/** Seam. */
|
||||
export abstract class Base {
|
||||
/**
|
||||
* Do it.
|
||||
* @param x - input.
|
||||
* @returns output.
|
||||
*/
|
||||
abstract run(x: number): number
|
||||
}
|
||||
/** Impl. */
|
||||
export class Impl extends Base {
|
||||
override run(x: number, verbose?: boolean): number { return verbose ? x : -x }
|
||||
}
|
||||
`))
|
||||
expect(violations).toEqual([expect.stringMatching(/exported class method 'Impl.run' .* is missing @param verbose\./)])
|
||||
})
|
||||
|
||||
it('does not exempt a public override of a protected-only base member', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/** Seam. */
|
||||
export abstract class Base {
|
||||
/** Subclass hook. */
|
||||
protected hook(): void {}
|
||||
}
|
||||
/** Impl. */
|
||||
export class Impl extends Base {
|
||||
override hook(): void {}
|
||||
}
|
||||
`))).toEqual([expect.stringMatching(/exported class method 'Impl.hook' .* has no JSDoc\./)])
|
||||
})
|
||||
|
||||
it('treats an underscore-prefixed rename of a base parameter as the same parameter', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/** Seam. */
|
||||
export abstract class Base {
|
||||
/**
|
||||
* Load it.
|
||||
* @param cwd - the working directory to scope the lookup.
|
||||
* @returns the loaded value.
|
||||
*/
|
||||
abstract load(cwd: string): number
|
||||
}
|
||||
/** Impl (ignores cwd). */
|
||||
export class Impl extends Base {
|
||||
load(_cwd: string): number { return 1 }
|
||||
}
|
||||
`))).toEqual([])
|
||||
})
|
||||
|
||||
it('flags a binding-pattern parameter an override adds beyond the base', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/** Seam. */
|
||||
export abstract class Base {
|
||||
/**
|
||||
* Do it.
|
||||
* @param x - input.
|
||||
* @returns output.
|
||||
*/
|
||||
abstract run(x: number): number
|
||||
}
|
||||
/** Impl. */
|
||||
export class Impl extends Base {
|
||||
override run(x: number, { verbose }: { verbose?: boolean } = {}): number { return verbose ? x : -x }
|
||||
}
|
||||
`))).toEqual([expect.stringMatching(/exported class method 'Impl.run' .* is a binding pattern/)])
|
||||
})
|
||||
|
||||
it('revives the @returns duty when an override grows a concrete result over a void base', () => {
|
||||
const voidBase = `
|
||||
/** Seam. */
|
||||
export abstract class Base {
|
||||
/** Do it (fire-and-forget). */
|
||||
abstract run(): void
|
||||
}
|
||||
`
|
||||
expect(collectExportJsdocViolations(make(`${voidBase}
|
||||
/** Impl. */
|
||||
export class Impl extends Base {
|
||||
override run(): number { return 1 }
|
||||
}
|
||||
`))).toEqual([expect.stringMatching(/exported class method 'Impl.run' .* is missing @returns \(return type: number\)\./)])
|
||||
expect(collectExportJsdocViolations(make(`${voidBase}
|
||||
/** Impl. */
|
||||
export class Impl extends Base {
|
||||
/**
|
||||
* Do it and count.
|
||||
* @returns how many were done.
|
||||
*/
|
||||
override run(): number { return 1 }
|
||||
}
|
||||
`))).toEqual([])
|
||||
})
|
||||
|
||||
it('classifies an unannotated override return over a void base via the checker', () => {
|
||||
const voidBase = `
|
||||
/** Seam. */
|
||||
export abstract class Base {
|
||||
/** Do it (fire-and-forget). */
|
||||
abstract run(): void
|
||||
}
|
||||
`
|
||||
expect(collectExportJsdocViolations(make(`${voidBase}
|
||||
/** Impl. */
|
||||
export class Impl extends Base {
|
||||
override run() { return 1 }
|
||||
}
|
||||
`))).toEqual([expect.stringMatching(/exported class method 'Impl.run' .* non-void result its heritage declaration does not document/)])
|
||||
expect(collectExportJsdocViolations(make(`${voidBase}
|
||||
/** Impl (faithful void, no annotation needed). */
|
||||
export class Impl extends Base {
|
||||
override run() {}
|
||||
}
|
||||
`))).toEqual([])
|
||||
})
|
||||
|
||||
it('keeps the full exemption when the base return already carries the @returns duty', () => {
|
||||
expect(collectExportJsdocViolations(make(`
|
||||
/** Seam. */
|
||||
export abstract class Base {
|
||||
/**
|
||||
* Count things.
|
||||
* @returns the count.
|
||||
*/
|
||||
abstract run(): number
|
||||
}
|
||||
/** Impl. */
|
||||
export class Impl extends Base {
|
||||
override run(): number { return 1 }
|
||||
}
|
||||
`))).toEqual([])
|
||||
})
|
||||
})
|
||||
@@ -9,6 +9,7 @@ Creates and holds event-sourced `Session` instances. Persistence is intentionall
|
||||
### Public API
|
||||
|
||||
- `ctx.sessions.create(id?: SessionId, options?: { seed?: SessionEvent[]; meta?: { cwd?: string; parentSession?: SessionId; createdAt?: number; seedLength?: number } }): Session` — Create a session. `options.seed` replays/forks an existing event log; `options.meta` attaches creation metadata (validated absolute `cwd`, `parentSession` lineage, seed boundary) as the immutable `SessionHeader`. The store fills `version`/`id` and defaults `createdAt` to now; a caller reconstructing a persisted session passes the original `createdAt` and persisted `seedLength` to preserve them. Disposed with the calling fiber.
|
||||
- `ctx.sessions.fork(source, boundary?, childSessionId?): Session` — Resolve a live session object or id, select a seed through the inclusive `boundary` event seq (default: current last event), require that boundary to be `turn/end`, and create a live child session with lineage metadata.
|
||||
- `ctx.sessions.get(id: SessionId): Session | undefined`
|
||||
- `ctx.sessions.list(): Session[]`
|
||||
|
||||
@@ -72,9 +73,9 @@ Every `SessionEvent` carries two optional top-level fields (structural metadata)
|
||||
### Extension points
|
||||
|
||||
- Persistence plugins: subscribe to `session/event` (write-behind) and drain on `session/flush` (awaited) and fiber dispose. A durable backend reads the log and reloads it into a live session; the metadata seam (`SessionHeader`, `session.header`) is what such a backend stores beside the log.
|
||||
- Replay/fork: `ctx.sessions.create(id, { seed })` seeds a new session with an existing event log. The surface rebuilds deterministically from `surfaceOp` markers in the seeded events. The seed is validated to the SAME invariants `append` enforces — including that every surface-eligible event (`SurfaceEventType`) carries a `surfaceOp` marker — so a marker-less message event is rejected at construction rather than silently vanishing from `deriveMessages()` (the surface is the sole derivation path) on resume.
|
||||
- Replay/fork: `ctx.sessions.create(id, { seed })` seeds a new session with an existing event log. The surface rebuilds deterministically from `surfaceOp` markers in the seeded events. The seed is validated to the SAME always-on invariants `append` enforces — contiguous seqs, JSON-serializable data, and required `surfaceOp` markers on surface-eligible events — so marker-less message events are rejected at construction rather than silently vanishing from `deriveMessages()`. Broader turn-enclosure checks stay in `dsh-invariants` and persistence repair. Ordinary live-session forks use `ctx.sessions.fork(source, boundary?, childSessionId?)`, where `boundary` is the inclusive source event seq to fork through.
|
||||
- Compaction: the `dsh-compact-basic` plugin appends a `user/message` with `surfaceOp: { op: 'replace', start, end }` to shadow old surface nodes behind a summary checkpoint.
|
||||
|
||||
### What is NOT here (TODO)
|
||||
|
||||
- **Session branching/tree** (pi-style entry tree) — deferred unless needed beyond seed-based forking.
|
||||
- **Session branching/tree** (pi-style entry tree) — deferred unless needed beyond boundary-based `fork()`.
|
||||
|
||||
@@ -153,10 +153,15 @@ export class Session {
|
||||
this.header = header ?? { version: SESSION_FORMAT_VERSION, id, createdAt: Date.now() }
|
||||
}
|
||||
|
||||
/**
|
||||
* The append-only event log, exposed live by reference (readonly-typed, not
|
||||
* a snapshot): later appends are visible through the same array.
|
||||
*/
|
||||
get events(): readonly SessionEvent[] {
|
||||
return this.log
|
||||
}
|
||||
|
||||
/** The next event's sequence number — always the log length (the `seq = log.length` contiguity contract). */
|
||||
get seq(): number {
|
||||
return this.log.length
|
||||
}
|
||||
@@ -175,6 +180,9 @@ export class Session {
|
||||
* declare how it joins the surface, the sole source of derived history) and
|
||||
* rejected by the compiler for non-surface types like `turn/start` or
|
||||
* `assistant/chunk`.
|
||||
* @returns the logged event — its assigned `seq`/`time` plus the SNAPSHOT of
|
||||
* `data` that entered the log, so reading `event.data` back sees the logged
|
||||
* value, never the caller's still-mutable input.
|
||||
* @throws if `data` is not losslessly JSON-serializable (BigInt, function,
|
||||
* symbol, undefined, non-finite number, circular ref, or an exotic object
|
||||
* like Map/Set/Date). The event log is the durable source of truth, so this
|
||||
@@ -362,6 +370,32 @@ export class Session {
|
||||
}
|
||||
}
|
||||
|
||||
/** A fork source: either the live session object or its live store id. */
|
||||
export type SessionForkSource = Session | SessionId
|
||||
|
||||
/**
|
||||
* Rejection codes for session forking: the fork source id is unknown to the
|
||||
* live store (`SESSION_NOT_FOUND`) or names a session object that is not the
|
||||
* store's live instance (`SESSION_NOT_LIVE`); the requested child id is
|
||||
* already taken (`SESSION_ALREADY_EXISTS`); the boundary is not a contiguous
|
||||
* existing seq (`INVALID_BOUNDARY`); or the boundary event is not a
|
||||
* `turn/end` — a fork must cut on a closed turn (`OPEN_TURN`).
|
||||
*/
|
||||
export type SessionForkErrorCode =
|
||||
| 'SESSION_NOT_FOUND'
|
||||
| 'SESSION_NOT_LIVE'
|
||||
| 'SESSION_ALREADY_EXISTS'
|
||||
| 'INVALID_BOUNDARY'
|
||||
| 'OPEN_TURN'
|
||||
|
||||
/** Typed error for session fork rejections. */
|
||||
export class SessionForkError extends Error {
|
||||
constructor(message: string, public readonly code: SessionForkErrorCode) {
|
||||
super(message)
|
||||
this.name = 'SessionForkError'
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* In-memory session store (`ctx.sessions`).
|
||||
*
|
||||
@@ -496,6 +530,92 @@ export class SessionStore extends Service {
|
||||
list(): Session[] {
|
||||
return [...this.store.values()]
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a live child session from a turn-enclosed prefix of a live source.
|
||||
* `boundary` is an inclusive source event seq; omitted means the source's
|
||||
* current last event. A non-empty selected slice must end at `turn/end`.
|
||||
*
|
||||
* @param source - Live source session object or id.
|
||||
* @param boundary - Inclusive source event seq to fork through; omitted means
|
||||
* the source's current last event, and omitted on an empty source forks an
|
||||
* empty child.
|
||||
* @param childSessionId - Optional child session id; omitted delegates to
|
||||
* `SessionStore`'s id policy.
|
||||
* @returns The created live child session.
|
||||
*/
|
||||
fork(source: SessionForkSource, boundary?: number, childSessionId?: SessionId): Session {
|
||||
if (childSessionId !== undefined && this.get(childSessionId) !== undefined) {
|
||||
throw new SessionForkError(`session "${childSessionId}" already exists`, 'SESSION_ALREADY_EXISTS')
|
||||
}
|
||||
const liveSource = this._resolveForkSource(source)
|
||||
const seed = this._forkSeed(liveSource, boundary)
|
||||
return this.create(childSessionId, {
|
||||
seed,
|
||||
meta: {
|
||||
...liveSource.header.cwd !== undefined ? { cwd: liveSource.header.cwd } : {},
|
||||
parentSession: liveSource.id,
|
||||
seedLength: seed.length,
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
private _forkSeed(session: Session, requestedBoundary: number | undefined): SessionEvent[] {
|
||||
const events = session.events
|
||||
const lastEvent = events.at(-1)
|
||||
let boundary: number
|
||||
if (requestedBoundary !== undefined) {
|
||||
boundary = requestedBoundary
|
||||
} else {
|
||||
if (lastEvent === undefined) return []
|
||||
boundary = lastEvent.seq
|
||||
}
|
||||
if (!Number.isSafeInteger(boundary) || boundary < 0) {
|
||||
throw new SessionForkError(
|
||||
`fork boundary for session "${session.id}" must be a non-negative safe integer, got ${String(boundary)}`,
|
||||
'INVALID_BOUNDARY',
|
||||
)
|
||||
}
|
||||
if (boundary >= events.length) {
|
||||
const lastSeq = events.at(-1)?.seq
|
||||
throw new SessionForkError(
|
||||
`fork boundary ${boundary} does not exist in session "${session.id}" (last seq: ${lastSeq ?? 'none'})`,
|
||||
'INVALID_BOUNDARY',
|
||||
)
|
||||
}
|
||||
|
||||
const boundaryEvent = events[boundary]
|
||||
if (boundaryEvent === undefined || boundaryEvent.seq !== boundary) {
|
||||
throw new SessionForkError(
|
||||
`fork boundary ${boundary} does not match a contiguous event seq in session "${session.id}"`,
|
||||
'INVALID_BOUNDARY',
|
||||
)
|
||||
}
|
||||
if (boundaryEvent.type !== 'turn/end') {
|
||||
throw new SessionForkError(
|
||||
`fork boundary ${boundary} in session "${session.id}" must be turn/end, got ${boundaryEvent.type}`,
|
||||
'OPEN_TURN',
|
||||
)
|
||||
}
|
||||
|
||||
return events.slice(0, boundary + 1).map(event => structuredClone(event))
|
||||
}
|
||||
|
||||
private _resolveForkSource(source: SessionForkSource): Session {
|
||||
if (typeof source === 'string') {
|
||||
const session = this.get(source)
|
||||
if (session === undefined) throw new SessionForkError(`session "${source}" not found`, 'SESSION_NOT_FOUND')
|
||||
return session
|
||||
}
|
||||
|
||||
const live = this.get(source.id)
|
||||
if (live === undefined) {
|
||||
throw new SessionForkError(`session "${source.id}" not found`, 'SESSION_NOT_FOUND')
|
||||
}
|
||||
if (live !== source) throw new SessionForkError(`session "${source.id}" is not the live store instance`, 'SESSION_NOT_LIVE')
|
||||
return source
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
export default SessionStore
|
||||
|
||||
@@ -40,6 +40,10 @@ export type JsonValue = null | boolean | number | string | JsonValue[] | { [key:
|
||||
* hiding under a symbol/non-enumerable key cannot make the round-trip lossy.
|
||||
* Getters are invoked during the check (again as `JSON.stringify` would), so the
|
||||
* contract is for plain data records, not objects with side-effecting accessors.
|
||||
* @param value - the candidate event data to test.
|
||||
* @param seen - objects on the current descent path, for circular-reference
|
||||
* detection; the recursion threads it — callers omit it.
|
||||
* @returns true when `value` survives a JSON round-trip losslessly.
|
||||
*/
|
||||
export function isJsonValue(value: unknown, seen: Set<object> = new Set()): boolean {
|
||||
if (value === null) return true
|
||||
|
||||
@@ -54,6 +54,8 @@ import type { SessionEvent } from './types.ts'
|
||||
* Only the LAST turn can be open: the invariants plugin guarantees a `turn/end`
|
||||
* before any later `turn/start`, so an interior open turn is impossible in a
|
||||
* valid committed log. Likewise at most one step is open within that turn.
|
||||
* @param events - the loaded durable log to scan (a valid committed prefix, possibly with a crash tail).
|
||||
* @returns the synthetic closer events to append after `events`, in order; empty when the log is already balanced.
|
||||
*/
|
||||
export function interruptedTurnClosers(events: readonly SessionEvent[]): SessionEvent[] {
|
||||
let openTurn: number | null = null
|
||||
|
||||
@@ -29,6 +29,8 @@ const SURFACE_EVENT_TYPES = new Set<string>([
|
||||
* surface-eligible event that is MISSING its mandatory marker (e.g. validating
|
||||
* a seed/load log); use {@link isSurfaceEvent} to narrow to a fully-formed
|
||||
* {@link SurfaceEvent} with `surfaceOp` present.
|
||||
* @param type - the event type string to test.
|
||||
* @returns true when the type is one of the five message-producing types.
|
||||
*/
|
||||
export function isSurfaceEligibleType(type: string): boolean {
|
||||
return SURFACE_EVENT_TYPES.has(type)
|
||||
@@ -38,6 +40,8 @@ export function isSurfaceEligibleType(type: string): boolean {
|
||||
* Narrow a {@link SessionEvent} to {@link SurfaceEvent}: checks that the
|
||||
* event's `type` is surface-eligible AND that `surfaceOp` is present.
|
||||
* The narrowed type has mandatory {@link SurfaceOp}.
|
||||
* @param event - the event to narrow.
|
||||
* @returns true when the event is surface-eligible and carries its `surfaceOp` marker.
|
||||
*/
|
||||
export function isSurfaceEvent(event: SessionEvent): event is SurfaceEvent {
|
||||
if (!SURFACE_EVENT_TYPES.has(event.type)) return false
|
||||
|
||||
@@ -74,6 +74,12 @@ function nodeDelta(event: SessionEvent): number {
|
||||
* surface successor (`SurfaceNode.next`), or `null` when `end` is the tail —
|
||||
* for the cut after `end`.
|
||||
*
|
||||
* @param nodes - the surface linked list in head→tail order.
|
||||
* @param events - the session log each node's `seq` indexes into.
|
||||
* @param beforeSeq - names the cut (the node it sits immediately before);
|
||||
* `null` — or any seq not on the surface — means the after-tail cut.
|
||||
* @returns true when every `tool-call` before the cut is answered before it
|
||||
* (the unanswered-call depth at the cut is zero).
|
||||
* @throws if the surface prefix drives the unanswered-call depth negative — a
|
||||
* `tool/result` with no preceding open `tool-call` on the surface. That is a
|
||||
* corrupt surface (a structural invariant violation), surfaced loudly here
|
||||
|
||||
@@ -4,7 +4,11 @@ import type { CallId, ContentBlock, LlmCallConfig, MessageSource, StreamChunk, T
|
||||
/** Identifies one session in the store (and its persistence artifacts). */
|
||||
export type SessionId = Branded<'SessionId'>
|
||||
|
||||
/** Brand a string as a {@link SessionId}. */
|
||||
/**
|
||||
* Brand a string as a {@link SessionId}.
|
||||
* @param id - the raw session id string.
|
||||
* @returns the same string, branded (a compile-time cast — no runtime cost).
|
||||
*/
|
||||
export function SessionId(id: string): SessionId {
|
||||
return id as SessionId
|
||||
}
|
||||
@@ -102,6 +106,7 @@ export interface TurnTriggerMap {
|
||||
injection: { kind: 'injection'; source: MessageSource }
|
||||
}
|
||||
|
||||
/** The union over {@link TurnTriggerMap} — what started a turn; plugins extend it by merging variants into the map. */
|
||||
export type TurnTrigger = TurnTriggerMap[keyof TurnTriggerMap]
|
||||
|
||||
/**
|
||||
@@ -156,6 +161,7 @@ export interface TurnEndReasonMap {
|
||||
interrupted: { kind: 'interrupted' }
|
||||
}
|
||||
|
||||
/** The union over {@link TurnEndReasonMap} — why a turn ended; plugins extend it by merging variants into the map. */
|
||||
export type TurnEndReason = TurnEndReasonMap[keyof TurnEndReasonMap]
|
||||
|
||||
/**
|
||||
@@ -361,6 +367,7 @@ export interface SessionEventMap {
|
||||
'request/header-delta': { system?: SystemDelta; tools?: ToolsDelta; config?: LlmCallConfig }
|
||||
}
|
||||
|
||||
/** The appendable event-type keys of {@link SessionEventMap}, plugin-merged extensions included. */
|
||||
export type SessionEventType = keyof SessionEventMap
|
||||
|
||||
/**
|
||||
|
||||
240
packages/core/session/tests/fork.spec.ts
Normal file
240
packages/core/session/tests/fork.spec.ts
Normal file
@@ -0,0 +1,240 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { Context } from 'cordis'
|
||||
import { CallId } from '@deepseek-ai/dsh-llm'
|
||||
import SessionStore, { Session, SessionForkError, SessionId } from '@deepseek-ai/dsh-session'
|
||||
import type { SessionEvent, TurnEndReason } from '@deepseek-ai/dsh-session'
|
||||
|
||||
async function setup(): Promise<{ ctx: Context; sessions: SessionStore }> {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
return { ctx, sessions: ctx.sessions }
|
||||
}
|
||||
|
||||
function appendClosedTurn(
|
||||
session: Session,
|
||||
turn: number,
|
||||
text = `hello ${turn}`,
|
||||
reason: TurnEndReason = { kind: 'completed' },
|
||||
): void {
|
||||
session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('user/message', {
|
||||
content: [{ type: 'text', text }],
|
||||
source: { kind: 'user' },
|
||||
}, { surfaceOp: 'append' })
|
||||
session.append('turn/end', { turn, reason })
|
||||
}
|
||||
|
||||
function appendOpenTurn(session: Session, turn: number): void {
|
||||
session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('user/message', {
|
||||
content: [{ type: 'text', text: `open ${turn}` }],
|
||||
source: { kind: 'user' },
|
||||
}, { surfaceOp: 'append' })
|
||||
}
|
||||
|
||||
function firstUserMessage(events: readonly SessionEvent[]): SessionEvent<'user/message'> {
|
||||
const event = events.find((e): e is SessionEvent<'user/message'> => e.type === 'user/message')
|
||||
if (event === undefined) throw new Error('missing user/message')
|
||||
return event
|
||||
}
|
||||
|
||||
function lastSeq(session: Session): number {
|
||||
const event = session.events.at(-1)
|
||||
if (event === undefined) throw new Error('missing last event')
|
||||
return event.seq
|
||||
}
|
||||
|
||||
describe('SessionStore.fork', () => {
|
||||
it('forks an empty live session as an empty child with lineage metadata', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const source = ctx.sessions.create(SessionId('empty-parent'), { meta: { cwd: '/workspace' } })
|
||||
|
||||
const child = sessions.fork(source, undefined, SessionId('empty-child'))
|
||||
|
||||
expect(child.events).toEqual([])
|
||||
expect(child.header).toMatchObject({
|
||||
id: SessionId('empty-child'),
|
||||
cwd: '/workspace',
|
||||
parentSession: SessionId('empty-parent'),
|
||||
seedLength: 0,
|
||||
})
|
||||
})
|
||||
|
||||
it('forks the latest completed boundary by default and deep-clones seed events', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const source = ctx.sessions.create(SessionId('parent'), { meta: { cwd: '/workspace' } })
|
||||
appendClosedTurn(source, 1, 'hello')
|
||||
|
||||
const child = sessions.fork(SessionId('parent'), undefined, SessionId('child'))
|
||||
|
||||
expect(child.events).toEqual(source.events)
|
||||
expect(child.events).not.toBe(source.events)
|
||||
expect(child.events[1]).not.toBe(source.events[1])
|
||||
firstUserMessage(child.events).data.content[0] = { type: 'text', text: 'child mutation' }
|
||||
expect(firstUserMessage(source.events).data.content).toEqual([{ type: 'text', text: 'hello' }])
|
||||
expect(child.header).toMatchObject({
|
||||
id: SessionId('child'),
|
||||
cwd: '/workspace',
|
||||
parentSession: SessionId('parent'),
|
||||
seedLength: source.events.length,
|
||||
})
|
||||
})
|
||||
|
||||
it('forks from an earlier turn boundary even when the source currently has an open tail', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const source = ctx.sessions.create(SessionId('parent'), { meta: { cwd: '/workspace' } })
|
||||
appendClosedTurn(source, 1, 'first')
|
||||
const firstBoundary = lastSeq(source)
|
||||
appendClosedTurn(source, 2, 'second')
|
||||
appendOpenTurn(source, 3)
|
||||
|
||||
const child = sessions.fork(source, firstBoundary, SessionId('child-from-first'))
|
||||
|
||||
expect(child.events).toEqual(source.events.slice(0, firstBoundary + 1))
|
||||
expect(child.header.seedLength).toBe(firstBoundary + 1)
|
||||
expect(child.deriveMessages()).toEqual([{ role: 'user', content: [{ type: 'text', text: 'first' }] }])
|
||||
})
|
||||
|
||||
it('accepts every turn/end reason as an explicit fork boundary', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const reasons: TurnEndReason[] = [
|
||||
{ kind: 'completed' },
|
||||
{ kind: 'aborted', reason: 'cancelled by user' },
|
||||
{ kind: 'error', step: 1, message: 'model failed', code: 'MODEL' },
|
||||
{ kind: 'disposed' },
|
||||
{ kind: 'max-tokens' },
|
||||
{ kind: 'interrupted' },
|
||||
]
|
||||
|
||||
for (const reason of reasons) {
|
||||
const source = ctx.sessions.create(SessionId(`parent-${reason.kind}`))
|
||||
appendClosedTurn(source, 1, reason.kind, reason)
|
||||
|
||||
const child = sessions.fork(source, lastSeq(source), SessionId(`child-${reason.kind}`))
|
||||
|
||||
expect(child.events.at(-1)?.type).toBe('turn/end')
|
||||
expect(child.header.seedLength).toBe(source.events.length)
|
||||
}
|
||||
})
|
||||
|
||||
it('rejects invalid boundaries before creating a child', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const empty = ctx.sessions.create(SessionId('empty'))
|
||||
expect(() => sessions.fork(empty, 0, SessionId('empty-child')))
|
||||
.toThrow(new SessionForkError('fork boundary 0 does not exist in session "empty" (last seq: none)', 'INVALID_BOUNDARY'))
|
||||
expect(ctx.sessions.get(SessionId('empty-child'))).toBeUndefined()
|
||||
|
||||
const source = ctx.sessions.create(SessionId('parent'))
|
||||
appendClosedTurn(source, 1)
|
||||
expect(() => sessions.fork(source, -1, SessionId('negative')))
|
||||
.toThrow(/non-negative safe integer/)
|
||||
expect(() => sessions.fork(source, 0.5, SessionId('fraction')))
|
||||
.toThrow(/non-negative safe integer/)
|
||||
expect(() => sessions.fork(source, Number.MAX_SAFE_INTEGER + 1, SessionId('unsafe')))
|
||||
.toThrow(/non-negative safe integer/)
|
||||
expect(() => sessions.fork(source, source.seq, SessionId('past-end')))
|
||||
.toThrow(new SessionForkError(`fork boundary ${source.seq} does not exist in session "parent" (last seq: ${source.seq - 1})`, 'INVALID_BOUNDARY'))
|
||||
})
|
||||
|
||||
it('rejects a corrupted live source whose array index no longer matches event seq', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const source = ctx.sessions.create(SessionId('corrupt-parent'))
|
||||
appendClosedTurn(source, 1)
|
||||
const mutableLog = (source as unknown as { log: SessionEvent[] }).log
|
||||
mutableLog[2] = { ...mutableLog[2]!, seq: 99 }
|
||||
|
||||
expect(() => sessions.fork(source, 2, SessionId('corrupt-child')))
|
||||
.toThrow(new SessionForkError('fork boundary 2 does not match a contiguous event seq in session "corrupt-parent"', 'INVALID_BOUNDARY'))
|
||||
expect(ctx.sessions.get(SessionId('corrupt-child'))).toBeUndefined()
|
||||
})
|
||||
|
||||
it('rejects an unknown live session id', async () => {
|
||||
const { sessions } = await setup()
|
||||
|
||||
expect(() => sessions.fork(SessionId('missing')))
|
||||
.toThrow(new SessionForkError('session "missing" not found', 'SESSION_NOT_FOUND'))
|
||||
})
|
||||
|
||||
it('rejects a detached Session object that is not live in ctx.sessions', async () => {
|
||||
const { sessions } = await setup()
|
||||
const detached = new Session(SessionId('detached'))
|
||||
|
||||
expect(() => sessions.fork(detached))
|
||||
.toThrow(new SessionForkError('session "detached" not found', 'SESSION_NOT_FOUND'))
|
||||
})
|
||||
|
||||
it('rejects a stale Session object whose id is live on a different instance', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
ctx.sessions.create(SessionId('same-id'))
|
||||
const stale = new Session(SessionId('same-id'))
|
||||
|
||||
expect(() => sessions.fork(stale))
|
||||
.toThrow(new SessionForkError('session "same-id" is not the live store instance', 'SESSION_NOT_LIVE'))
|
||||
})
|
||||
|
||||
it('rejects selected slices whose boundary is inside an open turn', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const cases: [string, (session: Session) => number][] = [
|
||||
['turn/start', (session) => {
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
return lastSeq(session)
|
||||
}],
|
||||
['step/start', (session) => {
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('step/start', { turn: 1, step: 1 })
|
||||
return lastSeq(session)
|
||||
}],
|
||||
['user/message', (session) => {
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('user/message', { content: [{ type: 'text', text: 'open' }], source: { kind: 'user' } }, { surfaceOp: 'append' })
|
||||
return lastSeq(session)
|
||||
}],
|
||||
['assistant/message', (session) => {
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('step/start', { turn: 1, step: 1 })
|
||||
session.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 'partial' }] }, { surfaceOp: 'append' })
|
||||
return lastSeq(session)
|
||||
}],
|
||||
['tool/call', (session) => {
|
||||
const callId = CallId('call-open')
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('step/start', { turn: 1, step: 1 })
|
||||
session.append('assistant/message', {
|
||||
turn: 1,
|
||||
step: 1,
|
||||
content: [{ type: 'tool-call', id: callId, name: 'bash', arguments: '{}' }],
|
||||
}, { surfaceOp: 'append' })
|
||||
session.append('tool/call', { turn: 1, step: 1, callId, name: 'bash', arguments: '{}' })
|
||||
return lastSeq(session)
|
||||
}],
|
||||
]
|
||||
|
||||
for (const [lastType, build] of cases) {
|
||||
const source = ctx.sessions.create(SessionId(`open-${lastType}`))
|
||||
const boundary = build(source)
|
||||
|
||||
expect(() => sessions.fork(source, boundary))
|
||||
.toThrow(new SessionForkError(`fork boundary ${boundary} in session "open-${lastType}" must be turn/end, got ${lastType}`, 'OPEN_TURN'))
|
||||
}
|
||||
})
|
||||
|
||||
it('rejects a child session id that is already live with a typed fork error', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const source = ctx.sessions.create(SessionId('parent'))
|
||||
appendClosedTurn(source, 1)
|
||||
ctx.sessions.create(SessionId('child'))
|
||||
|
||||
expect(() => sessions.fork(source, undefined, SessionId('child')))
|
||||
.toThrow(new SessionForkError('session "child" already exists', 'SESSION_ALREADY_EXISTS'))
|
||||
})
|
||||
|
||||
it('rejects a duplicate child session id before validating the boundary', async () => {
|
||||
const { ctx, sessions } = await setup()
|
||||
const source = ctx.sessions.create(SessionId('open-parent'))
|
||||
source.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
ctx.sessions.create(SessionId('child'))
|
||||
|
||||
expect(() => sessions.fork(source, undefined, SessionId('child')))
|
||||
.toThrow(new SessionForkError('session "child" already exists', 'SESSION_ALREADY_EXISTS'))
|
||||
})
|
||||
})
|
||||
@@ -61,8 +61,10 @@ describe('Session', () => {
|
||||
|
||||
it('replays identically from a seeded event log', () => {
|
||||
const original = new Session(SessionId('s3'))
|
||||
original.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
original.append('user/message', { content: [{ type: 'text', text: 'q' }], source: { kind: 'user' } }, { surfaceOp: 'append' })
|
||||
original.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 'a' }] }, { surfaceOp: 'append' })
|
||||
original.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
|
||||
const replayed = new Session(SessionId('s3-replay'), [...original.events])
|
||||
expect(replayed.deriveMessages()).toEqual(original.deriveMessages())
|
||||
@@ -423,7 +425,9 @@ describe('todo/write event', () => {
|
||||
|
||||
it('round-trips through a seeded replay identically (durable, no surfaceOp needed)', () => {
|
||||
const original = new Session(SessionId('t4'))
|
||||
original.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
original.append('todo/write', { todos: [{ content: 'only', status: 'completed' }] })
|
||||
original.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
// Seeding a non-surface event with no surfaceOp must not throw.
|
||||
const replayed = new Session(SessionId('t4-replay'), [...original.events])
|
||||
expect(replayed.events.findLast(e => e.type === 'todo/write')!.data.todos)
|
||||
|
||||
@@ -110,6 +110,7 @@ const VARIABLE_NAME = /^[a-z][a-z0-9_]*$/
|
||||
/** A complete `{{...}}` reference group at the scan position (validated after). */
|
||||
const GROUP_AT = /^\{\{([^{}]*)\}\}/
|
||||
|
||||
/** Plugin config: the deployment-authored fragment of the system prompt (see {@link Config.persona} for its contract). */
|
||||
export interface Config {
|
||||
/**
|
||||
* The deployment's persona — the ONE deployment-authored fragment of the
|
||||
@@ -138,6 +139,10 @@ export interface Config {
|
||||
* while a `}}` still follows (e.g. `{{{model}}}`, `{{a{b}}`) all throw. A
|
||||
* lone `{{` with no `}}` anywhere after it is ordinary prose and passes
|
||||
* through verbatim. Substituted values are never re-scanned.
|
||||
* @param assembly - the assembly to render (typically the awaited result of
|
||||
* {@link SystemPrompt.assemble}); only `sections` and `variables` are read.
|
||||
* @returns the full system prompt text; `''` when every section renders empty
|
||||
* (the caller then sends no system prompt at all).
|
||||
*/
|
||||
export function renderPrompt(assembly: PromptAssembly): string {
|
||||
return assembly.sections
|
||||
|
||||
@@ -256,6 +256,8 @@ function checkSchemaNode(node: unknown, path: string, violations: string[], seen
|
||||
* (`UNSUPPORTED_SCHEMA`) listing EVERY violation; returns (and narrows) on
|
||||
* success. Call this at the seam boundary, before any child is created.
|
||||
* @param schema - the caller-supplied schema (unknown until asserted).
|
||||
* @returns nothing — the assertion signature narrows `schema` to
|
||||
* {@link StructuredOutputSchema} in the caller's scope on normal return.
|
||||
*/
|
||||
export function assertSupportedOutputSchema(schema: unknown): asserts schema is StructuredOutputSchema {
|
||||
const violations: string[] = []
|
||||
|
||||
@@ -155,6 +155,9 @@ export interface JsonSchemaObject {
|
||||
* `properties`, `required` array).
|
||||
*
|
||||
* This is a plain function — no schemastery or other framework dependency.
|
||||
* @param spec - the author-facing per-property schema to convert.
|
||||
* @returns the wire-format JSON Schema; the top-level `required` array is
|
||||
* omitted entirely when no property is marked required.
|
||||
*/
|
||||
export function schemaSpecToJsonSchema(spec: SchemaSpec): JsonSchemaObject {
|
||||
const properties: Record<string, unknown> = {}
|
||||
@@ -269,6 +272,9 @@ function checkSpec(spec: SchemaSpec, value: unknown, path: string): string[] {
|
||||
* keys are allowed (no `additionalProperties: false`); `default` is not
|
||||
* applied; an `object`/`array` prop without `properties`/`items` only
|
||||
* type-checks; `enum` is membership (strings only).
|
||||
* @param spec - the declared parameter schema to validate against.
|
||||
* @param args - the model-generated arguments, however malformed.
|
||||
* @returns the violation messages in declaration order; empty means valid.
|
||||
*/
|
||||
export function validateArgs(spec: SchemaSpec, args: unknown): string[] {
|
||||
return checkSpec(spec, args, '')
|
||||
@@ -340,6 +346,13 @@ export interface DefineToolOptions<S extends SchemaSpec> {
|
||||
* Raw JSON-Schema tool definitions (from MCP servers) are still accepted
|
||||
* by `ToolRegistry.register()` directly — `defineTool` is sugar for
|
||||
* first-party plugin authors.
|
||||
* @param options - the tool's name, description, typed parameter schema,
|
||||
* execute body, and optional presenters.
|
||||
* @returns a registry-ready {@link ToolDefinition}: its `execute` validates the
|
||||
* raw args first (throwing {@link ToolArgsError} on mismatch, which the
|
||||
* registry turns into an isError result), and its presenters validate softly
|
||||
* (returning undefined on mismatch, since replay may feed them older-schema
|
||||
* args).
|
||||
*/
|
||||
export function defineTool<S extends SchemaSpec>(options: DefineToolOptions<S>): ToolDefinition {
|
||||
// Object-literal execute methods don't use `this`; the reference is safe.
|
||||
|
||||
Reference in New Issue
Block a user