Skip to main content

Lightweight asyncio task tracking as call tree and DAG

Project description

aionode

Lightweight asyncio task tracking as call tree and DAG.

Installation

pip install aionode

Quick Start

import asyncio
import aionode

async def fetch_data() -> list[int]:
    await asyncio.sleep(1)
    return list(range(100))

async def process(data: list[int]) -> int:
    await asyncio.sleep(0.5)
    return sum(data)

async def pipeline() -> None:
    async with asyncio.TaskGroup() as tg:
        fetch = tg.create_task(aionode.node(fetch_data)(), name="fetch")
        tg.create_task(aionode.node(process)(aionode.resolve(fetch)), name="process")

async def main() -> None:
    root = asyncio.create_task(aionode.node(pipeline)(), name="pipeline")
    await root
    for info in aionode.walk_dag():
        print(f"{info.name}: {info.status} ({info.duration():.2f}s)")

asyncio.run(main())

Core Concepts

node(func, wait_for=[], track=True, auto_progress=True)

Wraps an async function as a DAG node. Use resolve() to pass awaitables as arguments — they are gathered concurrently before the function is called. Sync functions must be wrapped with make_async first.

# Async function — pass upstream tasks with resolve()
fetch = tg.create_task(aionode.node(fetch_data)(), name="fetch")
process_task = tg.create_task(aionode.node(process)(aionode.resolve(fetch)), name="process")

# Sync function — wrap with make_async first
summarize = tg.create_task(
    aionode.node(aionode.make_async(my_sync_fn))(aionode.resolve(process_task)),
    name="summarize",
)

# Side-only dependency (no value passed): use wait_for
task_b = tg.create_task(aionode.node(cleanup, wait_for=[fetch])(), name="cleanup")

resolve(awaitable)

Marks an awaitable to be resolved before being passed as an argument to node(). This preserves type information — the type checker sees resolve(task: Task[T]) as returning T.

result = await aionode.node(process)(aionode.resolve(upstream_task))

walk_tree / walk_dag

Iterate over tracked tasks after they have been created.

# DFS pre-order through the call tree (parent → children)
for info in aionode.walk_tree():
    print(info.name, info.status)

# Topological order over the DAG (deps → dependents)
for info in aionode.walk_dag():
    print(info.name, info.status)

Both accept an optional root argument (an asyncio.Task or task id). Passing None (the default) includes all tasks in the current event loop.

Task Inspection

task_id = await aionode.get_task_id(asyncio_task)
info = aionode.get_task_info(task_id)

info.status        # TaskStatus: WAITING, RUNNING, DONE, FAILED, CANCELLED
info.duration()    # Elapsed seconds
info.name          # Task name
info.deps          # Upstream dependency IDs
info.dependents    # Downstream dependent IDs
info.logs          # Accumulated log output

await aionode.log("processing record 42")  # Append to current task's logs

# Get TaskInfo for the currently running tracked task
info = aionode.current_task_info()

API Reference

Function Description
node(func, wait_for, track, auto_progress) Wrap an async function as a DAG node
resolve(awaitable) Mark an awaitable to be resolved as a node argument
get_task_id(task, timeout) Get the task ID for an asyncio.Task
get_task_info(task_id) Get TaskInfo by ID
current_task_info() Get TaskInfo for the currently running tracked task
remove_task(task_id) Remove a task and its descendants
log(value, end) Append to the current task's logs
make_async(func) Run a sync function in a thread
make_async_generator(gen) Async iterate a sync iterator via threads
walk_tree(root) DFS pre-order through the call tree
walk_dag(root) Topological order over the DAG

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

aionode-0.1.2.tar.gz (7.2 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

aionode-0.1.2-py3-none-any.whl (8.0 kB view details)

Uploaded Python 3

File details

Details for the file aionode-0.1.2.tar.gz.

File metadata

  • Download URL: aionode-0.1.2.tar.gz
  • Upload date:
  • Size: 7.2 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.7

File hashes

Hashes for aionode-0.1.2.tar.gz
Algorithm Hash digest
SHA256 2b4358148c35096dc5d8992d9be6596f62188e134348c40c641becb1f4f5b567
MD5 668a79d3fb7cab71f6319577d586b594
BLAKE2b-256 6adba52dd619147236f1d1412b8fc9d2011475613745359797a8ebabf3cd7782

See more details on using hashes here.

Provenance

The following attestation bundles were made for aionode-0.1.2.tar.gz:

Publisher: publish.yml on MatteoDep/aionode

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file aionode-0.1.2-py3-none-any.whl.

File metadata

  • Download URL: aionode-0.1.2-py3-none-any.whl
  • Upload date:
  • Size: 8.0 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.7

File hashes

Hashes for aionode-0.1.2-py3-none-any.whl
Algorithm Hash digest
SHA256 4d2fd6cfb0cdc7e9ecc6afac6af0692090d46b4883327f69d8bd424e318d4a88
MD5 b04305fda65b7f09c62612158aa8b952
BLAKE2b-256 9389a46a904c6e16986f12a8715705d883dd375a4d66c4e8c9e5156b5424e25b

See more details on using hashes here.

Provenance

The following attestation bundles were made for aionode-0.1.2-py3-none-any.whl:

Publisher: publish.yml on MatteoDep/aionode

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page