Background Task Manager

In shortBackgroundTaskManager — launch and track asyncio tasks from node functions without blocking the graph response.

  • 4 min read
  • 7 sections
  • Updated
  • v0.9.2
  • Markdown

When to use this

Use BackgroundTaskManager when a node needs to do slow I/O (send a webhook, update a CRM, flush telemetry) without delaying the response to the caller. The graph returns immediately; the background task runs concurrently and is cleaned up on shutdown.

Import path

Python
from tenxgraph.utils.background_task_manager import BackgroundTaskManager, TaskMetadata
# also available from the top-level package:
from tenxgraph.utils import BackgroundTaskManager

Getting the instance in a node

BackgroundTaskManager is registered in the dependency container automatically by StateGraph. Declare it as a parameter and the framework injects it:

Python
from tenxgraph.utils.background_task_manager import BackgroundTaskManager

async def my_node(
    state,
    config: dict,
    task_manager: BackgroundTaskManager,  # auto-injected
) -> ...:
    task_manager.create_task(
        send_webhook(config.get("user_id")),
        name="send_webhook",
        timeout=15.0,
    )
    return state

BackgroundTaskManager

Constructor

Python
manager = BackgroundTaskManager(
    default_shutdown_timeout=30.0,
    max_pending_tasks=1000,
)
Parameter Type Default Description
default_shutdown_timeout float 30.0 Seconds to wait when draining tasks during shutdown().
max_pending_tasks int 1000 Cap on in-flight tasks. Beyond this, new tasks are dropped rather than queued. Set to 0 to disable the cap and restore unbounded growth.

Backpressure

Background tasks are fire-and-forget telemetry: publisher writes, webhook posts, memory syncs. When the sink cannot keep up, an unbounded queue grows until the process runs out of memory. max_pending_tasks bounds it instead.

Once max_pending_tasks tasks are in flight, create_task drops the newest task and returns None. Dropping the newest is deliberate: cancelling an already-running task would throw away work that is mid-flight, while shedding a new event only loses one event. The coroutine is closed so Python does not emit a “coroutine was never awaited” warning.

Drops are counted on the background_task_manager.tasks_dropped metric and logged at warning level, rate-limited so a sustained overload does not flood the log.

Because a task can be dropped, treat the return value as optional:

Python
task = manager.create_task(send_webhook(user_id), name="send_webhook")
if task is None:
    ...  # shed under backpressure

create_task

Python
task = manager.create_task(
    coro,
    name="my_task",
    timeout=None,
    context=None,
)
Parameter Type Default Description
coro Coroutine required The coroutine to execute in the background.
name str "background_task" Human-readable label. Appears in logs and get_task_info().
timeout float | None None Cancel the task after this many seconds. Logs a warning on timeout.
context dict | None None Extra key/value pairs attached to log entries and task info for debugging.

Returns the created asyncio.Task, or None when the task was dropped because max_pending_tasks was already in flight.

get_task_count

Python
n = manager.get_task_count()  # → int

Number of tasks that are still running.

get_task_info

Python
infos = manager.get_task_info()  # → list[dict]

Each dict contains:

Key Type Description
name str Task name.
age_seconds float Seconds since the task was created.
timeout float | None Configured timeout.
context dict Context passed at creation.
done bool Whether the task has finished.
cancelled bool Whether the task was cancelled (only meaningful when done=True).

wait_for_all

Python
await manager.wait_for_all(timeout=30.0, return_exceptions=False)

Wait for all tracked tasks to complete. Logs a warning if timeout is exceeded.

Parameter Type Default Description
timeout float | None None Max seconds to wait. None means wait forever.
return_exceptions bool False If True, exceptions are returned as results rather than raised.

cancel_all

Python
await manager.cancel_all()

Cancel every tracked task immediately. Does not wait for cancellation to propagate.

shutdown

Python
stats = await manager.shutdown(timeout=30.0)

Graceful shutdown: cancels all tasks, waits up to timeout seconds, then force-cancels any remaining. Returns a stats dict:

Key Type Description
status str "completed" or "timed_out".
initial_tasks int Number of tasks at shutdown start.
completed_tasks int Tasks that finished cleanly.
remaining_tasks int Tasks still alive after timeout.
duration_seconds float Total shutdown duration.

TaskMetadata

Dataclass holding per-task tracking info. Available via get_task_info().

Field Type Description
name str Task name.
created_at float Unix timestamp of creation.
timeout float | None Configured timeout.
context dict | None Extra context.

Lifecycle and shutdown

BackgroundTaskManager is created once per StateGraph during __init__. It is passed to CompiledGraph at compile(). When you call app.aclose(), the graph calls task_manager.shutdown(timeout=shutdown_timeout), which drains or cancels any running background tasks before the process exits.

Python
app = graph.compile(shutdown_timeout=30.0)  # shutdown_timeout flows to BackgroundTaskManager

# In your process teardown:
await app.aclose()   # or await app.astop()

Common errors

Error Cause Fix
Background task never starts Coroutine function passed instead of coroutine object. create_task(send(x)) not create_task(send).
task_manager is None in node Node is not inside a compiled graph. Ensure the node function runs inside a graph that was compiled.
Task silently fails, no log Exception swallowed. Check for background_task_manager.tasks_failed metric; errors are logged at ERROR level.
Tasks run after graph has conceptually “finished” aclose() not called. Always call await app.aclose() on process exit.
Last updated for v0.9.2Edit this page on GitHubReport an issue