From b7e73e5692f9a909d353d3a38683894d7724ca1f Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Thu, 23 Jul 2026 21:33:31 +0300 Subject: [PATCH 1/9] feat: add container status result models and NotAContainerError Add ContainerStatus (per-container breakdown with terminal-state percentage) and ContainersStatus (aggregate wrapper) to the signature status module, plus a NotAContainerError for the status lookups. Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- libs/third-magic/thirdmagic/errors.py | 4 ++ .../thirdmagic/signature/__init__.py | 10 +++- .../thirdmagic/signature/status.py | 56 +++++++++++++++++++ 3 files changed, 69 insertions(+), 1 deletion(-) diff --git a/libs/third-magic/thirdmagic/errors.py b/libs/third-magic/thirdmagic/errors.py index 8fe87817..c09e8ce9 100644 --- a/libs/third-magic/thirdmagic/errors.py +++ b/libs/third-magic/thirdmagic/errors.py @@ -28,3 +28,7 @@ class TaskAndMsgsDontMatchForSwarmError(SwarmError, RuntimeError): class UnrecognizedTaskError(MageflowError): pass + + +class NotAContainerError(MageflowError): + pass diff --git a/libs/third-magic/thirdmagic/signature/__init__.py b/libs/third-magic/thirdmagic/signature/__init__.py index afe301f7..cb5e3bed 100644 --- a/libs/third-magic/thirdmagic/signature/__init__.py +++ b/libs/third-magic/thirdmagic/signature/__init__.py @@ -1,6 +1,12 @@ from thirdmagic.signature.model import Signature, SignatureConfig from thirdmagic.signature.retry_cache import SignatureRetryCache, retry_cache_ctx -from thirdmagic.signature.status import PauseActionTypes, SignatureStatus, TaskStatus +from thirdmagic.signature.status import ( + ContainersStatus, + ContainerStatus, + PauseActionTypes, + SignatureStatus, + TaskStatus, +) __all__ = [ "Signature", @@ -9,5 +15,7 @@ "SignatureStatus", "PauseActionTypes", "TaskStatus", + "ContainerStatus", + "ContainersStatus", "retry_cache_ctx", ] diff --git a/libs/third-magic/thirdmagic/signature/status.py b/libs/third-magic/thirdmagic/signature/status.py index bc123efd..7ff2d570 100644 --- a/libs/third-magic/thirdmagic/signature/status.py +++ b/libs/third-magic/thirdmagic/signature/status.py @@ -1,8 +1,10 @@ from enum import Enum from typing import ClassVar +from pydantic import BaseModel from rapyer import AtomicRedisModel from rapyer.config import RedisConfig +from rapyer.fields import RapyerKey class SignatureStatus(str, Enum): @@ -33,3 +35,57 @@ def should_run(self): def is_done(self): return self.status in [SignatureStatus.DONE, SignatureStatus.FAILED] + + +class ContainerStatus(BaseModel): + signature_id: RapyerKey + task_name: str + status: SignatureStatus + total: int + finished: int + failed: int + running: int + pending: int + percentage: float + is_done: bool + + @classmethod + def from_counts( + cls, + signature_id: RapyerKey, + task_name: str, + status: SignatureStatus, + total: int, + finished: int, + failed: int, + running: int, + pending: int, + is_done: bool, + ) -> "ContainerStatus": + percentage = (finished + failed) / total * 100 if total else 0.0 + return cls( + signature_id=signature_id, + task_name=task_name, + status=status, + total=total, + finished=finished, + failed=failed, + running=running, + pending=pending, + percentage=percentage, + is_done=is_done, + ) + + +class ContainersStatus(BaseModel): + containers: list[ContainerStatus] + + @property + def overall_percentage(self) -> float: + total = sum(container.total for container in self.containers) + if not total: + return 0.0 + terminal = sum( + container.finished + container.failed for container in self.containers + ) + return terminal / total * 100 From 7f2453fe339a7ed996271ab89e7529fd3ce5c0b8 Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Thu, 23 Jul 2026 21:33:40 +0300 Subject: [PATCH 2/9] feat: compute container_status for swarm and chain Add an abstract container_status() to ContainerTaskSignature. Swarm derives counts from its bookkeeping lists; chain classifies children by their task_status. Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- libs/third-magic/thirdmagic/chain/model.py | 27 +++++++++++++++++++++- libs/third-magic/thirdmagic/container.py | 5 ++++ libs/third-magic/thirdmagic/swarm/model.py | 15 +++++++++++- 3 files changed, 45 insertions(+), 2 deletions(-) diff --git a/libs/third-magic/thirdmagic/chain/model.py b/libs/third-magic/thirdmagic/chain/model.py index 2cd16365..0ff24689 100644 --- a/libs/third-magic/thirdmagic/chain/model.py +++ b/libs/third-magic/thirdmagic/chain/model.py @@ -7,7 +7,7 @@ from thirdmagic.container import ContainerTaskSignature from thirdmagic.errors import MissingSignatureError -from thirdmagic.signature.status import SignatureStatus +from thirdmagic.signature.status import ContainerStatus, SignatureStatus from thirdmagic.task.model import TaskSignature from thirdmagic.utils import HAS_HATCHET @@ -47,6 +47,31 @@ async def sub_tasks(self) -> list[TaskSignature]: sub_tasks = await rapyer.afind(*self.tasks, skip_missing=True) return cast(list[TaskSignature], sub_tasks) + async def container_status(self) -> ContainerStatus: + sub_tasks = await self.sub_tasks() + finished = failed = running = 0 + for task in sub_tasks: + status = task.task_status.status + if status == SignatureStatus.DONE: + finished += 1 + elif status == SignatureStatus.FAILED: + failed += 1 + elif status == SignatureStatus.ACTIVE: + running += 1 + total = len(self.tasks) + pending = total - finished - failed - running + return ContainerStatus.from_counts( + signature_id=self.key, + task_name=self.task_name, + status=self.task_status.status, + total=total, + finished=finished, + failed=failed, + running=running, + pending=pending, + is_done=self.task_status.is_done(), + ) + async def acall(self, msg: Any, set_return_field: bool = True, **kwargs): first_task = await rapyer.afind_one(self.tasks[0]) if first_task is None: diff --git a/libs/third-magic/thirdmagic/container.py b/libs/third-magic/thirdmagic/container.py index 5d4f158b..10af4000 100644 --- a/libs/third-magic/thirdmagic/container.py +++ b/libs/third-magic/thirdmagic/container.py @@ -6,6 +6,7 @@ from rapyer.fields import RapyerKey from thirdmagic.signature import Signature +from thirdmagic.signature.status import ContainerStatus class ContainerTaskSignature(Signature, ABC): @@ -18,6 +19,10 @@ def task_ids(self) -> list[RapyerKey]: async def sub_tasks(self) -> list[Signature]: pass + @abc.abstractmethod + async def container_status(self) -> ContainerStatus: + pass + async def remove_references(self): sub_tasks = await self.sub_tasks() await asyncio.gather( diff --git a/libs/third-magic/thirdmagic/swarm/model.py b/libs/third-magic/thirdmagic/swarm/model.py index e2114cb9..d02a8759 100644 --- a/libs/third-magic/thirdmagic/swarm/model.py +++ b/libs/third-magic/thirdmagic/swarm/model.py @@ -14,7 +14,7 @@ TooManyTasksError, ) from thirdmagic.signature import Signature -from thirdmagic.signature.status import SignatureStatus +from thirdmagic.signature.status import ContainerStatus, SignatureStatus from thirdmagic.swarm.consts import SWARM_MESSAGE_PARAM_NAME from thirdmagic.swarm.state import PublishState from thirdmagic.task.creator import TaskSignatureConvertible, resolve_signatures @@ -187,6 +187,19 @@ async def is_swarm_done(self): finished_all_tasks = set(done_tasks) == set(self.tasks) return self.is_swarm_closed and finished_all_tasks + async def container_status(self) -> ContainerStatus: + return ContainerStatus.from_counts( + signature_id=self.key, + task_name=self.task_name, + status=self.task_status.status, + total=len(self.tasks), + finished=len(self.finished_tasks), + failed=len(self.failed_tasks), + running=self.current_running_tasks, + pending=len(self.tasks_left_to_run), + is_done=await self.is_swarm_done(), + ) + def has_published_callback(self): return self.task_status.status == SignatureStatus.DONE From a3a265f43879b0c450d0738ace47600141aa71df Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Thu, 23 Jul 2026 21:33:48 +0300 Subject: [PATCH 3/9] feat: add mageflow.astatus for batch container status Add astatus(*ids) that loads signatures, verifies each is a container (raising MissingSignatureError / NotAContainerError otherwise) and returns their ContainersStatus. Exposed as HatchetMageflow.astatus and re-exported from the mageflow package. Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- libs/mageflow/mageflow/__init__.py | 5 ++++ .../mageflow/clients/hatchet/mageflow.py | 6 ++++ libs/third-magic/thirdmagic/status.py | 29 +++++++++++++++++++ 3 files changed, 40 insertions(+) create mode 100644 libs/third-magic/thirdmagic/status.py diff --git a/libs/mageflow/mageflow/__init__.py b/libs/mageflow/mageflow/__init__.py index 95c8000a..a7dbe1ef 100644 --- a/libs/mageflow/mageflow/__init__.py +++ b/libs/mageflow/mageflow/__init__.py @@ -3,6 +3,8 @@ from thirdmagic import abounded_field from thirdmagic.chain.creator import chain as achain from thirdmagic.signature import Signature +from thirdmagic.signature.status import ContainersStatus, ContainerStatus +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 @@ -39,6 +41,9 @@ async def load_sign(key: RapyerKey) -> Signature: "SignatureTTLConfig", "achain", "aswarm", + "astatus", + "ContainerStatus", + "ContainersStatus", "start_mageflow", "abounded_field", ] diff --git a/libs/mageflow/mageflow/clients/hatchet/mageflow.py b/libs/mageflow/mageflow/clients/hatchet/mageflow.py index 7f114651..2db91102 100644 --- a/libs/mageflow/mageflow/clients/hatchet/mageflow.py +++ b/libs/mageflow/mageflow/clients/hatchet/mageflow.py @@ -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 @@ -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 diff --git a/libs/third-magic/thirdmagic/status.py b/libs/third-magic/thirdmagic/status.py new file mode 100644 index 00000000..8b190ba6 --- /dev/null +++ b/libs/third-magic/thirdmagic/status.py @@ -0,0 +1,29 @@ +import asyncio + +import rapyer +from rapyer.fields import RapyerKey + +from thirdmagic.container import ContainerTaskSignature +from thirdmagic.errors import MissingSignatureError, NotAContainerError +from thirdmagic.signature.status import ContainersStatus + + +async def _load_container(signature_id: RapyerKey) -> ContainerTaskSignature: + signature = await rapyer.afind_one(signature_id) + if signature is None: + raise MissingSignatureError(f"Signature {signature_id} not found") + if not isinstance(signature, ContainerTaskSignature): + raise NotAContainerError( + f"Signature {signature_id} is not a container signature" + ) + return signature + + +async def astatus(*signature_ids: RapyerKey) -> ContainersStatus: + containers = await asyncio.gather( + *[_load_container(signature_id) for signature_id in signature_ids] + ) + statuses = await asyncio.gather( + *[container.container_status() for container in containers] + ) + return ContainersStatus(containers=list(statuses)) From 5278c9a9f45426a9a8b039e647b2b153ed80a1dc Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Thu, 23 Jul 2026 21:33:58 +0300 Subject: [PATCH 4/9] test: cover container_status and mageflow.astatus Add unit tests for swarm/chain container_status (terminal percentage, done, empty) and for astatus (single, aggregate, and the missing / non-container raise paths). Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- .../tests/unit/workflows/test_astatus.py | 70 +++++++++++++++ .../tests/unit/test_container_status.py | 88 +++++++++++++++++++ 2 files changed, 158 insertions(+) create mode 100644 libs/mageflow/tests/unit/workflows/test_astatus.py create mode 100644 libs/third-magic/tests/unit/test_container_status.py diff --git a/libs/mageflow/tests/unit/workflows/test_astatus.py b/libs/mageflow/tests/unit/workflows/test_astatus.py new file mode 100644 index 00000000..d43d1d63 --- /dev/null +++ b/libs/mageflow/tests/unit/workflows/test_astatus.py @@ -0,0 +1,70 @@ +import pytest +from thirdmagic.errors import MissingSignatureError, NotAContainerError +from thirdmagic.signature.status import ContainersStatus + +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], + ) + + # Act + result = await mageflow.astatus(setup.swarm_task.key) + + # Assert + assert isinstance(result, ContainersStatus) + status = result.containers[0] + assert status.signature_id == setup.swarm_task.key + assert status.total == 4 + assert status.finished == 1 + assert status.failed == 1 + assert status.pending == 1 + assert status.percentage == 50.0 + + +@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) + + # Act + result = await mageflow.astatus(first.swarm_task.key, second.swarm_task.key) + + # Assert + assert len(result.containers) == 2 + # 2 terminal tasks out of 4 total across both swarms + assert result.overall_percentage == 50.0 + + +@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) diff --git a/libs/third-magic/tests/unit/test_container_status.py b/libs/third-magic/tests/unit/test_container_status.py new file mode 100644 index 00000000..27832e20 --- /dev/null +++ b/libs/third-magic/tests/unit/test_container_status.py @@ -0,0 +1,88 @@ +import pytest + +import thirdmagic +from thirdmagic.signature.status import SignatureStatus + + +@pytest.mark.asyncio +async def test_swarm_container_status_terminal_percentage(mock_task_def): + # Arrange + swarm = await thirdmagic.swarm(task_name="test_swarm") + tasks = [await thirdmagic.sign(f"test_task_{i}") for i in range(4)] + await swarm.add_tasks(tasks) + + async with swarm.apipeline(): + swarm.tasks_left_to_run.remove_range(0, len(swarm.tasks_left_to_run)) + swarm.finished_tasks.append(tasks[0].key) + swarm.failed_tasks.append(tasks[1].key) + swarm.current_running_tasks = 1 + swarm.tasks_left_to_run.append(tasks[3].key) + + # Act + status = await swarm.container_status() + + # Assert + assert status.signature_id == swarm.key + assert status.total == 4 + assert status.finished == 1 + assert status.failed == 1 + assert status.running == 1 + assert status.pending == 1 + assert status.percentage == 50.0 + assert status.is_done is False + + +@pytest.mark.asyncio +async def test_swarm_container_status_done_is_full(mock_task_def): + # Arrange + swarm = await thirdmagic.swarm(task_name="test_swarm") + tasks = [await thirdmagic.sign(f"test_task_{i}") for i in range(2)] + await swarm.add_tasks(tasks) + + async with swarm.apipeline(): + swarm.tasks_left_to_run.remove_range(0, len(swarm.tasks_left_to_run)) + swarm.finished_tasks.extend([task.key for task in tasks]) + swarm.is_swarm_closed = True + + # Act + status = await swarm.container_status() + + # Assert + assert status.percentage == 100.0 + assert status.is_done is True + + +@pytest.mark.asyncio +async def test_swarm_container_status_empty_is_zero(mock_task_def): + # Arrange + swarm = await thirdmagic.swarm(task_name="test_swarm") + + # Act + status = await swarm.container_status() + + # Assert + assert status.total == 0 + assert status.percentage == 0.0 + + +@pytest.mark.asyncio +async def test_chain_container_status_classifies_children(mock_task_def): + # Arrange + tasks = [await thirdmagic.sign(f"chain_task_{i}") for i in range(4)] + chain = await thirdmagic.chain([task.key for task in tasks]) + + await tasks[0].change_status(SignatureStatus.DONE) + await tasks[1].change_status(SignatureStatus.FAILED) + await tasks[2].change_status(SignatureStatus.ACTIVE) + # tasks[3] stays PENDING + + # Act + status = await chain.container_status() + + # Assert + assert status.total == 4 + assert status.finished == 1 + assert status.failed == 1 + assert status.running == 1 + assert status.pending == 1 + assert status.percentage == 50.0 From dd75d9ca625aa5e1112faa31b85401e61a782d7a Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Fri, 24 Jul 2026 12:37:25 +0300 Subject: [PATCH 5/9] refactor: load astatus signatures in a single afind Replace the per-id asyncio.gather of afind_one with one batch rapyer.afind, and short-circuit on empty input so an empty id list never triggers a full-DB extraction. Missing keys are detected via a length check and still raise MissingSignatureError. Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- .../tests/unit/workflows/test_astatus.py | 10 ++++++ libs/third-magic/thirdmagic/status.py | 33 ++++++++----------- 2 files changed, 24 insertions(+), 19 deletions(-) diff --git a/libs/mageflow/tests/unit/workflows/test_astatus.py b/libs/mageflow/tests/unit/workflows/test_astatus.py index d43d1d63..41f55b41 100644 --- a/libs/mageflow/tests/unit/workflows/test_astatus.py +++ b/libs/mageflow/tests/unit/workflows/test_astatus.py @@ -50,6 +50,16 @@ async def test_astatus_aggregates_multiple_containers(mock_adapter): 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 diff --git a/libs/third-magic/thirdmagic/status.py b/libs/third-magic/thirdmagic/status.py index 8b190ba6..92853300 100644 --- a/libs/third-magic/thirdmagic/status.py +++ b/libs/third-magic/thirdmagic/status.py @@ -1,5 +1,3 @@ -import asyncio - import rapyer from rapyer.fields import RapyerKey @@ -8,22 +6,19 @@ from thirdmagic.signature.status import ContainersStatus -async def _load_container(signature_id: RapyerKey) -> ContainerTaskSignature: - signature = await rapyer.afind_one(signature_id) - if signature is None: - raise MissingSignatureError(f"Signature {signature_id} not found") - if not isinstance(signature, ContainerTaskSignature): - raise NotAContainerError( - f"Signature {signature_id} is not a container signature" - ) - return signature +async def astatus(*signature_ids: RapyerKey) -> ContainersStatus: + if not signature_ids: + return ContainersStatus(containers=[]) + signatures = await rapyer.afind(*signature_ids, skip_missing=True) + if len(signatures) != len(signature_ids): + raise MissingSignatureError(f"Some signatures were not found: {signature_ids}") -async def astatus(*signature_ids: RapyerKey) -> ContainersStatus: - containers = await asyncio.gather( - *[_load_container(signature_id) for signature_id in signature_ids] - ) - statuses = await asyncio.gather( - *[container.container_status() for container in containers] - ) - return ContainersStatus(containers=list(statuses)) + statuses = [] + for signature in signatures: + if not isinstance(signature, ContainerTaskSignature): + raise NotAContainerError( + f"Signature {signature.key} is not a container signature" + ) + statuses.append(await signature.container_status()) + return ContainersStatus(containers=statuses) From 019941f510521f8b961fa869ac2bc596bec88b49 Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Fri, 24 Jul 2026 13:15:14 +0300 Subject: [PATCH 6/9] refactor: rename container method container_status to astatus The per-container progress method is now astatus() on ContainerTaskSignature (and swarm/chain), matching the client-facing mageflow.astatus naming. Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- ...ainer_status.py => test_container_astatus.py} | 16 ++++++++-------- libs/third-magic/thirdmagic/chain/model.py | 2 +- libs/third-magic/thirdmagic/container.py | 2 +- libs/third-magic/thirdmagic/status.py | 2 +- libs/third-magic/thirdmagic/swarm/model.py | 2 +- 5 files changed, 12 insertions(+), 12 deletions(-) rename libs/third-magic/tests/unit/{test_container_status.py => test_container_astatus.py} (82%) diff --git a/libs/third-magic/tests/unit/test_container_status.py b/libs/third-magic/tests/unit/test_container_astatus.py similarity index 82% rename from libs/third-magic/tests/unit/test_container_status.py rename to libs/third-magic/tests/unit/test_container_astatus.py index 27832e20..ef10e976 100644 --- a/libs/third-magic/tests/unit/test_container_status.py +++ b/libs/third-magic/tests/unit/test_container_astatus.py @@ -5,7 +5,7 @@ @pytest.mark.asyncio -async def test_swarm_container_status_terminal_percentage(mock_task_def): +async def test_swarm_astatus_terminal_percentage(mock_task_def): # Arrange swarm = await thirdmagic.swarm(task_name="test_swarm") tasks = [await thirdmagic.sign(f"test_task_{i}") for i in range(4)] @@ -19,7 +19,7 @@ async def test_swarm_container_status_terminal_percentage(mock_task_def): swarm.tasks_left_to_run.append(tasks[3].key) # Act - status = await swarm.container_status() + status = await swarm.astatus() # Assert assert status.signature_id == swarm.key @@ -33,7 +33,7 @@ async def test_swarm_container_status_terminal_percentage(mock_task_def): @pytest.mark.asyncio -async def test_swarm_container_status_done_is_full(mock_task_def): +async def test_swarm_astatus_done_is_full(mock_task_def): # Arrange swarm = await thirdmagic.swarm(task_name="test_swarm") tasks = [await thirdmagic.sign(f"test_task_{i}") for i in range(2)] @@ -45,7 +45,7 @@ async def test_swarm_container_status_done_is_full(mock_task_def): swarm.is_swarm_closed = True # Act - status = await swarm.container_status() + status = await swarm.astatus() # Assert assert status.percentage == 100.0 @@ -53,12 +53,12 @@ async def test_swarm_container_status_done_is_full(mock_task_def): @pytest.mark.asyncio -async def test_swarm_container_status_empty_is_zero(mock_task_def): +async def test_swarm_astatus_empty_is_zero(mock_task_def): # Arrange swarm = await thirdmagic.swarm(task_name="test_swarm") # Act - status = await swarm.container_status() + status = await swarm.astatus() # Assert assert status.total == 0 @@ -66,7 +66,7 @@ async def test_swarm_container_status_empty_is_zero(mock_task_def): @pytest.mark.asyncio -async def test_chain_container_status_classifies_children(mock_task_def): +async def test_chain_astatus_classifies_children(mock_task_def): # Arrange tasks = [await thirdmagic.sign(f"chain_task_{i}") for i in range(4)] chain = await thirdmagic.chain([task.key for task in tasks]) @@ -77,7 +77,7 @@ async def test_chain_container_status_classifies_children(mock_task_def): # tasks[3] stays PENDING # Act - status = await chain.container_status() + status = await chain.astatus() # Assert assert status.total == 4 diff --git a/libs/third-magic/thirdmagic/chain/model.py b/libs/third-magic/thirdmagic/chain/model.py index 0ff24689..3ac387b9 100644 --- a/libs/third-magic/thirdmagic/chain/model.py +++ b/libs/third-magic/thirdmagic/chain/model.py @@ -47,7 +47,7 @@ async def sub_tasks(self) -> list[TaskSignature]: sub_tasks = await rapyer.afind(*self.tasks, skip_missing=True) return cast(list[TaskSignature], sub_tasks) - async def container_status(self) -> ContainerStatus: + async def astatus(self) -> ContainerStatus: sub_tasks = await self.sub_tasks() finished = failed = running = 0 for task in sub_tasks: diff --git a/libs/third-magic/thirdmagic/container.py b/libs/third-magic/thirdmagic/container.py index 10af4000..0af80352 100644 --- a/libs/third-magic/thirdmagic/container.py +++ b/libs/third-magic/thirdmagic/container.py @@ -20,7 +20,7 @@ async def sub_tasks(self) -> list[Signature]: pass @abc.abstractmethod - async def container_status(self) -> ContainerStatus: + async def astatus(self) -> ContainerStatus: pass async def remove_references(self): diff --git a/libs/third-magic/thirdmagic/status.py b/libs/third-magic/thirdmagic/status.py index 92853300..a9e73e97 100644 --- a/libs/third-magic/thirdmagic/status.py +++ b/libs/third-magic/thirdmagic/status.py @@ -20,5 +20,5 @@ async def astatus(*signature_ids: RapyerKey) -> ContainersStatus: raise NotAContainerError( f"Signature {signature.key} is not a container signature" ) - statuses.append(await signature.container_status()) + statuses.append(await signature.astatus()) return ContainersStatus(containers=statuses) diff --git a/libs/third-magic/thirdmagic/swarm/model.py b/libs/third-magic/thirdmagic/swarm/model.py index d02a8759..6819f0f8 100644 --- a/libs/third-magic/thirdmagic/swarm/model.py +++ b/libs/third-magic/thirdmagic/swarm/model.py @@ -187,7 +187,7 @@ async def is_swarm_done(self): finished_all_tasks = set(done_tasks) == set(self.tasks) return self.is_swarm_closed and finished_all_tasks - async def container_status(self) -> ContainerStatus: + async def astatus(self) -> ContainerStatus: return ContainerStatus.from_counts( signature_id=self.key, task_name=self.task_name, From 2979243fcec2bb0490a6d5ad1b801e4ce3889ef4 Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Fri, 24 Jul 2026 13:15:14 +0300 Subject: [PATCH 7/9] docs: document mageflow.astatus and container astatus Add the Status section to the functions API reference, the client method, and the swarm/chain container astatus() method, and record the feature in the changelog (#134). Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- CHANGELOG.md | 1 + docs/api/chain.md | 8 +++++++ docs/api/client.md | 9 ++++++++ docs/api/functions.md | 54 +++++++++++++++++++++++++++++++++++++++++++ docs/api/swarm.md | 8 +++++++ 5 files changed, 80 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6d0241f8..ddb8a01e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/docs/api/chain.md b/docs/api/chain.md index 5d6c9bc0..dccf4928 100644 --- a/docs/api/chain.md +++ b/docs/api/chain.md @@ -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 +``` diff --git a/docs/api/client.md b/docs/api/client.md index b61ed1d5..94937195 100644 --- a/docs/api/client.md +++ b/docs/api/client.md @@ -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. diff --git a/docs/api/functions.md b/docs/api/functions.md index 64a35da8..cfb41a76 100644 --- a/docs/api/functions.md +++ b/docs/api/functions.md @@ -130,6 +130,60 @@ 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. + +#### `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)` diff --git a/docs/api/swarm.md b/docs/api/swarm.md index b40ebc45..8de82821 100644 --- a/docs/api/swarm.md +++ b/docs/api/swarm.md @@ -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 From 889ae5c8fe061e017b0c0b0d0f49f6fd2abb6d79 Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Fri, 24 Jul 2026 15:52:45 +0300 Subject: [PATCH 8/9] refactor: export status types from thirdmagic, not mageflow ContainerStatus and ContainersStatus are now imported from thirdmagic (top-level and thirdmagic.signature). The mageflow package only re-exports the astatus function, not the result types. Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- docs/api/functions.md | 6 ++++++ libs/mageflow/mageflow/__init__.py | 3 --- libs/third-magic/thirdmagic/__init__.py | 10 +++++++++- 3 files changed, 15 insertions(+), 4 deletions(-) diff --git a/docs/api/functions.md b/docs/api/functions.md index cfb41a76..59d81d60 100644 --- a/docs/api/functions.md +++ b/docs/api/functions.md @@ -160,6 +160,12 @@ 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`. diff --git a/libs/mageflow/mageflow/__init__.py b/libs/mageflow/mageflow/__init__.py index a7dbe1ef..88bf4f81 100644 --- a/libs/mageflow/mageflow/__init__.py +++ b/libs/mageflow/mageflow/__init__.py @@ -3,7 +3,6 @@ from thirdmagic import abounded_field from thirdmagic.chain.creator import chain as achain from thirdmagic.signature import Signature -from thirdmagic.signature.status import ContainersStatus, ContainerStatus from thirdmagic.status import astatus from thirdmagic.swarm.creator import swarm as aswarm from thirdmagic.task import TaskSignature @@ -42,8 +41,6 @@ async def load_sign(key: RapyerKey) -> Signature: "achain", "aswarm", "astatus", - "ContainerStatus", - "ContainersStatus", "start_mageflow", "abounded_field", ] diff --git a/libs/third-magic/thirdmagic/__init__.py b/libs/third-magic/thirdmagic/__init__.py index 5daa2b19..3d9eca64 100644 --- a/libs/third-magic/thirdmagic/__init__.py +++ b/libs/third-magic/thirdmagic/__init__.py @@ -1,9 +1,17 @@ import rapyer from thirdmagic.chain.creator import chain +from thirdmagic.signature.status import ContainersStatus, ContainerStatus from thirdmagic.swarm.creator import swarm from thirdmagic.task.creator import sign abounded_field = rapyer.apipeline -__all__ = ["sign", "chain", "swarm", "abounded_field"] +__all__ = [ + "sign", + "chain", + "swarm", + "abounded_field", + "ContainerStatus", + "ContainersStatus", +] From 98c63d0adffec1a7746867856e8694ef61b39edc Mon Sep 17 00:00:00 2001 From: YedidyaHKfir Date: Fri, 24 Jul 2026 16:05:10 +0300 Subject: [PATCH 9/9] test: assert astatus against a full expected model Build the entire expected ContainerStatus / ContainersStatus in the arrange step and compare with equality, instead of asserting individual fields. Refs #134 Co-Authored-By: Claude Opus 4.8 (1M context) --- .../tests/unit/workflows/test_astatus.py | 58 ++++++++++++--- .../tests/unit/test_container_astatus.py | 74 ++++++++++++++----- 2 files changed, 104 insertions(+), 28 deletions(-) diff --git a/libs/mageflow/tests/unit/workflows/test_astatus.py b/libs/mageflow/tests/unit/workflows/test_astatus.py index 41f55b41..0b286738 100644 --- a/libs/mageflow/tests/unit/workflows/test_astatus.py +++ b/libs/mageflow/tests/unit/workflows/test_astatus.py @@ -1,6 +1,7 @@ import pytest +from thirdmagic import ContainersStatus, ContainerStatus from thirdmagic.errors import MissingSignatureError, NotAContainerError -from thirdmagic.signature.status import ContainersStatus +from thirdmagic.signature.status import SignatureStatus import mageflow from tests.integration.hatchet.models import ContextMessage @@ -18,19 +19,28 @@ async def test_astatus_returns_status_for_single_container(mock_adapter): 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 isinstance(result, ContainersStatus) - status = result.containers[0] - assert status.signature_id == setup.swarm_task.key - assert status.total == 4 - assert status.finished == 1 - assert status.failed == 1 - assert status.pending == 1 - assert status.percentage == 50.0 + assert result == expected @pytest.mark.asyncio @@ -40,12 +50,40 @@ async def test_astatus_aggregates_multiple_containers(mock_adapter): 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 len(result.containers) == 2 + assert result == expected # 2 terminal tasks out of 4 total across both swarms assert result.overall_percentage == 50.0 diff --git a/libs/third-magic/tests/unit/test_container_astatus.py b/libs/third-magic/tests/unit/test_container_astatus.py index ef10e976..526bea0f 100644 --- a/libs/third-magic/tests/unit/test_container_astatus.py +++ b/libs/third-magic/tests/unit/test_container_astatus.py @@ -1,6 +1,7 @@ import pytest import thirdmagic +from thirdmagic import ContainerStatus from thirdmagic.signature.status import SignatureStatus @@ -18,18 +19,24 @@ async def test_swarm_astatus_terminal_percentage(mock_task_def): swarm.current_running_tasks = 1 swarm.tasks_left_to_run.append(tasks[3].key) + expected = ContainerStatus( + signature_id=swarm.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 status = await swarm.astatus() # Assert - assert status.signature_id == swarm.key - assert status.total == 4 - assert status.finished == 1 - assert status.failed == 1 - assert status.running == 1 - assert status.pending == 1 - assert status.percentage == 50.0 - assert status.is_done is False + assert status == expected @pytest.mark.asyncio @@ -44,25 +51,48 @@ async def test_swarm_astatus_done_is_full(mock_task_def): swarm.finished_tasks.extend([task.key for task in tasks]) swarm.is_swarm_closed = True + expected = ContainerStatus( + signature_id=swarm.key, + task_name="test_swarm", + status=SignatureStatus.PENDING, + total=2, + finished=2, + failed=0, + running=0, + pending=0, + percentage=100.0, + is_done=True, + ) + # Act status = await swarm.astatus() # Assert - assert status.percentage == 100.0 - assert status.is_done is True + assert status == expected @pytest.mark.asyncio async def test_swarm_astatus_empty_is_zero(mock_task_def): # Arrange swarm = await thirdmagic.swarm(task_name="test_swarm") + expected = ContainerStatus( + signature_id=swarm.key, + task_name="test_swarm", + status=SignatureStatus.PENDING, + total=0, + finished=0, + failed=0, + running=0, + pending=0, + percentage=0.0, + is_done=False, + ) # Act status = await swarm.astatus() # Assert - assert status.total == 0 - assert status.percentage == 0.0 + assert status == expected @pytest.mark.asyncio @@ -76,13 +106,21 @@ async def test_chain_astatus_classifies_children(mock_task_def): await tasks[2].change_status(SignatureStatus.ACTIVE) # tasks[3] stays PENDING + expected = ContainerStatus( + signature_id=chain.key, + task_name=chain.task_name, + status=SignatureStatus.PENDING, + total=4, + finished=1, + failed=1, + running=1, + pending=1, + percentage=50.0, + is_done=False, + ) + # Act status = await chain.astatus() # Assert - assert status.total == 4 - assert status.finished == 1 - assert status.failed == 1 - assert status.running == 1 - assert status.pending == 1 - assert status.percentage == 50.0 + assert status == expected