index.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409
  1. /**
  2. * Runtime listeners that fail loudly when cross-event contracts are broken:
  3. * turn and step nesting, scoped dispatch, status transitions, and request
  4. * reconstruction. The plugin has no environment guard and is active wherever
  5. * mounted, including the default `dsh-agent-spine-demo` bundle; custom compositions
  6. * may omit it. Sessions own immutable, surface-valid event storage; this plugin
  7. * checks only relationships that event acceptance cannot express.
  8. * @module @deepseek-ai/dsh-invariants
  9. */
  10. import type { Context } from 'cordis'
  11. import { carrierKeyOf, isScopeCarrier } from '@deepseek-ai/dsh-scope'
  12. import { assertNever, HarnessError } from '@deepseek-ai/dsh-llm'
  13. import type { CallId, GenerateOptions } from '@deepseek-ai/dsh-llm'
  14. import type { Agent, AgentStatus } from '@deepseek-ai/dsh-agent'
  15. import { Session, SessionId, foldRequestHeader } from '@deepseek-ai/dsh-session'
  16. import type { SessionEvent } from '@deepseek-ai/dsh-session'
  17. import { scopedSubjectResolverFor } from './scoped-events.generated.ts'
  18. export const name = 'invariants'
  19. export const inject = ['sessions']
  20. /**
  21. * Thrown when a harness event-contract invariant is violated. Extends
  22. * {@link HarnessError} (`code: 'INVARIANT'`) so a violation is routable like
  23. * any other harness failure.
  24. */
  25. export class InvariantError extends HarnessError {
  26. constructor(message: string) {
  27. super(`invariant violated: ${message}`, 'INVARIANT')
  28. this.name = 'InvariantError'
  29. }
  30. }
  31. /** Per-session bookkeeping for the session-log invariants. */
  32. interface SessionTrace {
  33. /** Highest `seq` seen so far (must strictly increase). */
  34. lastSeq: number
  35. /** Open turn number, or null between turns. */
  36. openTurn: number | null
  37. /** Open step within the current turn, or null between steps. */
  38. openStep: number | null
  39. /** The next turn number expected in this session log. */
  40. nextTurn: number
  41. /** The next step number expected within the open turn. */
  42. nextStep: number
  43. /**
  44. * Tool-call ids issued in the OPEN step awaiting a result. Cleared at
  45. * `step/end` — a result must arrive in the same step as its call.
  46. */
  47. pendingCalls: Set<CallId>
  48. }
  49. /** One accepted event's deferred mutation of a live session trace. */
  50. interface SessionTraceTransition {
  51. /** Scalar state after the event commits. */
  52. scalars: Pick<SessionTrace, 'lastSeq' | 'openTurn' | 'openStep' | 'nextTurn' | 'nextStep'>
  53. /** The event's mutation of the open step's pending call set. */
  54. pendingCalls:
  55. | { kind: 'none' }
  56. | { kind: 'add' | 'delete'; callId: CallId }
  57. | { kind: 'clear' }
  58. }
  59. /** Assert that a step-scoped event names the currently open turn and step. */
  60. function requireOpenStep(trace: SessionTrace, kind: string, turn: number, step: number): void {
  61. if (trace.openTurn !== turn || trace.openStep !== step) {
  62. throw new InvariantError(
  63. `${kind} names turn ${turn}/step ${step} but open is turn ${trace.openTurn}/step ${trace.openStep}`,
  64. )
  65. }
  66. }
  67. /** Validate one candidate event without mutating the committed session trace. */
  68. function validateEvent(trace: SessionTrace, event: SessionEvent): SessionTraceTransition {
  69. // seq is strictly monotonic — the spine of replay equivalence. lastSeq
  70. // starts at -1, so the first event (seq 0) passes.
  71. if (event.seq <= trace.lastSeq) {
  72. throw new InvariantError(`seq must strictly increase: saw ${event.seq} after ${trace.lastSeq}`)
  73. }
  74. let openTurn = trace.openTurn
  75. let openStep = trace.openStep
  76. let nextTurn = trace.nextTurn
  77. let nextStep = trace.nextStep
  78. let pendingCalls: SessionTraceTransition['pendingCalls'] = { kind: 'none' }
  79. // Boundary/step-scoped events have explicit cases; every OTHER event type —
  80. // including plugin-added (merge-extensible) SessionEventMap keys — is caught
  81. // by the `default` and must be turn-enclosed (the turn-enclosure RFC). No assertNever: an
  82. // unknown variant is valid, not a compile error.
  83. switch (event.type) {
  84. case 'turn/start': {
  85. if (trace.openTurn !== null) {
  86. throw new InvariantError(`turn/start ${event.data.turn} while turn ${trace.openTurn} is still open`)
  87. }
  88. // Current sessions replay full logs, so numbering starts at 1 and remains
  89. // contiguous. If a future compaction/fork stores a partial log, it must
  90. // seed `nextTurn` from retained metadata before this check runs.
  91. if (event.data.turn !== trace.nextTurn) {
  92. throw new InvariantError(`turn/start expected turn ${trace.nextTurn}, got ${event.data.turn}`)
  93. }
  94. openTurn = event.data.turn
  95. nextStep = 1
  96. break
  97. }
  98. case 'turn/end': {
  99. if (trace.openTurn !== event.data.turn) {
  100. throw new InvariantError(`turn/end ${event.data.turn} does not match open turn ${trace.openTurn}`)
  101. }
  102. if (trace.openStep !== null) {
  103. throw new InvariantError(`turn/end ${event.data.turn} while step ${trace.openStep} is still open`)
  104. }
  105. openTurn = null
  106. nextTurn += 1
  107. break
  108. }
  109. case 'step/start': {
  110. if (trace.openTurn !== event.data.turn) {
  111. throw new InvariantError(`step/start in turn ${event.data.turn} but open turn is ${trace.openTurn}`)
  112. }
  113. if (trace.openStep !== null) {
  114. throw new InvariantError(`step/start ${event.data.step} while step ${trace.openStep} is still open`)
  115. }
  116. // Steps are checked under the same full-log assumption as turns above.
  117. if (event.data.step !== trace.nextStep) {
  118. throw new InvariantError(`step/start expected step ${trace.nextStep} in turn ${event.data.turn}, got ${event.data.step}`)
  119. }
  120. openStep = event.data.step
  121. break
  122. }
  123. case 'step/end': {
  124. requireOpenStep(trace, 'step/end', event.data.turn, event.data.step)
  125. // A result must arrive in the step that issued the call; orphan calls
  126. // (a step that errored before its result) do not carry to the next step.
  127. pendingCalls = { kind: 'clear' }
  128. openStep = null
  129. nextStep += 1
  130. break
  131. }
  132. case 'assistant/chunk': {
  133. requireOpenStep(trace, 'assistant/chunk', event.data.turn, event.data.step)
  134. break
  135. }
  136. case 'assistant/message': {
  137. requireOpenStep(trace, 'assistant/message', event.data.turn, event.data.step)
  138. break
  139. }
  140. case 'tool/call': {
  141. requireOpenStep(trace, 'tool/call', event.data.turn, event.data.step)
  142. pendingCalls = { kind: 'add', callId: event.data.callId }
  143. break
  144. }
  145. case 'tool/result': {
  146. requireOpenStep(trace, 'tool/result', event.data.turn, event.data.step)
  147. // A result needs a prior matching call in the same step. (The converse
  148. // does NOT hold: a call may have no result — a throwing tool-execution
  149. // pipeline step ends the turn with no tool/result, which is legal.)
  150. const syntheticInterrupted = event.data.isError && event.data.error?.code === 'interrupted'
  151. if (!trace.pendingCalls.has(event.data.callId) && !syntheticInterrupted) {
  152. throw new InvariantError(`tool/result for ${event.data.callId} with no prior tool/call in this step`)
  153. }
  154. pendingCalls = { kind: 'delete', callId: event.data.callId }
  155. break
  156. }
  157. // Turn-enclosure (the turn-enclosure RFC): EVERY session event not handled by a boundary
  158. // case above must sit inside an open turn. The durable session log uses the
  159. // turn as its commit/replay boundary (the JSONL backend treats anything
  160. // after the last turn/end as a crash tail), so a bare event between turns is
  161. // silently dropped on reload. The loop records queued user messages after
  162. // turn/start, and an idle agent.inject() wraps its context/message in a
  163. // one-shot turn. A `default`
  164. // (not an enumerated list) is deliberate: SessionEventMap is
  165. // merge-extensible, so a PLUGIN-added event type appended while idle must
  166. // also fail here rather than fall through and be dropped on resume.
  167. default: {
  168. if (trace.openTurn === null) {
  169. throw new InvariantError(`${event.type} appended outside any open turn (every event must be turn-enclosed)`)
  170. }
  171. break
  172. }
  173. }
  174. return {
  175. scalars: { lastSeq: event.seq, openTurn, openStep, nextTurn, nextStep },
  176. pendingCalls,
  177. }
  178. }
  179. /** Apply one already-validated transition after its event commits. */
  180. function applyTransition(trace: SessionTrace, transition: SessionTraceTransition): void {
  181. Object.assign(trace, transition.scalars)
  182. switch (transition.pendingCalls.kind) {
  183. case 'none':
  184. break
  185. case 'add':
  186. trace.pendingCalls.add(transition.pendingCalls.callId)
  187. break
  188. case 'delete':
  189. trace.pendingCalls.delete(transition.pendingCalls.callId)
  190. break
  191. case 'clear':
  192. trace.pendingCalls.clear()
  193. break
  194. /* v8 ignore next -- validateEvent produces this closed transition union */
  195. default:
  196. assertNever(transition.pendingCalls, 'session trace pending-call transition')
  197. }
  198. }
  199. /** Validate and apply one event while rebuilding an already-committed log. */
  200. function replayEvent(trace: SessionTrace, event: SessionEvent): void {
  201. applyTransition(trace, validateEvent(trace, event))
  202. }
  203. /** Allow an initial observation, idle/running transitions, and terminal disposal; reject repeats and leaving disposed. */
  204. function checkTransition(from: AgentStatus | undefined, to: AgentStatus): void {
  205. if (from === undefined) return
  206. if (from === to) {
  207. throw new InvariantError(`agent/status repeated ${to} (no-op transition)`)
  208. }
  209. if (from === 'disposed') {
  210. throw new InvariantError(`agent/status left terminal state disposed → ${to}`)
  211. }
  212. }
  213. /**
  214. * Register the runtime invariants. Contributions are effect-scoped, so
  215. * disposing the plugin fiber removes all listeners (HMR-safe). On (re-)apply
  216. * the trace state is rebuilt by replaying each existing session's log, so a
  217. * hot reload mid-turn does not falsely reject the next event.
  218. *
  219. * @param ctx - Cordis context that receives the invariant listeners.
  220. */
  221. export function apply(ctx: Context): void {
  222. const traces = new WeakMap<Session, SessionTrace>()
  223. const stagedTransitions = new WeakMap<SessionEvent, {
  224. session: Session
  225. trace: SessionTrace
  226. transition: SessionTraceTransition
  227. }>()
  228. // Agent status has no stored history to replay; the first observation after
  229. // (re-)apply seeds the baseline, so a reload never produces a false positive.
  230. const lastStatus = new WeakMap<Agent, AgentStatus>()
  231. const freshTrace = (): SessionTrace => ({
  232. lastSeq: -1,
  233. openTurn: null,
  234. openStep: null,
  235. nextTurn: 1,
  236. nextStep: 1,
  237. pendingCalls: new Set(),
  238. })
  239. /** Build (or rebuild) a session's trace by replaying its whole log. */
  240. const seedSession = (session: Session): SessionTrace => {
  241. const trace = freshTrace()
  242. traces.set(session, trace)
  243. for (const event of session.events) {
  244. replayEvent(trace, event)
  245. }
  246. return trace
  247. }
  248. // Every store-created session (the only kind that emits session/event) is
  249. // seeded first — via ctx.sessions.list() at apply or session/created — so the
  250. // fallback is a defensive guard, never hit in practice.
  251. /* v8 ignore next -- traceFor's fallback: session/event always follows a seed */
  252. const traceFor = (session: Session): SessionTrace => traces.get(session) ?? seedSession(session)
  253. // Rebuild state for sessions that already exist at (re-)apply time — HMR
  254. // reload starts a fresh fiber, and a mid-turn session would otherwise look
  255. // like it began with a stray chunk/step-end.
  256. for (const session of ctx.sessions.list()) seedSession(session)
  257. // A newly created session may arrive seeded/forked (the constructor copies
  258. // the seed WITHOUT emitting session/event), so replay its log here too.
  259. ctx.on('session/created', (session) => { seedSession(session) }, { global: true })
  260. ctx.on('session/event', (session, event) => {
  261. // Session resolves dispatch before committing, so internal/dispatch has
  262. // already staged this exact event. A later dispatch veto skips every
  263. // session/event callback and therefore leaves the live trace unchanged.
  264. const staged = stagedTransitions.get(event)
  265. /* v8 ignore next 2 -- internal/dispatch stages the exact callback arguments */
  266. if (staged === undefined || staged.session !== session) {
  267. throw new InvariantError('session/event reached publication without matching pre-commit validation')
  268. }
  269. stagedTransitions.delete(event)
  270. applyTransition(staged.trace, staged.transition)
  271. }, { global: true })
  272. ctx.on('agent/status', (agent, status) => {
  273. checkTransition(lastStatus.get(agent), status)
  274. lastStatus.set(agent, status)
  275. }, { global: true })
  276. // --- Scoped-dispatch invariants (the agent-scoping seam) ---------------
  277. //
  278. // Every scope-filtered event family must dispatch with a scope carrier
  279. // (scopeTarget) whose key IS the subject the event's arguments name —
  280. // a dispatch without one silently reverts that event to global delivery
  281. // (agent-scoped listeners over-hear foreign agents), and a mis-keyed one
  282. // delivers to the wrong agent's listeners. `internal/dispatch` fires
  283. // synchronously before listener delivery, so a violation throws at the
  284. // dispatching call site. The generated table maps each family to the unique
  285. // payload path whose Program type matches the real scopeTarget routing key;
  286. // `null` means the key is external to the payload, so only carrier presence
  287. // can be asserted.
  288. ctx.on('internal/dispatch', (_mode, name, args, thisArg) => {
  289. const subjectOf = scopedSubjectResolverFor(name)
  290. if (subjectOf === undefined) return
  291. if (!isScopeCarrier(thisArg)) {
  292. throw new InvariantError(
  293. `"${name}" is a scope-filtered event but was dispatched without a scope carrier — `
  294. + 'pass scopeTarget(base, subject) as the dispatch thisArg (agent events: use agentEvents(ctx, agent))')
  295. }
  296. if (subjectOf !== null && carrierKeyOf(thisArg) !== subjectOf(args)) {
  297. throw new InvariantError(
  298. `"${name}" was dispatched with a scope carrier keyed to a DIFFERENT subject than its arguments name — `
  299. + 'the carrier key and the event\'s subject must be the same object (use agentEvents(ctx, agent))')
  300. }
  301. if (name === 'session/event') {
  302. const [session, event] = args as [Session, SessionEvent]
  303. const trace = traceFor(session)
  304. const transition = validateEvent(trace, event)
  305. // The exact event identity reaches the contained post-commit listener.
  306. // A later internal/dispatch listener may still veto; because validation
  307. // is pure, abandoning this weakly keyed transition does not advance the
  308. // committed trace or retain the session.
  309. stagedTransitions.set(event, { session, trace, transition })
  310. }
  311. }, { global: true })
  312. // Request-reconstruction cross-check (the reconstructability RFC): a
  313. // loop-built request — frozen envelope + live sessionId is the marker; a
  314. // hand-built one-shot (compaction summarize) is unfrozen and skipped — must
  315. // be EXACTLY what the session log reconstructs:
  316. //
  317. // - messages: the folded header's session prefix (messagePrefix — the
  318. // `agent/session-prefix` product, logged on the header because no
  319. // session event carries it) followed by the
  320. // derivation over the log prefix strictly before the in-flight step's
  321. // `step/start` (the reconstruction boundary). The derivation is compared
  322. // against a FRESH Session built over that prefix — the same projection
  323. // code with zero shared state, so the live cache under test cannot vouch
  324. // for itself. Boundary-correct by construction: content appended after
  325. // the boundary (an `agent/request`-window inject) is legitimately absent
  326. // from this request, and a current-surface comparison would false-fire.
  327. // - header: every non-content field must equal the fold of the log's
  328. // `request/header` events — the loop logs the header event BEFORE
  329. // dispatch, so the fold already covers this request.
  330. //
  331. // Registered with `prepend: true` so a short-circuiting llm/stream listener
  332. // (the replay adapter returns its chunks without calling next()) cannot
  333. // silence the check by registering first. Prepend beats APPEND-registered
  334. // listeners only — two prepended listeners have no defined mutual order
  335. // (cordis unshift) — which is fine: correctness rests on the seq-bounded
  336. // fold below, never on listener timing.
  337. ctx.on('llm/stream', (options: GenerateOptions, next) => {
  338. if (options.sessionId === undefined || !Object.isFrozen(options)) return next()
  339. // GenerateOptions types sessionId as Branded<'SessionId'>, which IS
  340. // SessionId (dsh-llm cannot import it without a cycle) — no cast needed.
  341. const session = ctx.sessions.get(options.sessionId)
  342. if (!session) return next()
  343. if (!Object.isFrozen(options.messages)) {
  344. throw new InvariantError('a loop-built request must carry a frozen messages array')
  345. }
  346. const events = session.events
  347. // seq === index (checked above), so the last step/start's seq bounds the
  348. // prefix directly. The in-flight step's step/start is necessarily the
  349. // last one: the loop cannot open another step while this call streams.
  350. let boundary = -1
  351. for (let i = events.length - 1; i >= 0; i -= 1) {
  352. if (events[i]?.type === 'step/start') {
  353. boundary = i
  354. break
  355. }
  356. }
  357. if (boundary === -1) {
  358. throw new InvariantError('a loop-built request with no step/start in its session log')
  359. }
  360. const header = foldRequestHeader(events)
  361. if (header === undefined) {
  362. throw new InvariantError('a loop-built request with no request/header event in its session log')
  363. }
  364. const rebuilt = new Session(SessionId(`${String(session.id)}-invariant-rebuild`), structuredClone(events.slice(0, boundary)))
  365. // The reconstruction equation: the folded header's session prefix, then
  366. // the boundary derivation — the loop
  367. // logs the header event BEFORE dispatch, so the fold already covers this
  368. // request's prefix. JSON equality is sound here: both sides are
  369. // structuredClones produced by the same projection/build code path, so key
  370. // insertion order matches when the values do.
  371. const expected = [...header.messagePrefix ?? [], ...rebuilt.deriveMessages()]
  372. if (JSON.stringify(options.messages) !== JSON.stringify(expected)) {
  373. throw new InvariantError(`llm request for session "${String(session.id)}" diverges from the boundary derivation (log-reconstruction desync)`)
  374. }
  375. const headerMatches = options.model === header.config.model
  376. && options.system === header.system
  377. && options.temperature === header.config.temperature
  378. && options.maxTokens === header.config.maxTokens
  379. && JSON.stringify(options.stop) === JSON.stringify(header.config.stop)
  380. && JSON.stringify(options.tools ?? []) === JSON.stringify(header.tools ?? [])
  381. if (!headerMatches) {
  382. throw new InvariantError(`llm request for session "${String(session.id)}" diverges from the folded request header`)
  383. }
  384. return next()
  385. }, { global: true, prepend: true })
  386. }