3.9 KiB
Async Resource Safety
Overview
CleverAgents uses asynchronous resources throughout its stack — event subscriptions, LangGraph tasks, checkpoint file I/O, and reactive stream connections. The async-cleanup subsystem ensures that these resources are tracked, closed deterministically on shutdown, and that leaks are detected and logged.
Core Component: AsyncResourceTracker
cleveragents.core.async_cleanup.AsyncResourceTracker is the central
registry for any resource that implements the AsyncResource protocol
(i.e.\ exposes an async def close() -> None method).
Registration
from cleveragents.core.async_cleanup import AsyncResourceTracker
tracker = AsyncResourceTracker()
tracker.register("db-pool", db_pool)
tracker.register("event-queue", event_queue)
- Names must be unique and non-empty.
- Duplicate registrations raise
ValueError. Noneresources are rejected immediately.
Shutdown
await tracker.close_all(timeout=30.0)
close_all() iterates over every registered resource and awaits its
close() coroutine with asyncio.wait_for(). Resources that exceed
the deadline are logged as forced terminations and their names are
collected in tracker.timed_out_resources.
close_all() is idempotent — calling it multiple times is safe.
After close_all(), calling register() raises RuntimeError.
Resources cannot be added to a closed tracker.
Leak Detection
If the tracker is garbage-collected without close_all() having been
called, the __del__ finalizer logs a warning for each unclosed
resource by name. This makes it straightforward to identify leaks
during development and in CI logs.
Async Context Manager
async with AsyncResourceTracker() as tracker:
tracker.register("conn", connection)
# ... use connection ...
# connection.close() is awaited automatically
Enhanced Bridge Cleanup
cleveragents.langgraph.bridge.RxPyLangGraphBridge now provides:
cleanup_tasks_async(timeout)— cancels all in-flight asyncio tasks and awaits their completion within timeout seconds. Tasks that do not finish are logged as warnings.cancel_task_with_reason(task, reason)— cancels a specific task and records the human-readable reason inbridge.cancellation_reasons.
The synchronous cleanup_tasks() remains for best-effort cleanup in
__del__.
Subscription Cleanup
cleveragents.a2a.events.A2aEventQueue.close() removes all local
subscriptions, clears the event buffer, and logs the count of
subscriptions that were active. After close(), calling publish()
raises RuntimeError. The is_closed property exposes the closed
state for callers to check.
Checkpoint File Safety
cleveragents.langgraph.state.StateManager.close() marks the manager
as closed, completes the underlying RxPY BehaviorSubject, and
prevents further state updates. After close(), calling
update_state(), reset(), load_checkpoint(), or time_travel()
raises RuntimeError. Checkpoint files written via
_save_checkpoint() use Path.write_text(), which handles file-handle
closing internally.
Thread Safety
AsyncResourceTracker.register() and close_all() are protected by a
threading.Lock so that resources can be registered from any thread
without races. After close_all() completes, register() raises
RuntimeError to prevent silently leaked resources.
Design Decisions
| Decision | Rationale |
|---|---|
Protocol-based AsyncResource |
Structural typing avoids coupling to a specific base class. |
| Per-resource timeout | One slow resource should not block the entire shutdown sequence. |
Idempotent close_all |
Prevents double-close errors in complex shutdown paths. |
__del__ leak warning |
Best-effort; relies on CPython deterministic GC but degrades safely. |
| Cancellation reason dict | Lightweight tracing for debugging cancelled tasks without heavy instrumentation. |