feat(invariants): add package-owned service seam
This commit is contained in:
@@ -1,409 +1,199 @@
|
||||
/**
|
||||
* Runtime listeners that fail loudly when cross-event contracts are broken:
|
||||
* turn and step nesting, scoped dispatch, status transitions, and request
|
||||
* reconstruction. The plugin has no environment guard and is active wherever
|
||||
* mounted, including the default `dsh-agent-spine-demo` bundle; custom compositions
|
||||
* may omit it. Sessions own immutable, surface-valid event storage; this plugin
|
||||
* checks only relationships that event acceptance cannot express.
|
||||
* Configurable registry for package-owned runtime invariant contributions.
|
||||
* Packages register checks from optional `./invariant` companion plugins;
|
||||
* ordinary package entrypoints stay independent of diagnostics.
|
||||
*
|
||||
* @module @deepseek-ai/dsh-invariants
|
||||
*/
|
||||
|
||||
import type { Context } from 'cordis'
|
||||
import { carrierKeyOf, isScopeCarrier } from '@deepseek-ai/dsh-scope'
|
||||
import { assertNever, HarnessError } from '@deepseek-ai/dsh-llm'
|
||||
import type { CallId, GenerateOptions } from '@deepseek-ai/dsh-llm'
|
||||
import type { Agent, AgentStatus } from '@deepseek-ai/dsh-agent'
|
||||
import { Session, SessionId, foldRequestHeader } from '@deepseek-ai/dsh-session'
|
||||
import type { SessionEvent } from '@deepseek-ai/dsh-session'
|
||||
import { scopedSubjectResolverFor } from './scoped-events.generated.ts'
|
||||
import { Context, Service } from 'cordis'
|
||||
import type { Inject } from 'cordis'
|
||||
import z from 'schemastery'
|
||||
import type Schema from 'schemastery'
|
||||
|
||||
export const name = 'invariants'
|
||||
export const inject = ['sessions']
|
||||
|
||||
/**
|
||||
* Thrown when a harness event-contract invariant is violated. Extends
|
||||
* {@link HarnessError} (`code: 'INVARIANT'`) so a violation is routable like
|
||||
* any other harness failure.
|
||||
*/
|
||||
export class InvariantError extends HarnessError {
|
||||
constructor(message: string) {
|
||||
super(`invariant violated: ${message}`, 'INVARIANT')
|
||||
this.name = 'InvariantError'
|
||||
}
|
||||
/** Runtime invariant selection configured on the service plugin. */
|
||||
export interface Config {
|
||||
/** Global switch; defaults to `true`. */
|
||||
readonly enabled?: boolean
|
||||
/** Case-sensitive JavaScript regex sources that admit package names; empty admits all. */
|
||||
readonly package_allowlist?: string[]
|
||||
/** Case-sensitive JavaScript regex sources that exclude package names after allowlist matching. */
|
||||
readonly package_blocklist?: string[]
|
||||
}
|
||||
|
||||
/** Per-session bookkeeping for the session-log invariants. */
|
||||
interface SessionTrace {
|
||||
/** Highest `seq` seen so far (must strictly increase). */
|
||||
lastSeq: number
|
||||
/** Open turn number, or null between turns. */
|
||||
openTurn: number | null
|
||||
/** Open step within the current turn, or null between steps. */
|
||||
openStep: number | null
|
||||
/** The next turn number expected in this session log. */
|
||||
nextTurn: number
|
||||
/** The next step number expected within the open turn. */
|
||||
nextStep: number
|
||||
/**
|
||||
* Throw a package-attributed invariant failure.
|
||||
* @param message - violated package contract without the standard prefix.
|
||||
* @returns never because reporting a violation throws.
|
||||
*/
|
||||
export type InvariantFailure = (message: string) => never
|
||||
|
||||
/** Install one package's listeners into the registration's child context. */
|
||||
export interface InvariantInstaller {
|
||||
/**
|
||||
* Tool-call ids issued in the OPEN step awaiting a result. Cleared at
|
||||
* `step/end` — a result must arrive in the same step as its call.
|
||||
* Install the package contribution.
|
||||
* @param ctx - child context owned by this invariant registration.
|
||||
* @param fail - reporter bound to the registering package name.
|
||||
* @returns nothing after synchronous listener installation completes.
|
||||
*/
|
||||
pendingCalls: Set<CallId>
|
||||
(ctx: Context, fail: InvariantFailure): void
|
||||
/** Services the child installer fiber may access. */
|
||||
readonly inject?: Inject
|
||||
}
|
||||
|
||||
/** One accepted event's deferred mutation of a live session trace. */
|
||||
interface SessionTraceTransition {
|
||||
/** Scalar state after the event commits. */
|
||||
scalars: Pick<SessionTrace, 'lastSeq' | 'openTurn' | 'openStep' | 'nextTurn' | 'nextStep'>
|
||||
/** The event's mutation of the open step's pending call set. */
|
||||
pendingCalls:
|
||||
| { kind: 'none' }
|
||||
| { kind: 'add' | 'delete'; callId: CallId }
|
||||
| { kind: 'clear' }
|
||||
/** Internal effect shape used to join child startup before a companion loads. */
|
||||
interface PendingInvariantRegistration extends PromiseLike<() => void> {
|
||||
(): void | Promise<void>
|
||||
}
|
||||
|
||||
/** Assert that a step-scoped event names the currently open turn and step. */
|
||||
function requireOpenStep(trace: SessionTrace, kind: string, turn: number, step: number): void {
|
||||
if (trace.openTurn !== turn || trace.openStep !== step) {
|
||||
throw new InvariantError(
|
||||
`${kind} names turn ${turn}/step ${step} but open is turn ${trace.openTurn}/step ${trace.openStep}`,
|
||||
)
|
||||
/** Thrown when a package-owned runtime invariant is violated. */
|
||||
export class InvariantError extends Error {
|
||||
/** Stable machine-readable invariant failure code. */
|
||||
readonly code = 'INVARIANT' as const
|
||||
/** Full npm package name that owns the violated invariant. */
|
||||
readonly packageName: string
|
||||
|
||||
/**
|
||||
* Construct a package-attributed invariant failure.
|
||||
* @param packageName - full npm package name that registered the check.
|
||||
* @param message - violated contract, without the standard error prefix.
|
||||
*/
|
||||
constructor(packageName: string, message: string) {
|
||||
super(`invariant violated by "${packageName}": ${message}`)
|
||||
this.name = 'InvariantError'
|
||||
this.packageName = packageName
|
||||
}
|
||||
}
|
||||
|
||||
/** Validate one candidate event without mutating the committed session trace. */
|
||||
function validateEvent(trace: SessionTrace, event: SessionEvent): SessionTraceTransition {
|
||||
// seq is strictly monotonic — the spine of replay equivalence. lastSeq
|
||||
// starts at -1, so the first event (seq 0) passes.
|
||||
if (event.seq <= trace.lastSeq) {
|
||||
throw new InvariantError(`seq must strictly increase: saw ${event.seq} after ${trace.lastSeq}`)
|
||||
}
|
||||
let openTurn = trace.openTurn
|
||||
let openStep = trace.openStep
|
||||
let nextTurn = trace.nextTurn
|
||||
let nextStep = trace.nextStep
|
||||
let pendingCalls: SessionTraceTransition['pendingCalls'] = { kind: 'none' }
|
||||
|
||||
// Boundary/step-scoped events have explicit cases; every OTHER event type —
|
||||
// including plugin-added (merge-extensible) SessionEventMap keys — is caught
|
||||
// by the `default` and must be turn-enclosed (the turn-enclosure RFC). No assertNever: an
|
||||
// unknown variant is valid, not a compile error.
|
||||
switch (event.type) {
|
||||
case 'turn/start': {
|
||||
if (trace.openTurn !== null) {
|
||||
throw new InvariantError(`turn/start ${event.data.turn} while turn ${trace.openTurn} is still open`)
|
||||
}
|
||||
// Current sessions replay full logs, so numbering starts at 1 and remains
|
||||
// contiguous. If a future compaction/fork stores a partial log, it must
|
||||
// seed `nextTurn` from retained metadata before this check runs.
|
||||
if (event.data.turn !== trace.nextTurn) {
|
||||
throw new InvariantError(`turn/start expected turn ${trace.nextTurn}, got ${event.data.turn}`)
|
||||
}
|
||||
openTurn = event.data.turn
|
||||
nextStep = 1
|
||||
break
|
||||
}
|
||||
case 'turn/end': {
|
||||
if (trace.openTurn !== event.data.turn) {
|
||||
throw new InvariantError(`turn/end ${event.data.turn} does not match open turn ${trace.openTurn}`)
|
||||
}
|
||||
if (trace.openStep !== null) {
|
||||
throw new InvariantError(`turn/end ${event.data.turn} while step ${trace.openStep} is still open`)
|
||||
}
|
||||
openTurn = null
|
||||
nextTurn += 1
|
||||
break
|
||||
}
|
||||
case 'step/start': {
|
||||
if (trace.openTurn !== event.data.turn) {
|
||||
throw new InvariantError(`step/start in turn ${event.data.turn} but open turn is ${trace.openTurn}`)
|
||||
}
|
||||
if (trace.openStep !== null) {
|
||||
throw new InvariantError(`step/start ${event.data.step} while step ${trace.openStep} is still open`)
|
||||
}
|
||||
// Steps are checked under the same full-log assumption as turns above.
|
||||
if (event.data.step !== trace.nextStep) {
|
||||
throw new InvariantError(`step/start expected step ${trace.nextStep} in turn ${event.data.turn}, got ${event.data.step}`)
|
||||
}
|
||||
openStep = event.data.step
|
||||
break
|
||||
}
|
||||
case 'step/end': {
|
||||
requireOpenStep(trace, 'step/end', event.data.turn, event.data.step)
|
||||
// A result must arrive in the step that issued the call; orphan calls
|
||||
// (a step that errored before its result) do not carry to the next step.
|
||||
pendingCalls = { kind: 'clear' }
|
||||
openStep = null
|
||||
nextStep += 1
|
||||
break
|
||||
}
|
||||
case 'assistant/chunk': {
|
||||
requireOpenStep(trace, 'assistant/chunk', event.data.turn, event.data.step)
|
||||
break
|
||||
}
|
||||
case 'assistant/message': {
|
||||
requireOpenStep(trace, 'assistant/message', event.data.turn, event.data.step)
|
||||
break
|
||||
}
|
||||
case 'tool/call': {
|
||||
requireOpenStep(trace, 'tool/call', event.data.turn, event.data.step)
|
||||
pendingCalls = { kind: 'add', callId: event.data.callId }
|
||||
break
|
||||
}
|
||||
case 'tool/result': {
|
||||
requireOpenStep(trace, 'tool/result', event.data.turn, event.data.step)
|
||||
// A result needs a prior matching call in the same step. (The converse
|
||||
// does NOT hold: a call may have no result — a throwing tool-execution
|
||||
// pipeline step ends the turn with no tool/result, which is legal.)
|
||||
const syntheticInterrupted = event.data.isError && event.data.error?.code === 'interrupted'
|
||||
if (!trace.pendingCalls.has(event.data.callId) && !syntheticInterrupted) {
|
||||
throw new InvariantError(`tool/result for ${event.data.callId} with no prior tool/call in this step`)
|
||||
}
|
||||
pendingCalls = { kind: 'delete', callId: event.data.callId }
|
||||
break
|
||||
}
|
||||
// Turn-enclosure (the turn-enclosure RFC): EVERY session event not handled by a boundary
|
||||
// case above must sit inside an open turn. The durable session log uses the
|
||||
// turn as its commit/replay boundary (the JSONL backend treats anything
|
||||
// after the last turn/end as a crash tail), so a bare event between turns is
|
||||
// silently dropped on reload. The loop records queued user messages after
|
||||
// turn/start, and an idle agent.inject() wraps its context/message in a
|
||||
// one-shot turn. A `default`
|
||||
// (not an enumerated list) is deliberate: SessionEventMap is
|
||||
// merge-extensible, so a PLUGIN-added event type appended while idle must
|
||||
// also fail here rather than fall through and be dropped on resume.
|
||||
default: {
|
||||
if (trace.openTurn === null) {
|
||||
throw new InvariantError(`${event.type} appended outside any open turn (every event must be turn-enclosed)`)
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
return {
|
||||
scalars: { lastSeq: event.seq, openTurn, openStep, nextTurn, nextStep },
|
||||
pendingCalls,
|
||||
declare module 'cordis' {
|
||||
interface Context {
|
||||
invariants: InvariantService
|
||||
}
|
||||
}
|
||||
|
||||
/** Apply one already-validated transition after its event commits. */
|
||||
function applyTransition(trace: SessionTrace, transition: SessionTraceTransition): void {
|
||||
Object.assign(trace, transition.scalars)
|
||||
switch (transition.pendingCalls.kind) {
|
||||
case 'none':
|
||||
break
|
||||
case 'add':
|
||||
trace.pendingCalls.add(transition.pendingCalls.callId)
|
||||
break
|
||||
case 'delete':
|
||||
trace.pendingCalls.delete(transition.pendingCalls.callId)
|
||||
break
|
||||
case 'clear':
|
||||
trace.pendingCalls.clear()
|
||||
break
|
||||
/* v8 ignore next -- validateEvent produces this closed transition union */
|
||||
default:
|
||||
assertNever(transition.pendingCalls, 'session trace pending-call transition')
|
||||
}
|
||||
/** Compile and validate one package-filter list. */
|
||||
function compilePatterns(field: 'package_allowlist' | 'package_blocklist', values: readonly string[]): RegExp[] {
|
||||
const seen = new Set<string>()
|
||||
return values.map((value) => {
|
||||
if (value.length === 0 || value.trim() !== value) {
|
||||
throw new Error(`invariants: ${field} entries must be non-blank and have no surrounding whitespace`)
|
||||
}
|
||||
if (seen.has(value)) {
|
||||
throw new Error(`invariants: ${field} contains duplicate regex ${JSON.stringify(value)}`)
|
||||
}
|
||||
seen.add(value)
|
||||
try {
|
||||
return new RegExp(value)
|
||||
} catch (cause) {
|
||||
throw new Error(`invariants: ${field} contains invalid regex ${JSON.stringify(value)}`, { cause })
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/** Validate and apply one event while rebuilding an already-committed log. */
|
||||
function replayEvent(trace: SessionTrace, event: SessionEvent): void {
|
||||
applyTransition(trace, validateEvent(trace, event))
|
||||
}
|
||||
|
||||
/** Allow an initial observation, idle/running transitions, and terminal disposal; reject repeats and leaving disposed. */
|
||||
function checkTransition(from: AgentStatus | undefined, to: AgentStatus): void {
|
||||
if (from === undefined) return
|
||||
if (from === to) {
|
||||
throw new InvariantError(`agent/status repeated ${to} (no-op transition)`)
|
||||
}
|
||||
if (from === 'disposed') {
|
||||
throw new InvariantError(`agent/status left terminal state disposed → ${to}`)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Register the runtime invariants. Contributions are effect-scoped, so
|
||||
* disposing the plugin fiber removes all listeners (HMR-safe). On (re-)apply
|
||||
* the trace state is rebuilt by replaying each existing session's log, so a
|
||||
* hot reload mid-turn does not falsely reject the next event.
|
||||
*
|
||||
* @param ctx - Cordis context that receives the invariant listeners.
|
||||
*/
|
||||
export function apply(ctx: Context): void {
|
||||
const traces = new WeakMap<Session, SessionTrace>()
|
||||
const stagedTransitions = new WeakMap<SessionEvent, {
|
||||
session: Session
|
||||
trace: SessionTrace
|
||||
transition: SessionTraceTransition
|
||||
}>()
|
||||
// Agent status has no stored history to replay; the first observation after
|
||||
// (re-)apply seeds the baseline, so a reload never produces a false positive.
|
||||
const lastStatus = new WeakMap<Agent, AgentStatus>()
|
||||
|
||||
const freshTrace = (): SessionTrace => ({
|
||||
lastSeq: -1,
|
||||
openTurn: null,
|
||||
openStep: null,
|
||||
nextTurn: 1,
|
||||
nextStep: 1,
|
||||
pendingCalls: new Set(),
|
||||
/** Package-owned invariant registry with global and regex-based selection. */
|
||||
export class InvariantService extends Service {
|
||||
static Config: Schema<Config> = z.object({
|
||||
enabled: z.boolean().default(true),
|
||||
package_allowlist: z.array(z.string()).default([]),
|
||||
package_blocklist: z.array(z.string()).default([]),
|
||||
})
|
||||
|
||||
/** Build (or rebuild) a session's trace by replaying its whole log. */
|
||||
const seedSession = (session: Session): SessionTrace => {
|
||||
const trace = freshTrace()
|
||||
traces.set(session, trace)
|
||||
for (const event of session.events) {
|
||||
replayEvent(trace, event)
|
||||
}
|
||||
return trace
|
||||
private readonly enabled: boolean
|
||||
private readonly ownerCtx: Context
|
||||
private readonly packageAllowlist: readonly RegExp[]
|
||||
private readonly packageBlocklist: readonly RegExp[]
|
||||
private readonly registrations = new Set<string>()
|
||||
|
||||
/**
|
||||
* Create and install the invariant registry.
|
||||
* @param ctx - Cordis context that owns the service.
|
||||
* @param config - global enablement and package-name regex filters.
|
||||
*/
|
||||
constructor(ctx: Context, config: Config = {}) {
|
||||
super(ctx, 'invariants')
|
||||
this.ownerCtx = ctx
|
||||
this.enabled = config.enabled ?? true
|
||||
this.packageAllowlist = compilePatterns('package_allowlist', config.package_allowlist ?? [])
|
||||
this.packageBlocklist = compilePatterns('package_blocklist', config.package_blocklist ?? [])
|
||||
}
|
||||
|
||||
// Every store-created session (the only kind that emits session/event) is
|
||||
// seeded first — via ctx.sessions.list() at apply or session/created — so the
|
||||
// fallback is a defensive guard, never hit in practice.
|
||||
/* v8 ignore next -- traceFor's fallback: session/event always follows a seed */
|
||||
const traceFor = (session: Session): SessionTrace => traces.get(session) ?? seedSession(session)
|
||||
/** Return whether one full package name passes the configured filters. */
|
||||
private selected(packageName: string): boolean {
|
||||
if (!this.enabled) return false
|
||||
if (this.packageAllowlist.length > 0
|
||||
&& !this.packageAllowlist.some(pattern => pattern.test(packageName))) return false
|
||||
return !this.packageBlocklist.some(pattern => pattern.test(packageName))
|
||||
}
|
||||
|
||||
// Rebuild state for sessions that already exist at (re-)apply time — HMR
|
||||
// reload starts a fresh fiber, and a mid-turn session would otherwise look
|
||||
// like it began with a stray chunk/step-end.
|
||||
for (const session of ctx.sessions.list()) seedSession(session)
|
||||
|
||||
// A newly created session may arrive seeded/forked (the constructor copies
|
||||
// the seed WITHOUT emitting session/event), so replay its log here too.
|
||||
ctx.on('session/created', (session) => { seedSession(session) }, { global: true })
|
||||
|
||||
ctx.on('session/event', (session, event) => {
|
||||
// Session resolves dispatch before committing, so internal/dispatch has
|
||||
// already staged this exact event. A later dispatch veto skips every
|
||||
// session/event callback and therefore leaves the live trace unchanged.
|
||||
const staged = stagedTransitions.get(event)
|
||||
/* v8 ignore next 2 -- internal/dispatch stages the exact callback arguments */
|
||||
if (staged === undefined || staged.session !== session) {
|
||||
throw new InvariantError('session/event reached publication without matching pre-commit validation')
|
||||
/**
|
||||
* Register one package's invariant installer. The package name is reserved
|
||||
* even when filtering disables its checks. Enabled installers run in a child
|
||||
* fiber; failure disposes that fiber and releases the reservation.
|
||||
* @param packageName - full npm package name that owns the contribution.
|
||||
* @param installer - synchronous listener installer for the child context.
|
||||
* @returns an effect-scoped disposer for the registration.
|
||||
*/
|
||||
register(packageName: string, installer: InvariantInstaller): () => void {
|
||||
if (packageName.length === 0 || packageName.trim() !== packageName || /\s/.test(packageName)) {
|
||||
throw new Error('invariants: packageName must be non-blank and contain no whitespace')
|
||||
}
|
||||
stagedTransitions.delete(event)
|
||||
applyTransition(staged.trace, staged.transition)
|
||||
}, { global: true })
|
||||
|
||||
ctx.on('agent/status', (agent, status) => {
|
||||
checkTransition(lastStatus.get(agent), status)
|
||||
lastStatus.set(agent, status)
|
||||
}, { global: true })
|
||||
|
||||
// --- Scoped-dispatch invariants (the agent-scoping seam) ---------------
|
||||
//
|
||||
// Every scope-filtered event family must dispatch with a scope carrier
|
||||
// (scopeTarget) whose key IS the subject the event's arguments name —
|
||||
// a dispatch without one silently reverts that event to global delivery
|
||||
// (agent-scoped listeners over-hear foreign agents), and a mis-keyed one
|
||||
// delivers to the wrong agent's listeners. `internal/dispatch` fires
|
||||
// synchronously before listener delivery, so a violation throws at the
|
||||
// dispatching call site. The generated table maps each family to the unique
|
||||
// payload path whose Program type matches the real scopeTarget routing key;
|
||||
// `null` means the key is external to the payload, so only carrier presence
|
||||
// can be asserted.
|
||||
ctx.on('internal/dispatch', (_mode, name, args, thisArg) => {
|
||||
const subjectOf = scopedSubjectResolverFor(name)
|
||||
if (subjectOf === undefined) return
|
||||
if (!isScopeCarrier(thisArg)) {
|
||||
throw new InvariantError(
|
||||
`"${name}" is a scope-filtered event but was dispatched without a scope carrier — `
|
||||
+ 'pass scopeTarget(base, subject) as the dispatch thisArg (agent events: use agentEvents(ctx, agent))')
|
||||
}
|
||||
if (subjectOf !== null && carrierKeyOf(thisArg) !== subjectOf(args)) {
|
||||
throw new InvariantError(
|
||||
`"${name}" was dispatched with a scope carrier keyed to a DIFFERENT subject than its arguments name — `
|
||||
+ 'the carrier key and the event\'s subject must be the same object (use agentEvents(ctx, agent))')
|
||||
}
|
||||
if (name === 'session/event') {
|
||||
const [session, event] = args as [Session, SessionEvent]
|
||||
const trace = traceFor(session)
|
||||
const transition = validateEvent(trace, event)
|
||||
// The exact event identity reaches the contained post-commit listener.
|
||||
// A later internal/dispatch listener may still veto; because validation
|
||||
// is pure, abandoning this weakly keyed transition does not advance the
|
||||
// committed trace or retain the session.
|
||||
stagedTransitions.set(event, { session, trace, transition })
|
||||
}
|
||||
}, { global: true })
|
||||
|
||||
// Request-reconstruction cross-check (the reconstructability RFC): a
|
||||
// loop-built request — frozen envelope + live sessionId is the marker; a
|
||||
// hand-built one-shot (compaction summarize) is unfrozen and skipped — must
|
||||
// be EXACTLY what the session log reconstructs:
|
||||
//
|
||||
// - messages: the folded header's session prefix (messagePrefix — the
|
||||
// `agent/session-prefix` product, logged on the header because no
|
||||
// session event carries it) followed by the
|
||||
// derivation over the log prefix strictly before the in-flight step's
|
||||
// `step/start` (the reconstruction boundary). The derivation is compared
|
||||
// against a FRESH Session built over that prefix — the same projection
|
||||
// code with zero shared state, so the live cache under test cannot vouch
|
||||
// for itself. Boundary-correct by construction: content appended after
|
||||
// the boundary (an `agent/request`-window inject) is legitimately absent
|
||||
// from this request, and a current-surface comparison would false-fire.
|
||||
// - header: every non-content field must equal the fold of the log's
|
||||
// `request/header` events — the loop logs the header event BEFORE
|
||||
// dispatch, so the fold already covers this request.
|
||||
//
|
||||
// Registered with `prepend: true` so a short-circuiting llm/stream listener
|
||||
// (the replay adapter returns its chunks without calling next()) cannot
|
||||
// silence the check by registering first. Prepend beats APPEND-registered
|
||||
// listeners only — two prepended listeners have no defined mutual order
|
||||
// (cordis unshift) — which is fine: correctness rests on the seq-bounded
|
||||
// fold below, never on listener timing.
|
||||
ctx.on('llm/stream', (options: GenerateOptions, next) => {
|
||||
if (options.sessionId === undefined || !Object.isFrozen(options)) return next()
|
||||
// GenerateOptions types sessionId as Branded<'SessionId'>, which IS
|
||||
// SessionId (dsh-llm cannot import it without a cycle) — no cast needed.
|
||||
const session = ctx.sessions.get(options.sessionId)
|
||||
if (!session) return next()
|
||||
if (!Object.isFrozen(options.messages)) {
|
||||
throw new InvariantError('a loop-built request must carry a frozen messages array')
|
||||
if (this.registrations.has(packageName)) {
|
||||
throw new Error(`invariants: package "${packageName}" is already registered`)
|
||||
}
|
||||
|
||||
const events = session.events
|
||||
// seq === index (checked above), so the last step/start's seq bounds the
|
||||
// prefix directly. The in-flight step's step/start is necessarily the
|
||||
// last one: the loop cannot open another step while this call streams.
|
||||
let boundary = -1
|
||||
for (let i = events.length - 1; i >= 0; i -= 1) {
|
||||
if (events[i]?.type === 'step/start') {
|
||||
boundary = i
|
||||
break
|
||||
}
|
||||
}
|
||||
if (boundary === -1) {
|
||||
throw new InvariantError('a loop-built request with no step/start in its session log')
|
||||
}
|
||||
const header = foldRequestHeader(events)
|
||||
if (header === undefined) {
|
||||
throw new InvariantError('a loop-built request with no request/header event in its session log')
|
||||
}
|
||||
const rebuilt = new Session(SessionId(`${String(session.id)}-invariant-rebuild`), structuredClone(events.slice(0, boundary)))
|
||||
// The reconstruction equation: the folded header's session prefix, then
|
||||
// the boundary derivation — the loop
|
||||
// logs the header event BEFORE dispatch, so the fold already covers this
|
||||
// request's prefix. JSON equality is sound here: both sides are
|
||||
// structuredClones produced by the same projection/build code path, so key
|
||||
// insertion order matches when the values do.
|
||||
const expected = [...header.messagePrefix ?? [], ...rebuilt.deriveMessages()]
|
||||
if (JSON.stringify(options.messages) !== JSON.stringify(expected)) {
|
||||
throw new InvariantError(`llm request for session "${String(session.id)}" diverges from the boundary derivation (log-reconstruction desync)`)
|
||||
}
|
||||
// Service method tracing binds `this.ctx` to the caller. This explicit
|
||||
// origin keeps registrations and their child fibers owned by the service;
|
||||
// companion disposal is covered independently by the returned disposer.
|
||||
const ctx = this.ownerCtx
|
||||
const registrations = this.registrations
|
||||
registrations.add(packageName)
|
||||
|
||||
const headerMatches = options.model === header.config.model
|
||||
&& options.system === header.system
|
||||
&& options.temperature === header.config.temperature
|
||||
&& options.maxTokens === header.config.maxTokens
|
||||
&& JSON.stringify(options.stop) === JSON.stringify(header.config.stop)
|
||||
&& JSON.stringify(options.tools ?? []) === JSON.stringify(header.tools ?? [])
|
||||
if (!headerMatches) {
|
||||
throw new InvariantError(`llm request for session "${String(session.id)}" diverges from the folded request header`)
|
||||
let registration: PendingInvariantRegistration
|
||||
try {
|
||||
registration = ctx.effect(async () => {
|
||||
if (!this.selected(packageName)) {
|
||||
return () => {
|
||||
registrations.delete(packageName)
|
||||
}
|
||||
}
|
||||
|
||||
const installInvariant = (childCtx: Context) => {
|
||||
installer(childCtx, (message): never => {
|
||||
throw new InvariantError(packageName, message)
|
||||
})
|
||||
}
|
||||
const child = ctx.plugin(installer.inject === undefined
|
||||
? installInvariant
|
||||
: Object.assign(installInvariant, { inject: installer.inject }))
|
||||
|
||||
try {
|
||||
await child
|
||||
} catch (error) {
|
||||
try {
|
||||
await child.dispose()
|
||||
} finally {
|
||||
registrations.delete(packageName)
|
||||
}
|
||||
throw error
|
||||
}
|
||||
|
||||
return async () => {
|
||||
try {
|
||||
await child.dispose()
|
||||
} finally {
|
||||
registrations.delete(packageName)
|
||||
}
|
||||
}
|
||||
}, `invariants.register(${JSON.stringify(packageName)})`)
|
||||
} catch (error) {
|
||||
registrations.delete(packageName)
|
||||
throw error
|
||||
}
|
||||
return next()
|
||||
}, { global: true, prepend: true })
|
||||
// Cordis attaches setup thenability and async teardown to this callable;
|
||||
// the service seam intentionally exposes only the conventional disposer.
|
||||
// eslint-disable-next-line @typescript-eslint/no-misused-promises -- the extra runtime shape stays private.
|
||||
return registration
|
||||
}
|
||||
}
|
||||
|
||||
export default InvariantService
|
||||
|
||||
@@ -1,71 +0,0 @@
|
||||
/**
|
||||
* Generated scoped-event routing-subject resolvers for dsh-invariants.
|
||||
* Do not edit by hand; run `pnpm run gen-scoped-events`.
|
||||
*
|
||||
* @module @deepseek-ai/dsh-invariants/scoped-events.generated
|
||||
*/
|
||||
|
||||
import type { Events } from 'cordis'
|
||||
import type { Scoped } from '@deepseek-ai/dsh-scope'
|
||||
import type {} from '@deepseek-ai/dsh-agent'
|
||||
import type {} from '@deepseek-ai/dsh-session'
|
||||
import type {} from '@deepseek-ai/dsh-subagent'
|
||||
import type {} from '@deepseek-ai/dsh-system-prompt'
|
||||
import type {} from '@deepseek-ai/dsh-tools'
|
||||
import type {} from '@deepseek-ai/dsh-user-approval'
|
||||
|
||||
type ScopedEventName = {
|
||||
[K in keyof Events]: ThisParameterType<Events[K]> extends Scoped<object> ? K : never
|
||||
}[keyof Events]
|
||||
|
||||
type ScopedSubjectResolver = (args: readonly unknown[]) => unknown
|
||||
|
||||
function adapt<K extends ScopedEventName>(
|
||||
resolver: (args: Parameters<Events[K]>) => unknown,
|
||||
): ScopedSubjectResolver {
|
||||
return args => resolver(args as Parameters<Events[K]>)
|
||||
}
|
||||
|
||||
const scopedSubjectResolvers = Object.freeze({
|
||||
'agent/created': adapt<'agent/created'>(args => args[0]),
|
||||
'agent/disposed': adapt<'agent/disposed'>(args => args[0]),
|
||||
'agent/error': adapt<'agent/error'>(args => args[0]),
|
||||
'agent/post-step': adapt<'agent/post-step'>(args => args[0]),
|
||||
'agent/pre-step': adapt<'agent/pre-step'>(args => args[0]),
|
||||
'agent/prompt-submit': adapt<'agent/prompt-submit'>(args => args[0]),
|
||||
'agent/queued': adapt<'agent/queued'>(args => args[0]),
|
||||
'agent/request': adapt<'agent/request'>(args => args[0]),
|
||||
'agent/request-error': adapt<'agent/request-error'>(args => args[0]),
|
||||
'agent/session-prefix': adapt<'agent/session-prefix'>(args => args[0]),
|
||||
'agent/session-start': adapt<'agent/session-start'>(args => args[0]),
|
||||
'agent/status': adapt<'agent/status'>(args => args[0]),
|
||||
'agent/step-result': adapt<'agent/step-result'>(args => args[0]),
|
||||
'agent/turn-continuation': adapt<'agent/turn-continuation'>(args => args[0]),
|
||||
'agent/turn-stop': adapt<'agent/turn-stop'>(args => args[0]),
|
||||
'approval/request': adapt<'approval/request'>(args => args[0].agent),
|
||||
'session/created': null,
|
||||
'session/disposed': null,
|
||||
'session/event': null,
|
||||
'session/flush': null,
|
||||
'subagent/end': null,
|
||||
'subagent/start': null,
|
||||
'system-prompt/assemble': adapt<'system-prompt/assemble'>(args => args[1].scope),
|
||||
'tools/execute': adapt<'tools/execute'>(args => args[0].agent),
|
||||
'tools/post-execute': adapt<'tools/post-execute'>(args => args[0].agent),
|
||||
'tools/pre-execute': adapt<'tools/pre-execute'>(args => args[0].agent),
|
||||
'tools/result': adapt<'tools/result'>(args => args[0].agent),
|
||||
} as const satisfies Readonly<Record<ScopedEventName, ScopedSubjectResolver | null>>)
|
||||
|
||||
const scopedSubjectResolverIndex: Readonly<Record<string, ScopedSubjectResolver | null>> = scopedSubjectResolvers
|
||||
|
||||
/**
|
||||
* Resolve the routing key named by one scoped event payload. A null
|
||||
* resolver means the payload cannot expose its external routing key, so the
|
||||
* invariant checks carrier presence only.
|
||||
* @param event - runtime Cordis event name.
|
||||
* @returns the generated subject resolver, null for presence-only,
|
||||
* or undefined when the event is not scope-filtered.
|
||||
*/
|
||||
export function scopedSubjectResolverFor(event: string): ScopedSubjectResolver | null | undefined {
|
||||
return scopedSubjectResolverIndex[event]
|
||||
}
|
||||
Reference in New Issue
Block a user