From abfda37fec40eba7acedaccf8ed19ed62bc65ebf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Cle=CC=81ment=20Doumouro?= Date: Thu, 27 Aug 2026 18:55:00 +0200 Subject: [PATCH 1/2] fix(workflows-worker): version bump --- workers/workflows-worker/main.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/workers/workflows-worker/main.py b/workers/workflows-worker/main.py index 0c181840..c6582158 100644 --- a/workers/workflows-worker/main.py +++ b/workers/workflows-worker/main.py @@ -57,11 +57,11 @@ def _bump_version(current: Version, *, breaking: bool) -> tuple[Version, BumpTyp patch = release[2] if current < _1_0_0: if breaking: - return Version(f"0.{minor + 1}.{patch}"), BumpType.MINOR + return Version(f"0.{minor + 1}.0"), BumpType.MINOR return Version(f"0.{minor}.{patch + 1}"), BumpType.PATCH if breaking: - return Version(f"{major + 1}.{0}.{0}"), BumpType.MAJOR - return Version(f"{major}.{minor + 1}.{0}"), BumpType.MINOR + return Version(f"{major + 1}.0.0"), BumpType.MAJOR + return Version(f"{major}.{minor + 1}.0"), BumpType.MINOR def _validate_version(current: Version) -> None: From a4d6ef6afbdc4f86ddd9f99c62278a760da0d460 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Cle=CC=81ment=20Doumouro?= Date: Fri, 28 Aug 2026 09:48:50 +0200 Subject: [PATCH 2/2] feature(datashare-python): add task_id to manifest entry --- datashare-python/datashare_python/objects.py | 12 ++++- datashare-python/tests/test_objects.py | 51 +++++++++++++++++++ datashare-python/tests/test_utils.py | 3 ++ datashare-python/uv.lock | 2 +- .../passport-worker/tests/test_workflows.py | 4 ++ 5 files changed, 70 insertions(+), 2 deletions(-) diff --git a/datashare-python/datashare_python/objects.py b/datashare-python/datashare_python/objects.py index b37d6911..6aa7270c 100644 --- a/datashare-python/datashare_python/objects.py +++ b/datashare-python/datashare_python/objects.py @@ -14,7 +14,7 @@ from pydantic_core import PydanticCustomError, ValidationError, core_schema from pydantic_core.core_schema import PlainValidatorFunctionSchema from pydantic_extra_types.language_code import LanguageName -from temporalio import workflow +from temporalio import activity, workflow from .constants import TIKA_METADATA_RESOURCENAME @@ -325,6 +325,8 @@ def as_manifest_task_input(self) -> dict[str, Any]: class ManifestEntry[A](DatashareModel, ABC): status: ManifestEntryStatus + # TODO: make this one non optional in the next major ! + task_id: str | None label: str | None = None input: Annotated[ dict[str, Any] | None, @@ -336,7 +338,11 @@ class ManifestEntry[A](DatashareModel, ABC): @classmethod def complete(cls, args: A, label: str | None = None, **kwargs) -> Self: + task_id = None + if activity.in_activity(): + task_id = activity.info().workflow_id return cls( + task_id=task_id, input=args.as_manifest_task_input(), label=label, status=ManifestEntryStatus.COMPLETE, @@ -345,7 +351,11 @@ def complete(cls, args: A, label: str | None = None, **kwargs) -> Self: @classmethod def partial(cls, args: A, label: str | None = None, **kwargs) -> Self: + task_id = None + if activity.in_activity(): + task_id = activity.info().workflow_id return cls( + task_id=task_id, input=args.as_manifest_task_input(), label=label, status=ManifestEntryStatus.PARTIAL, diff --git a/datashare-python/tests/test_objects.py b/datashare-python/tests/test_objects.py index 110b5b28..18288132 100644 --- a/datashare-python/tests/test_objects.py +++ b/datashare-python/tests/test_objects.py @@ -2,8 +2,10 @@ import re from datetime import datetime from pathlib import Path +from unittest.mock import MagicMock, PropertyMock import pytest +from _pytest.monkeypatch import MonkeyPatch from datashare_python.conftest import TEST_PROJECT from datashare_python.constants import TIKA_METADATA_RESOURCENAME from datashare_python.objects import ( @@ -12,12 +14,21 @@ Document, DocumentLocation, FilesystemPagination, + ManifestEntry, Pages, ProcessedFile, Task, + TaskArgs, TaskState, ) from pydantic import TypeAdapter, ValidationError +from temporalio import activity + + +class MockedManifestEntry(ManifestEntry): ... + + +class MockedArgs(TaskArgs): ... def test_task_ser() -> None: @@ -139,3 +150,43 @@ def test_pages_validation_should_raise_for_inconsistent_byte_ranges() -> None: pagination=ByteRangesPagination(byte_ranges=[(0, 1), (1, 2), (2, 3)]), total=2, ) + + +@pytest.mark.parametrize("in_activity", [True, False]) +def test_manifest_entry_complete_task_id( + *, in_activity: bool, monkeypatch: MonkeyPatch +) -> None: + # Given + args = MockedArgs() + mocked_info = MagicMock() + type(mocked_info).workflow_id = PropertyMock(return_value="some_value") + if in_activity: + monkeypatch.setattr(activity, "in_activity", lambda: True) + monkeypatch.setattr(activity, "info", lambda: mocked_info) + # When + manifest_entry = MockedManifestEntry.complete(args) + # Then + if in_activity: + assert manifest_entry.task_id is not None + else: + assert manifest_entry.task_id is None + + +@pytest.mark.parametrize("in_activity", [True, False]) +def test_manifest_entry_partial_task_id( + *, in_activity: bool, monkeypatch: MonkeyPatch +) -> None: + # Given + args = MockedArgs() + mocked_info = MagicMock() + type(mocked_info).workflow_id = PropertyMock(return_value="some_value") + if in_activity: + monkeypatch.setattr(activity, "in_activity", lambda: True) + monkeypatch.setattr(activity, "info", lambda: mocked_info) + # When + manifest_entry = MockedManifestEntry.partial(args) + # Then + if in_activity: + assert manifest_entry.task_id is not None + else: + assert manifest_entry.task_id is None diff --git a/datashare-python/tests/test_utils.py b/datashare-python/tests/test_utils.py index 6d50afe6..1225379a 100644 --- a/datashare-python/tests/test_utils.py +++ b/datashare-python/tests/test_utils.py @@ -179,6 +179,7 @@ def test_write_artifact(tmp_path: Path) -> None: expected_manifest = { "structure": { "status": "complete", + "taskId": None, "taskInput": {"someValue": "value"}, "label": None, } @@ -219,6 +220,7 @@ def test_write_artifact_with_existing_metadata(tmp_path: Path) -> None: expected_manifest = { "structure": { "status": "complete", + "taskId": None, "taskInput": {"someValue": "value"}, "label": None, }, @@ -303,6 +305,7 @@ def test_overwrite_artifact(tmp_path: Path) -> None: expected_manifest = { "structure": { "status": "complete", + "taskId": None, "taskInput": {"someValue": "value"}, "label": None, }, diff --git a/datashare-python/uv.lock b/datashare-python/uv.lock index 1601b0c5..60f399a4 100644 --- a/datashare-python/uv.lock +++ b/datashare-python/uv.lock @@ -477,7 +477,7 @@ wheels = [ [[package]] name = "datashare-python" -version = "0.10.2" +version = "0.10.4" source = { editable = "." } dependencies = [ { name = "aiofile" }, diff --git a/workers/passport-worker/tests/test_workflows.py b/workers/passport-worker/tests/test_workflows.py index ad0cd307..1e364c1a 100644 --- a/workers/passport-worker/tests/test_workflows.py +++ b/workers/passport-worker/tests/test_workflows.py @@ -6,6 +6,7 @@ from datashare_python.conftest import TEST_PROJECT from datashare_python.objects import ProcessedFile from datashare_python.utils import safe_dir +from icij_common.pydantic_utils import safe_copy from passport_worker.config import PassportWorkerConfig from passport_worker.objects import ( PassportDetectionArgs, @@ -71,6 +72,9 @@ async def test_passport_detection_workflow( # noqa: PLR0917 has_passport = "not_a" not in f.name expected.append((f"e2e-doc-{len(expected)}", has_manifest, has_passport)) expected_manifest_entry = PassportManifestEntry.complete(args) + expected_manifest_entry = safe_copy( + expected_manifest_entry, update={"task_id": wf_id} + ) for doc_id, has_manifest, has_passport in expected: artifacts_path = ( worker_paths.artifacts / TEST_PROJECT / safe_dir(doc_id) / doc_id