Skip to content

Python: fix(core): Distributed agent orchestration with pluggable BackgroundTaskRuntimeStore - #9040

Open
Karthik Thota (karthik-0306) wants to merge 7 commits into
microsoft:mainfrom
karthik-0306:fix/8760-background-task-runtime-store
Open

Karthik Thota (karthik-0306) wants to merge 7 commits into
microsoft:mainfrom
karthik-0306:fix/8760-background-task-runtime-store

Conversation

@karthik-0306

Copy link
Copy Markdown
Contributor

Motivation & Context

BackgroundAgentsProvider tracks running tasks using only in-process asyncio.Task references. 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:

  1. Cross-process invisibility: A task running on Process A is invisible to Process B, which immediately marks it LOST with no path to recovery.
  2. Sticky LOST state: _refresh_task_state skipped records already marked LOST, so once a task was marked lost, its result was permanently discarded even if the owning process finished successfully.
  3. Result never delivered cross-process: There was no mechanism to publish a completed task's result out of the owning process. A non-owning process had no way to read the outcome.
  4. No liveness heartbeat: No renewal mechanism existed. A long-running task's owner would be considered dead the moment the lease TTL elapsed.
  5. Expired lease not surfaced: A task whose owner process died had no way to surface as LOST to another process.
  6. release_session on non-owner is a silent no-op: The early-return guard prevented the cancel signal from reaching the owning process.
  7. Wait tool blocks only on local tasks: background_agents_wait_for_first_completion had no store-polling path, so a non-owning process would immediately error instead of waiting.

Description & Review Guide

What are the major changes?

1. BackgroundTaskRuntimeStore protocol (_background_agents.py)

A new @experimental Protocol class with six async methods:

Method Purpose
acquire(key, task_id, ttl_seconds) Write initial lease when a task starts or is continued. Clears any stale outcome for that task_id so a continued task can write a fresh result (first-write-wins is reset on re-acquire).
renew(key, task_id, ttl_seconds) Extend lease before it expires. Called by the renewal wrapper at ttl/2 intervals.
is_alive(key, task_id) Returns True if the lease exists and has not expired (wall-clock time.time()).
publish_outcome(key, task_id, info) Persist the terminal BackgroundTaskInfo. First-write-wins; subsequent calls are silently ignored.
fetch_outcome(key, task_id) Return the stored outcome, or None if not yet published.
request_cancel(key) Signal the owning process to cancel all tasks for this session.

The default _InMemoryBackgroundTaskRuntimeStore reproduces the existing single-process behaviour exactly — no change to deployments that don't supply a store.

2. _run_agent_with_renewal wrapper

Replaces the bare asyncio.create_task(_run_agent(...)) call. Runs the agent coroutine as a sub-task while periodically calling store.renew() at ttl_seconds / 2. Cancellation (local or remote) propagates correctly via BaseException re-raise + work_task.cancel().

3. _refresh_task_state — fixed LOST recovery

The loop guard changed from:

if task_info.status != BackgroundTaskStatus.RUNNING:
    continue

to processing both RUNNING and LOST records. For tasks with no local asyncio.Task, the function now:

  1. Calls store.is_alive() — if True, leaves status as RUNNING (cross-process task still healthy).
  2. Calls store.fetch_outcome() — if a result is present, finalizes via _finalize_task_from_info().
  3. Only falls back to LOST if neither check yields a result.

4. Done-callback via _make_done_callback

Attached to every asyncio.Task at start and continue time. Fires the instant the task settles and schedules store.publish_outcome() via loop.create_task(). The outcome is available to any process on the very next turn.

5. release_session — cancel signal fires before early-return guard

The store.request_cancel() call is now placed before runtime = 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 on cancel_running=True so release_session(..., cancel_running=False) does not issue a remote kill.

6. background_agents_wait_for_first_completion — store-polling fallback

Rewritten as a time-bounded poll loop. On each iteration it calls _refresh_task_state (which now queries the store), then either waits on local asyncio.Task objects or falls back to asyncio.sleep(1.0) when no local tasks are present — correctly handling the cross-process wait case.


What is the impact of these changes?


What do you want reviewers to focus on?

  • _refresh_task_state loop 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_callback first-write-wins semantics — ensures a late or duplicate done-callback cannot overwrite an already-committed outcome.
  • release_session ordering — request_cancel must fire before the early-return guard; confirm the placement is correct.
  • _run_agent_with_renewal cancellation path — verify BaseException catch + work_task.cancel() + await work_task is the right cleanup sequence.

Related Issue

Fixes #8760


Contribution Checklist

  • The code builds clean without any errors or warnings
  • All unit tests pass, and I have added new tests where possible
  • The PR follows the Contribution Guidelines
  • This PR is linked to an issue and there is no other open PR for this issue (see Related Issue above).
  • This is not a breaking change.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Distributed cancellation, continuation, failure handling, and timeout logic contain correctness and concurrency issues.

Review effort: Balanced
Findings: 4 High severity · 5 Medium severity · 1 Low severity

Open (10)
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.

Comment on lines +240 to +243
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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread python/packages/core/agent_framework/_harness/_background_agents.py Outdated
Comment thread python/packages/core/agent_framework/_harness/_background_agents.py
Comment on lines +1155 to +1158
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)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread python/packages/core/agent_framework/_harness/_background_agents.py
Comment thread python/packages/core/agent_framework/_harness/_background_agents.py Outdated
Comment thread python/packages/core/agent_framework/_harness/_background_agents.py Outdated
Comment thread python/packages/core/agent_framework/_harness/_background_agents.py Outdated
Comment thread python/packages/core/agent_framework/_harness/_background_agents.py Outdated
Comment thread python/packages/core/agent_framework/_harness/_background_agents.py Outdated

This branch was successfully deployed

1 active deployment
github-app-auth — bbde0364 Deployed Oct 10, 2026 by karthik-0306 via add_label #24923
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

python Usage: [Issues, PRs], Target: Python

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Feature]: BackgroundAgentsProvider cannot tell a task that is running on another process from one whose owner died

2 participants