subagent-durability-failure.ts 4.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110
  1. import type { Context } from 'cordis'
  2. import { SessionId } from '@deepseek-ai/dsh-session'
  3. export const name = 'subagent-durability-failure'
  4. export const inject = ['agents', 'sessionPersistence', 'subagents']
  5. /**
  6. * The authored parent transcript names the background child by a stable
  7. * placeholder id, but the live continuable child is minted with a fresh random
  8. * session id at run time. This snapshot-only overlay bridges that gap and forces
  9. * deterministic ordering plus authored child durability failures:
  10. *
  11. * - `PLACEHOLDER_CHILD_ID` in a scripted `send_message` is remapped to the real
  12. * child so both follow-ups queue onto the same live inbox in FIFO order.
  13. * - The unknown-id `send_message` (`UNKNOWN_CHILD_ID`) resolves through a
  14. * persistence load fenced behind both accepted follow-ups, so the transcript
  15. * records the same order on every runner.
  16. * - The child's final continuation turn fails its durability checkpoint with a
  17. * fixed message, so the scenario proves child-first disposal survives a failed
  18. * last flush.
  19. * - Under `DSH_SUBAGENT_PUBLISHED_FAILURE`, a one-shot child's first
  20. * follow-up fails after publication, so its model prompt never runs; its
  21. * published handle then also fails disposal, so the parent observes both
  22. * independent failures.
  23. */
  24. const PLACEHOLDER_CHILD_ID = '33333333-3333-4333-8333-333333333333'
  25. const UNKNOWN_CHILD_ID = '22222222-2222-4222-8222-222222222222'
  26. /** The child continuation turn whose durability checkpoint is forced to fail. */
  27. const FAILED_CHECKPOINT_TURN = 3
  28. /** Fail the child checkpoint and stabilize the authored follow-up failure ordering. */
  29. export function apply(ctx: Context): void {
  30. const followupsAccepted = Promise.withResolvers<undefined>()
  31. const publishedFailure = process.env.DSH_SUBAGENT_PUBLISHED_FAILURE === '1'
  32. const persistence = ctx.sessionPersistence
  33. const load = persistence.load.bind(persistence)
  34. const agents = ctx.agents
  35. const create = agents.create.bind(agents)
  36. agents.create = async (options) => {
  37. const handle = await create(options)
  38. if (!publishedFailure || options.meta?.parentSession === undefined) return handle
  39. handle.agent.followup = () => {
  40. throw new Error('snapshot published run failed')
  41. }
  42. return {
  43. ...handle,
  44. async dispose() {
  45. await handle.dispose()
  46. throw new Error('snapshot published handle disposal failed')
  47. },
  48. }
  49. }
  50. // The unavailable-child lookup is real asynchronous I/O. Fence it behind both
  51. // authored follow-ups so runner speed cannot reorder the exact log.
  52. persistence.load = async (id) => {
  53. if (id === UNKNOWN_CHILD_ID) await followupsAccepted.promise
  54. return load.call(persistence, id)
  55. }
  56. ctx.effect(() => () => {
  57. agents.create = create
  58. persistence.load = load
  59. followupsAccepted.resolve(undefined)
  60. }, 'subagent snapshot ordering')
  61. // Remap the placeholder child id in a follow-up to the live child. The child
  62. // id the model "knows" is authored into the transcript, while the running
  63. // child is minted with a random id, so without this the follow-ups would
  64. // never reach the live inbox.
  65. let realChildId: string | undefined
  66. const subagents = ctx.subagents as unknown as {
  67. followup: (authority: unknown, childId: SessionId, content: unknown, options: unknown) => Promise<unknown>
  68. }
  69. const deliver = subagents.followup.bind(subagents)
  70. subagents.followup = (authority, childId, content, options) => {
  71. const mapped = childId === PLACEHOLDER_CHILD_ID && realChildId !== undefined
  72. ? SessionId(realChildId)
  73. : childId
  74. return deliver(authority, mapped, content, options)
  75. }
  76. // Both authored follow-ups reach the child inbox before the unknown-id lookup
  77. // runs, so the queued FIFO order is what the transcript records. The first
  78. // child enqueue is the initial delegation, which also pins the real child id.
  79. let accepted = 0
  80. ctx.on('agent/inbox/inserted', (agent) => {
  81. if (agent.session.header.parentSession === undefined) return
  82. if (realChildId === undefined) realChildId = agent.session.header.id
  83. accepted += 1
  84. if (accepted >= 3) followupsAccepted.resolve(undefined)
  85. })
  86. ctx.on('agent/pre-step', async (agent, _messages, _context, next) => {
  87. if (agent.session.header.parentSession !== undefined) await followupsAccepted.promise
  88. return next()
  89. })
  90. // The child's ordinary per-turn flushes succeed; only the final continuation
  91. // turn's durability checkpoint fails, turning that turn/end into a durable
  92. // error the parent never sees.
  93. const childTurn = new WeakMap<object, number>()
  94. ctx.on('session/event', (session, event) => {
  95. if (session.header.parentSession === undefined || event.type !== 'turn/start') return
  96. childTurn.set(session, event.data.turn)
  97. })
  98. ctx.on('session/flush', (session) => {
  99. if (session.header.parentSession === undefined) return
  100. if (childTurn.get(session) === FAILED_CHECKPOINT_TURN) throw new Error('snapshot disk full')
  101. })
  102. }