Repository navigation
Python: fix(core): Distributed agent orchestration with pluggable BackgroundTaskRuntimeStore - #9040
Conversation
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Distributed cancellation, continuation, failure handling, and timeout logic contain correctness and concurrency issues.
Review effort: Balanced
Findings: 4
Open (10)
Use per-task cancellation generations · New Do not mark tasks LOST during store outages · New Prevent release races during task acquisition · New Fence stale publishers with attempt generations · New Re-export the pluggable store protocol · New Use shared-store TTL or server time for distributed expiry · New Require atomic lease acquisition and stale-outcome removal · New Require atomic outcome publication and lease removal · New Handle publication failures with retries or durable handoff · New Avoid repeated LOST-state writes during refresh · New
What changed in this PR
Adds distributed background-agent orchestration through a pluggable runtime store.
Changes:
- Introduces lease renewal, outcome publication, and remote cancellation.
- Adds cross-process task refresh and polling.
- Expands distributed-runtime tests.
| File | Description |
|---|---|
_background_agents.py |
Implements distributed task coordination. |
test_harness_background_agents.py |
Adds cross-process runtime-store tests. |
💡 Add a code-review agent skill for context-aware, tailored reviews. Learn more in the docs.
| Polled by :func:`_run_agent_with_renewal` on each renewal tick so that a | ||
| cancel signal from another process causes the owning process to abort the | ||
| task. Implementations may clear the flag after returning ``True`` to | ||
| avoid repeated cancellations, or leave it set. |
There was a problem hiding this comment.
A session-scoped cancel flag is by design for this initial implementation. The intention is that when release_session is called, all currently running tasks under that session ID are aborted. While a generation-based token system would allow safe reuse of the exact same session ID for future tasks, it requires much heavier distributed state tracking. For now, the assumption is that a cancelled session is discarded, and new work should use a new session ID.
| try: | ||
| await runtime_store.acquire(qualified_session_key, task_id, ttl_seconds=DEFAULT_LEASE_TTL_SECONDS) | ||
| except Exception: | ||
| logger.debug("Failed to acquire lease for continued task %s", task_id) |
There was a problem hiding this comment.
This is a valid theoretical race condition where a delayed publication from a dying worker could overwrite the slot of a new continuation. However, introducing a fencing token or attempt-generation mechanism requires pervasive changes to how the framework tracks task identifiers across all process boundaries. For this initial distributed store implementation, we accept this edge case in order to keep the protocol simple.
…kgroundTaskRuntimeStore
1918856 to
78a8d9c
Compare
…ub.com/karthik-0306/microsoft-agent-framework into fix/8760-background-task-runtime-store
…-task-runtime-store



Motivation & Context
BackgroundAgentsProvidertracks running tasks using only in-processasyncio.Taskreferences. In any multi-replica deployment — where successive turns of the same session can land on different worker processes — this causes several cascading failures described in #8760:LOSTwith no path to recovery._refresh_task_stateskipped records already markedLOST, so once a task was marked lost, its result was permanently discarded even if the owning process finished successfully.LOSTto another process.release_sessionon non-owner is a silent no-op: The early-return guard prevented the cancel signal from reaching the owning process.background_agents_wait_for_first_completionhad no store-polling path, so a non-owning process would immediately error instead of waiting.Description & Review Guide
What are the major changes?
1.
BackgroundTaskRuntimeStoreprotocol (_background_agents.py)A new
@experimentalProtocolclass with six async methods:acquire(key, task_id, ttl_seconds)task_idso a continued task can write a fresh result (first-write-wins is reset on re-acquire).renew(key, task_id, ttl_seconds)ttl/2intervals.is_alive(key, task_id)Trueif the lease exists and has not expired (wall-clocktime.time()).publish_outcome(key, task_id, info)BackgroundTaskInfo. First-write-wins; subsequent calls are silently ignored.fetch_outcome(key, task_id)Noneif not yet published.request_cancel(key)The default
_InMemoryBackgroundTaskRuntimeStorereproduces the existing single-process behaviour exactly — no change to deployments that don't supply a store.2.
_run_agent_with_renewalwrapperReplaces the bare
asyncio.create_task(_run_agent(...))call. Runs the agent coroutine as a sub-task while periodically callingstore.renew()atttl_seconds / 2. Cancellation (local or remote) propagates correctly viaBaseExceptionre-raise +work_task.cancel().3.
_refresh_task_state— fixed LOST recoveryThe loop guard changed from:
to processing both
RUNNINGandLOSTrecords. For tasks with no localasyncio.Task, the function now:store.is_alive()— ifTrue, leaves status asRUNNING(cross-process task still healthy).store.fetch_outcome()— if a result is present, finalizes via_finalize_task_from_info().LOSTif neither check yields a result.4. Done-callback via
_make_done_callbackAttached to every
asyncio.Taskat start and continue time. Fires the instant the task settles and schedulesstore.publish_outcome()vialoop.create_task(). The outcome is available to any process on the very next turn.5.
release_session— cancel signal fires before early-return guardThe
store.request_cancel()call is now placed beforeruntime = self._runtime.get(session_id), so a non-owning process (which has no local runtime entry) still propagates the signal to the process that owns the tasks. The call is gated oncancel_running=Truesorelease_session(..., cancel_running=False)does not issue a remote kill.6.
background_agents_wait_for_first_completion— store-polling fallbackRewritten as a time-bounded poll loop. On each iteration it calls
_refresh_task_state(which now queries the store), then either waits on localasyncio.Taskobjects or falls back toasyncio.sleep(1.0)when no local tasks are present — correctly handling the cross-process wait case.What is the impact of these changes?
_InMemoryBackgroundTaskRuntimeStoreis the default and replicates existing semantics exactly.BackgroundTaskRuntimeStore(e.g. Redis-backed) viaBackgroundAgentsProvider(agents, runtime_store=my_store). All six failure scenarios described in [Feature]: BackgroundAgentsProvider cannot tell a task that is running on another process from one whose owner died #8760 are resolved.runtime_storekwarg is optional and defaults to the in-memory implementation.What do you want reviewers to focus on?
_refresh_task_stateloop guard change (the LOST-state recovery fix) — this is the highest-risk edit since it changes which records are processed on every refresh._make_done_callbackfirst-write-wins semantics — ensures a late or duplicate done-callback cannot overwrite an already-committed outcome.release_sessionordering —request_cancelmust fire before the early-return guard; confirm the placement is correct._run_agent_with_renewalcancellation path — verifyBaseExceptioncatch +work_task.cancel()+await work_taskis the right cleanup sequence.Related Issue
Fixes #8760
Contribution Checklist