Skip to content
Merged
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

### ✨ Added

- **Container Status / Progress (`mageflow.astatus`)** (#134): Container signatures (swarms and chains) now expose an `astatus()` method returning a structured `ContainerStatus` (total / finished / failed / running / pending, terminal-state percentage, and completion flag). The new `mageflow.astatus(*ids)` loads several containers in a single Redis lookup and returns a `ContainersStatus` with an aggregate `overall_percentage`, raising on missing or non-container ids.
- **Signing Hatchet Workflows (`MageWorkflow`)**: Native Hatchet `Workflow` objects can now be tracked by mageflow's signature lifecycle, enabling status callbacks (success/failure) without wrapping tasks in mageflow decorators.

### 🐛 Fixed
Expand Down
8 changes: 8 additions & 0 deletions docs/api/chain.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,3 +81,11 @@ async def interrupt()
```

Interrupts all tasks in the chain and sets the status to `INTERRUPTED`.

#### `astatus()`

Return a `ContainerStatus` describing the chain's progress. Child tasks are classified by their `task_status` (done / failed / running / pending) and the percentage is terminal-state based. See [`mageflow.astatus`](functions.md#mageflowastatussignature_ids) for the model fields and the batch helper that covers several containers at once.

```python
async def astatus() -> ContainerStatus
```
9 changes: 9 additions & 0 deletions docs/api/client.md
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,15 @@ Create a task swarm.
swarm = await client.aswarm(tasks=[task1, task2], task_name="my-swarm")
```

#### `astatus(*signature_ids)`

Report the progress of one or more container signatures (swarms / chains). See [`mageflow.astatus`](functions.md#mageflowastatussignature_ids) for the returned `ContainersStatus` / `ContainerStatus` models.

```python
status = await client.astatus(swarm.key, chain.key)
print(status.overall_percentage)
```

#### `with_ctx`

Override the default parameter configuration to enable context for a specific task.
Expand Down
60 changes: 60 additions & 0 deletions docs/api/functions.md
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,66 @@ async def load_signature(key: RapyerKey) -> Optional[Signature]
signature = await mageflow.load_signature(task_key)
```

## Status

### `mageflow.astatus(*signature_ids)`

Report the progress of one or more **container** signatures (swarms / chains) in a single Redis lookup.

```python
async def astatus(*signature_ids: RapyerKey) -> ContainersStatus
```

**Parameters:**

- `*signature_ids` (RapyerKey): Keys of the container signatures to inspect.

**Returns:** a `ContainersStatus` model describing every requested container. Calling it with no ids returns an empty `ContainersStatus` without touching Redis.

**Raises:**

- `MissingSignatureError`: one of the keys does not exist.
- `NotAContainerError`: one of the keys points to a non-container signature (e.g. a plain `TaskSignature`).

```python
status = await mageflow.astatus(swarm.key, chain.key)
for container in status.containers:
print(container.task_name, container.percentage, container.is_done)
print("overall", status.overall_percentage)
```

The percentage is terminal-state based — `(finished + failed) / total * 100` — so it reflects how many child tasks have reached a final state.

The result models are defined in `thirdmagic` and imported from there:

```python
from thirdmagic import ContainerStatus, ContainersStatus
```

#### `ContainerStatus`

Per-container breakdown returned inside `ContainersStatus.containers`.

| Field | Type | Description |
| --- | --- | --- |
| `signature_id` | RapyerKey | Key of the container signature |
| `task_name` | str | The container's task name |
| `status` | SignatureStatus | The container's own lifecycle status |
| `total` | int | Total child tasks |
| `finished` | int | Children that completed successfully |
| `failed` | int | Children that errored |
| `running` | int | Children currently running |
| `pending` | int | Children not yet started |
| `percentage` | float | `(finished + failed) / total * 100`, `0.0` when `total == 0` |
| `is_done` | bool | Whether the container itself has completed |

#### `ContainersStatus`

| Member | Type | Description |
| --- | --- | --- |
| `containers` | list[ContainerStatus] | One entry per requested container |
| `overall_percentage` | float (property) | Terminal-state percentage aggregated across all containers |

## Atomic Operations

### `mageflow.abounded_field(ignore_redis_error=False)`
Expand Down
8 changes: 8 additions & 0 deletions docs/api/swarm.md
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,14 @@ Check if swarm has completed all tasks.
async def is_swarm_done() -> bool
```

#### `astatus()`

Return a `ContainerStatus` describing the swarm's progress (counts + terminal-state percentage), computed from its bookkeeping lists without loading child tasks. See [`mageflow.astatus`](functions.md#mageflowastatussignature_ids) for the model fields and the batch helper that covers several containers at once.

```python
async def astatus() -> ContainerStatus
```

## Error Classes

### TooManyTasksError
Expand Down
2 changes: 2 additions & 0 deletions libs/mageflow/mageflow/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
from thirdmagic import abounded_field
from thirdmagic.chain.creator import chain as achain
from thirdmagic.signature import Signature
from thirdmagic.status import astatus
from thirdmagic.swarm.creator import swarm as aswarm
from thirdmagic.task import TaskSignature
from thirdmagic.task import sign as asign
Expand Down Expand Up @@ -39,6 +40,7 @@ async def load_sign(key: RapyerKey) -> Signature:
"SignatureTTLConfig",
"achain",
"aswarm",
"astatus",
"start_mageflow",
"abounded_field",
]
6 changes: 6 additions & 0 deletions libs/mageflow/mageflow/clients/hatchet/mageflow.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,13 @@
)
from hatchet_sdk.runnables.workflow import BaseWorkflow, Standalone
from hatchet_sdk.worker.worker import LifespanFn
from rapyer.fields import RapyerKey
from redis.asyncio import Redis
from thirdmagic import chain, sign
from thirdmagic.chain import ChainTaskSignature
from thirdmagic.signature import Signature
from thirdmagic.signature.status import ContainersStatus
from thirdmagic.status import astatus
from thirdmagic.swarm import SwarmTaskSignature
from thirdmagic.swarm.creator import SignatureOptions, swarm
from thirdmagic.task import TaskInputType, TaskSignature, TaskSignatureConvertible
Expand Down Expand Up @@ -329,6 +332,9 @@ async def aswarm(
):
return await swarm(tasks, task_name, **kwargs)

async def astatus(self, *signature_ids: RapyerKey) -> ContainersStatus:
return await astatus(*signature_ids)

def with_ctx(self, func):
func.__user_ctx__ = True
return func
Expand Down
118 changes: 118 additions & 0 deletions libs/mageflow/tests/unit/workflows/test_astatus.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
import pytest
from thirdmagic import ContainersStatus, ContainerStatus
from thirdmagic.errors import MissingSignatureError, NotAContainerError
from thirdmagic.signature.status import SignatureStatus

import mageflow
from tests.integration.hatchet.models import ContextMessage
from tests.unit.workflows.conftest import create_swarm_item_test_setup


@pytest.mark.asyncio
async def test_astatus_returns_status_for_single_container(mock_adapter):
# Arrange
setup = await create_swarm_item_test_setup(
num_tasks=4,
stop_after_n_failures=None,
current_running=1,
tasks_left_indices=[3],
finished_indices=[0],
failed_indices=[1],
)
expected = ContainersStatus(
containers=[
ContainerStatus(
signature_id=setup.swarm_task.key,
task_name="test_swarm",
status=SignatureStatus.PENDING,
total=4,
finished=1,
failed=1,
running=1,
pending=1,
percentage=50.0,
is_done=False,
)
]
)

# Act
result = await mageflow.astatus(setup.swarm_task.key)

# Assert
assert result == expected


@pytest.mark.asyncio
async def test_astatus_aggregates_multiple_containers(mock_adapter):
# Arrange
first = await create_swarm_item_test_setup(
num_tasks=2, stop_after_n_failures=None, finished_indices=[0, 1]
)
second = await create_swarm_item_test_setup(num_tasks=2, stop_after_n_failures=None)
expected = ContainersStatus(
containers=[
ContainerStatus(
signature_id=first.swarm_task.key,
task_name="test_swarm",
status=SignatureStatus.PENDING,
total=2,
finished=2,
failed=0,
running=1,
pending=0,
percentage=100.0,
is_done=False,
),
ContainerStatus(
signature_id=second.swarm_task.key,
task_name="test_swarm",
status=SignatureStatus.PENDING,
total=2,
finished=0,
failed=0,
running=1,
pending=0,
percentage=0.0,
is_done=False,
),
]
)

# Act
result = await mageflow.astatus(first.swarm_task.key, second.swarm_task.key)

# Assert
assert result == expected
# 2 terminal tasks out of 4 total across both swarms
assert result.overall_percentage == 50.0


@pytest.mark.asyncio
async def test_astatus_empty_ids_returns_empty_without_db_scan(mock_adapter):
# Act
result = await mageflow.astatus()

# Assert
assert isinstance(result, ContainersStatus)
assert result.containers == []


@pytest.mark.asyncio
async def test_astatus_raises_for_non_container(mock_adapter):
# Arrange
task = await mageflow.asign("plain_task", model_validators=ContextMessage)

# Act / Assert
with pytest.raises(NotAContainerError):
await mageflow.astatus(task.key)


@pytest.mark.asyncio
async def test_astatus_raises_for_missing_signature(mock_adapter):
# Arrange
missing_key = "SwarmTaskSignature:00000000-0000-0000-0000-000000000000"

# Act / Assert
with pytest.raises(MissingSignatureError):
await mageflow.astatus(missing_key)
Loading
Loading