invariant.ts 7.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167
  1. /** Package-owned durable workflow-record invariants. @module @deepseek-ai/dsh-tool-workflow/invariant */
  2. import type { Context } from '@deepseek-ai/cordis'
  3. import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
  4. import type { InvariantFailure, InvariantInstaller } from '@deepseek-ai/dsh-invariants'
  5. import type {} from './types.ts'
  6. const PACKAGE_NAME = '@deepseek-ai/dsh-tool-workflow'
  7. /** Cordis companion plugin name. */
  8. export const name = 'tool-workflow-invariant'
  9. /** Services required to validate existing and newly appended Session logs. */
  10. export const inject = ['invariants']
  11. interface RunTrace {
  12. ended: boolean
  13. readonly members: Map<number, boolean>
  14. }
  15. type WorkflowTrace = Map<string, RunTrace>
  16. /** Whether this package owns the candidate Session event. */
  17. function isWorkflowRecordEvent(event: SessionEvent): boolean {
  18. return event.type.startsWith('tool-workflow/')
  19. }
  20. /** Require a durable opaque identity to be a non-empty string. */
  21. function stringId(value: unknown, label: string, fail: InvariantFailure): string {
  22. if (typeof value !== 'string' || value.length === 0) fail(`${label} must be a non-empty string`)
  23. return value
  24. }
  25. /** Require one workflow member's 1-based sequence identity. */
  26. function memberSeq(value: unknown, fail: InvariantFailure): number {
  27. if (!Number.isSafeInteger(value) || (value as number) < 1) {
  28. fail('tool-workflow member seq must be a positive safe integer')
  29. }
  30. return value as number
  31. }
  32. /** Read one plain payload field without trusting restored plugin data. */
  33. function recordOf(event: SessionEvent, fail: InvariantFailure): Record<string, unknown> {
  34. const data: unknown = event.data
  35. if (data === null || typeof data !== 'object' || Array.isArray(data)) {
  36. fail(`${event.type} data must be a JSON object`)
  37. }
  38. return data as Record<string, unknown>
  39. }
  40. /** Copy only the run one candidate can mutate; other committed states stay shared. */
  41. function cloneTraceForEvent(
  42. source: WorkflowTrace,
  43. event: SessionEvent,
  44. fail: InvariantFailure,
  45. ): WorkflowTrace {
  46. const trace = new Map(source)
  47. if (event.type === 'tool-workflow/run-start') return trace
  48. const data = recordOf(event, fail)
  49. const runId = stringId(data.runId, `${event.type} runId`, fail)
  50. const run = source.get(runId)
  51. if (run !== undefined) {
  52. trace.set(runId, { ended: run.ended, members: new Map(run.members) })
  53. }
  54. return trace
  55. }
  56. /** Require the named run to exist and remain open. */
  57. function openRun(trace: WorkflowTrace, runId: string, eventType: string, fail: InvariantFailure): RunTrace {
  58. const run = trace.get(runId)
  59. if (run === undefined) fail(`${eventType} has no matching tool-workflow/run-start for run ${runId}`)
  60. if (run.ended) fail(`${eventType} appears after tool-workflow/run-end for run ${runId}`)
  61. return run
  62. }
  63. /** Advance the workflow-record fold with one relevant Session event. */
  64. function applyEvent(trace: WorkflowTrace, event: SessionEvent, fail: InvariantFailure): void {
  65. const data = recordOf(event, fail)
  66. const runId = stringId(data.runId, `${event.type} runId`, fail)
  67. switch (event.type) {
  68. case 'tool-workflow/run-start': {
  69. if (typeof data.name !== 'string' || data.name.length === 0) {
  70. fail('tool-workflow/run-start name must be a non-empty string')
  71. }
  72. if (trace.has(runId)) fail(`tool-workflow/run-start repeats run ${runId}`)
  73. trace.set(runId, { ended: false, members: new Map() })
  74. return
  75. }
  76. case 'tool-workflow/agent-start': {
  77. const run = openRun(trace, runId, event.type, fail)
  78. const seq = memberSeq(data.seq, fail)
  79. if (typeof data.label !== 'string') fail('tool-workflow/agent-start label must be a string')
  80. if (data.phase !== undefined && typeof data.phase !== 'string') {
  81. fail('tool-workflow/agent-start phase must be a string when present')
  82. }
  83. stringId(data.childId, 'tool-workflow/agent-start childId', fail)
  84. if (run.members.has(seq)) fail(`tool-workflow/agent-start repeats member seq ${seq} in run ${runId}`)
  85. run.members.set(seq, false)
  86. return
  87. }
  88. case 'tool-workflow/agent-end': {
  89. const run = openRun(trace, runId, event.type, fail)
  90. const seq = memberSeq(data.seq, fail)
  91. if (data.outcome !== 'completed' && data.outcome !== 'failed' && data.outcome !== 'cancelled') {
  92. fail(`tool-workflow/agent-end outcome ${String(data.outcome)} is invalid`)
  93. }
  94. const ended = run.members.get(seq)
  95. if (ended === undefined) fail(`tool-workflow/agent-end has no matching member seq ${seq} in run ${runId}`)
  96. if (ended) fail(`tool-workflow/agent-end repeats member seq ${seq} in run ${runId}`)
  97. run.members.set(seq, true)
  98. return
  99. }
  100. case 'tool-workflow/run-end': {
  101. const run = openRun(trace, runId, event.type, fail)
  102. if (data.stopReason !== 'completed' && data.stopReason !== 'cancelled' && data.stopReason !== 'error') {
  103. fail(`tool-workflow/run-end stopReason ${String(data.stopReason)} is invalid`)
  104. }
  105. const openMembers = [...run.members].filter(([, ended]) => !ended).map(([seq]) => seq)
  106. if (openMembers.length > 0) {
  107. fail(`tool-workflow/run-end leaves member seq ${openMembers.join(', ')} open in run ${runId}`)
  108. }
  109. run.ended = true
  110. run.members.clear()
  111. return
  112. }
  113. default:
  114. fail(`unknown tool-workflow event type ${event.type}`)
  115. }
  116. }
  117. /** Install an independent incremental fold over every attached Session. */
  118. const install: InvariantInstaller = Object.assign((ctx: Context, fail: InvariantFailure) => {
  119. const traces = new WeakMap<Session, WorkflowTrace>()
  120. const staged = new WeakMap<SessionEvent, { session: Session; trace: WorkflowTrace }>()
  121. const seed = (session: Session): WorkflowTrace => {
  122. const trace: WorkflowTrace = new Map()
  123. for (const event of session.events.filter(isWorkflowRecordEvent)) applyEvent(trace, event, fail)
  124. traces.set(session, trace)
  125. return trace
  126. }
  127. ctx.sessions.list().forEach(seed)
  128. ctx.on('session/created', (session) => { seed(session) }, { global: true })
  129. ctx.on('internal/dispatch', (_mode, eventName, args) => {
  130. if (eventName !== 'session/event') return
  131. const [session, event] = args as [Session, SessionEvent]
  132. if (!isWorkflowRecordEvent(event)) return
  133. // session/event dispatch follows list() or session/created seeding.
  134. const trace = cloneTraceForEvent(traces.get(session) as WorkflowTrace, event, fail)
  135. applyEvent(trace, event, fail)
  136. staged.set(event, { session, trace })
  137. }, { global: true })
  138. ctx.on('session/event', (session, event) => {
  139. if (!isWorkflowRecordEvent(event)) return
  140. const candidate = staged.get(event)
  141. /* v8 ignore next 2 -- internal/dispatch stages the exact session/event callback arguments. */
  142. if (candidate === undefined || candidate.session !== session) {
  143. return fail('session/event reached publication without matching workflow-record validation')
  144. }
  145. staged.delete(event)
  146. traces.set(session, candidate.trace)
  147. }, { global: true })
  148. }, { inject: ['sessions'] })
  149. /** Register this package's invariant companion. */
  150. export const apply = (ctx: Context): Promise<() => void> =>
  151. Promise.resolve(ctx.invariants.register(PACKAGE_NAME, install))