1
0

python-snapshot-workflow-order.mjs 1.9 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647
  1. /** Hold the advanced workflow child's first step until its parent records membership. */
  2. export const name = 'python-snapshot-workflow-order'
  3. /**
  4. * @param {import('@deepseek-ai/cordis').Context} ctx - Scenario-local host context.
  5. * @param {{ parentSessionId: string, prompt: string }} config - Exact advanced scenario identities.
  6. */
  7. export function apply(ctx, config) {
  8. const started = new Set()
  9. const pending = new Map()
  10. let disposed = false
  11. ctx.effect(() => async () => {
  12. disposed = true
  13. const waits = [...pending.values()]
  14. for (const wait of waits) wait.reject(new Error('workflow snapshot barrier disposed'))
  15. await Promise.allSettled(waits.map(wait => wait.done))
  16. started.clear()
  17. })
  18. ctx.on('session/event', (session, event) => {
  19. if (disposed || session.id !== config.parentSessionId || event.type !== 'tool-workflow/agent-start') return
  20. started.add(event.data.childId)
  21. pending.get(event.data.childId)?.resolve()
  22. })
  23. ctx.on('agent/pre-step', async ({ agent, messages, turn, step, signal }, next) => {
  24. if (agent.session.header.parentSession !== config.parentSessionId || turn !== 1 || step !== 1
  25. || !messages.some(message => message.content.some(block => block.type === 'text' && block.text === config.prompt))) {
  26. return next()
  27. }
  28. signal.throwIfAborted()
  29. if (disposed) throw new Error('workflow snapshot barrier disposed')
  30. if (!started.has(agent.id)) {
  31. const wait = Promise.withResolvers()
  32. const abort = () => { wait.reject(signal.reason) }
  33. signal.addEventListener('abort', abort, { once: true })
  34. wait.done = wait.promise.finally(() => {
  35. signal.removeEventListener('abort', abort)
  36. pending.delete(agent.id)
  37. })
  38. pending.set(agent.id, wait)
  39. await wait.done
  40. }
  41. signal.throwIfAborted()
  42. if (disposed) throw new Error('workflow snapshot barrier disposed')
  43. return next()
  44. })
  45. }