feat: add DBOS integration for durable Agent execution - #3767
Closed
julian-risch wants to merge 2 commits into
Closed
julian-risch wants to merge 2 commits into
julian-risch wants to merge 2 commits into
Conversation
Adds `dbos-haystack`, which checkpoints Haystack Agent runs with DBOS (https://docs.dbos.dev) so they survive crashes and can suspend for days awaiting human approval. - `DBOSChatGenerator` wraps any Chat Generator and runs each model call inside a DBOS step, recording the replies in Haystack's dictionary format. Outside a workflow it is transparent, so the same object works without DBOS launched. - `durable_agent()` returns a copy of an Agent with its generator wrapped. The user writes their own `@DBOS.workflow`, following the shape Pydantic AI moved to after deprecating its `DBOSAgent` wrapper: DBOS persists workflow inputs, and a streaming callback or a lambda-backed Tool cannot be pickled. - `DBOSConfirmationStrategy` plugs into the stock `ConfirmationHook` and waits on `DBOS.recv`, so a run survives a restart while it waits and a recovered run reuses the decision instead of re-prompting. Tool calls are not checkpointed in this version and re-execute on recovery; the README documents that alongside the other replay caveats. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Contributor
Coverage report (dbos)Click to see where and how coverage changed
This report was generated by python-coverage-comment-action |
||||||||||||||||||||||||||||||||||||||||||
|
Review the following changes in direct dependencies. Learn more about Socket for GitHub.
|
…ersion-tolerant CI's lowest-direct resolution installed dbos 2.0.0, which does not export `StepOptions` or `DBOS.run_step`, so collection failed with an ImportError. 2.9.0 is the first release that has both. Testing against that floor also showed two assertions that only held on recent dbos: whether draining a cancelled workflow's handle raises, and whether resuming re-executes the workflow body. Neither is what those tests are about, so they now assert the durable property itself - the resumed run does not wait for a fresh decision. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Related Issues
Proposed Changes:
Adds
dbos-haystack, a new integration that makes Haystack Agent runs durable with DBOS — an MIT-licensed, in-process library backed by Postgres or SQLite, with no external orchestrator to deploy.DBOSChatGeneratorwraps any Chat Generator and runs each model call inside a@DBOS.step(). The step records the replies asChatMessage.to_dict()rather than objects, so the checkpoint stays readable, survives DBOS's portable-JSON serializer, and doesn't unpickle into a changed dataclass ifhaystack-aiis upgraded between the crash and the recovery. Outside a workflow the wrapper is transparent, so the same object works in tests and in applications that never launch DBOS.durable_agent()returns a copy of an Agent with its generator wrapped, and warns when the Agent carries a toolset that discovers its tools while running.DBOSConfirmationStrategyplugs into the stockConfirmationHookwith no new hook class. It publishes the pending tool call as a DBOS event and waits onDBOS.recv, so a run can suspend for hours or days awaiting approval and survive a restart while it waits.The user writes their own
@DBOS.workflowaroundagent.run(). There is deliberately noDBOSAgentwrapper that owns the decorator: DBOS persists workflow inputs, so promotingAgent.runto a workflow would require picklingmessages,streaming_callbackand runtimetools— a plain callback or a lambda-backedToolmakes the workflow fail to start, with the error surfacing deep inside DBOS. Pydantic AI hit exactly this and deprecated its ownDBOSAgentin favour of this shape.Recovery in DBOS is replay, not forward-resume: the workflow body re-executes and completed steps return their recorded output. Tool calls are not checkpointed in this version and therefore run again, so tools must be idempotent or be steps themselves. That is documented in a durability-guarantees table in the README, together with the other consequences (streaming does not replay, tracing spans are re-created, dynamic toolsets are unsupported, sync and async must not be mixed). Auto-wrapping tools is deliberately left out for now — see "Notes for the reviewer".
Also included: a README with a quickstart and the caveats, a runnable
examples/durable_agent.pythat crashes and then recovers, pydoc config, and the scaffold-generated CI workflow, labeler entry and README inventory row.How did you test it?
From
integrations/dbos:hatch run fmt-check,hatch run test:types,hatch run test:unit(54 passed) andhatch run docsall pass locally.The whole suite runs against in-process SQLite, so no service container and no secrets are needed — the workflow keeps the "no integration tests yet" comment in place of the integration-test steps.
Replay is tested two ways, both taking the same path a crash-restart would, and neither needing a subprocess or a sleep:
fork_workflowrestarts a failed workflow from a chosen step.cancel_workflow→resume_workflowpicks up a workflow interrupted mid-flight.Both assert the fake generator was called exactly once across the original and recovered runs. A companion test asserts an unwrapped tool's counter reaches 2, so the at-least-once contract is pinned as tested behaviour rather than only a documentation claim. The human-in-the-loop suite covers confirm, reject, modify, timeout, and — the decisive one — that a recovered run reuses the recorded decision and executes the tool without a second
DBOS.send.Not covered by automated tests: a real
kill -9against a live process, and Postgres as the system database.examples/durable_agent.pyexercises the former by hand against a realOpenAIChatGenerator.Notes for the reviewer
Opening as a draft — the proposal issue is still open, and two decisions are worth your input before this lands:
Agent.cloneis not inhaystack-ai==3.0.0, sodurable_agentrebuilds the Agent from its init signature. That surfaced a second divergence: on 3.0.0,Agent.state_schemaholds the resolved schema including reserved keys, which cannot be passed back to__init__; the value that was passed in lives in_state_schema. Handled by a documented attribute map indurability.py, which keeps working oncecloneships. Alternatively we wait for that release and raise thehaystack-aifloor.ThreadPoolExecutorwith a copied context, and DBOS assigns step ids by mutating a counter on a single shared context object. Withtool_concurrency_limit > 1two@DBOS.steptools can be handed the same id, which either fails recovery loudly or — when two calls to the same tool collide — returns one call's cached result for the other. The async path is safe, since tool coroutines are created in call order and the semaphore wakes waiters FIFO. Shipping an auto-wrapper needs that resolved first; the README documents the manual@DBOS.stepescape hatch and its constraints in the meantime.Related, though solving the problem from the other direction: deepset-ai/hayhooks#253 adds a Redis-backed durable engine inside the Hayhooks server, checkpointing Agent
Statethrough the hook seams and resuming forward. This integration is the library-shaped counterpart — plain Python, Postgres or SQLite, replay-based, no engine to operate. Worth a view on whether the READMEs should cross-reference each other once that lands.Suggested review order:
chat_generator.pyfor the step boundary,durability.pyfor the rebuild and guard rails,human_in_the_loop.pyfor the durable wait, thentests/test_recovery.py.Checklist
fix:,feat:,build:,chore:,ci:,docs:,style:,refactor:,perf:,test:.🤖 Generated with Claude Code