From 9715a57eb4a5eb36f94ec8a9931c73ce339bb45b Mon Sep 17 00:00:00 2001 From: Asuka Minato Date: Thu, 13 Aug 2026 16:11:59 +0900 Subject: [PATCH] test: run workflow archiver coverage on SQLite --- .../test_workflow_run_archiver_sqlite.py} | 223 +++++++++--------- 1 file changed, 107 insertions(+), 116 deletions(-) rename api/tests/{integration_tests/services/retention/test_workflow_run_archiver.py => unit_tests/services/retention/workflow_run/test_workflow_run_archiver_sqlite.py} (82%) diff --git a/api/tests/integration_tests/services/retention/test_workflow_run_archiver.py b/api/tests/unit_tests/services/retention/workflow_run/test_workflow_run_archiver_sqlite.py similarity index 82% rename from api/tests/integration_tests/services/retention/test_workflow_run_archiver.py rename to api/tests/unit_tests/services/retention/workflow_run/test_workflow_run_archiver_sqlite.py index 098069776f2..bc512466aac 100644 --- a/api/tests/integration_tests/services/retention/test_workflow_run_archiver.py +++ b/api/tests/unit_tests/services/retention/workflow_run/test_workflow_run_archiver_sqlite.py @@ -1,18 +1,27 @@ +"""SQLite and in-process coverage for workflow-run bundle archiving.""" + import datetime import json import uuid +from types import SimpleNamespace from unittest.mock import ANY, MagicMock, patch import pyarrow as pa import pyarrow.parquet as pq import pytest +from sqlalchemy import select +from sqlalchemy.engine import Engine from sqlalchemy.exc import OperationalError +from sqlalchemy.orm import Session, sessionmaker from enums import DeploymentEdition -from models.workflow import WorkflowRunArchiveBundle +from graphon.enums import WorkflowExecutionStatus +from models.enums import CreatorUserRole, WorkflowRunTriggeredFrom +from models.workflow import WorkflowRun, WorkflowRunArchiveBundle, WorkflowType from services.retention.workflow_run.archive_paid_plan_workflow_run import ( ArchiveResult, ArchiveSummary, + TableStats, WorkflowRunArchiver, ) from services.retention.workflow_run.constants import ARCHIVE_BUNDLE_FORMAT, ARCHIVE_BUNDLE_SCHEMA_VERSION @@ -45,19 +54,29 @@ def _db_disconnect_error() -> OperationalError: ) -def _run(run_id: str = "run-1"): - run = MagicMock() - run.id = run_id - run.tenant_id = "tenant-1" - run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) - return run - - -def _session_context(session): - context = MagicMock() - context.__enter__.return_value = session - context.__exit__.return_value = False - return context +def _run(run_id: str = "run-1") -> WorkflowRun: + """Build the mapped entity consumed by the archiver without persisting it.""" + return WorkflowRun( + id=run_id, + tenant_id="tenant-1", + app_id="app-1", + workflow_id="workflow-1", + type=WorkflowType.WORKFLOW, + triggered_from=WorkflowRunTriggeredFrom.APP_RUN, + version="1", + graph="{}", + inputs="{}", + status=WorkflowExecutionStatus.SUCCEEDED, + outputs="{}", + error=None, + elapsed_time=0, + total_tokens=0, + total_steps=0, + created_by_role=CreatorUserRole.ACCOUNT, + created_by="account-1", + created_at=datetime.datetime(2025, 3, 15, 10, 0, 0), + finished_at=datetime.datetime(2025, 3, 15, 10, 1, 0), + ) class TestWorkflowRunArchiverInit: @@ -214,10 +233,8 @@ class TestWorkflowRunArchiverInit: class TestBuildArchiveBundle: def test_bundle_contains_manifest_and_all_table_objects(self): archiver = WorkflowRunArchiver(days=90) - run = MagicMock() - run.id = str(uuid.uuid4()) + run = _run(str(uuid.uuid4())) run.tenant_id = str(uuid.uuid4()) - run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) identity = archiver._build_bundle_identity([run]) table_data = {"workflow_runs": [{"id": run.id, "tenant_id": run.tenant_id}]} @@ -237,12 +254,8 @@ class TestGenerateManifest: start = datetime.datetime(2025, 1, 1, tzinfo=datetime.UTC) end = datetime.datetime(2025, 4, 1, tzinfo=datetime.UTC) archiver = WorkflowRunArchiver(start_from=start, end_before=end, run_shard_index=1, run_shard_total=4) - from services.retention.workflow_run.archive_paid_plan_workflow_run import TableStats - - run = MagicMock() - run.id = str(uuid.uuid4()) + run = _run(str(uuid.uuid4())) run.tenant_id = str(uuid.uuid4()) - run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) identity = archiver._build_bundle_identity([run]) stats = [ @@ -351,31 +364,32 @@ class TestFilterPaidTenants: class TestDryRunArchive: @patch("services.retention.workflow_run.archive_paid_plan_workflow_run.get_archive_storage") - def test_dry_run_does_not_call_storage(self, mock_get_storage, flask_req_ctx): + def test_dry_run_does_not_call_storage(self, mock_get_storage, sqlite_engine: Engine): archiver = WorkflowRunArchiver(days=90, dry_run=True) - with patch.object(archiver, "_get_runs_batch", return_value=[]): + with ( + patch( + "services.retention.workflow_run.archive_paid_plan_workflow_run.db", + SimpleNamespace(engine=sqlite_engine), + ), + patch.object(archiver, "_get_runs_batch", return_value=[]), + ): summary = archiver.run() mock_get_storage.assert_not_called() assert isinstance(summary, ArchiveSummary) assert summary.runs_failed == 0 - def test_dry_run_estimates_table_and_object_sizes(self): + def test_dry_run_estimates_table_and_object_sizes(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90, dry_run=True) - run = MagicMock() - run.id = "run-1" - run.tenant_id = "tenant-1" - run.app_id = "app-1" - run.workflow_id = "workflow-1" - run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) + run = _run() table_data = { "workflow_runs": [{"id": "run-1", "tenant_id": "tenant-1"}], "workflow_app_logs": [{"id": "log-1", "workflow_run_id": "run-1"}], } with patch.object(archiver, "_extract_bundle_data", return_value=table_data): - result = archiver._archive_bundle(MagicMock(), None, [run]) + result = archiver._archive_bundle(sqlite_session, None, [run]) stats_by_table = {stat.table_name: stat for stat in result.tables} assert result.success is True @@ -389,12 +403,18 @@ class TestDryRunArchive: def test_summary_merges_dry_run_estimates(self): summary = ArchiveSummary() - result = MagicMock() - result.object_size_bytes = 128 - result.tables = [ - MagicMock(table_name="workflow_runs", row_count=1, size_bytes=64), - MagicMock(table_name="workflow_app_logs", row_count=2, size_bytes=32), - ] + result = ArchiveResult( + bundle_id="bundle-1", + tenant_id="tenant-1", + object_prefix="prefix", + run_count=1, + success=True, + object_size_bytes=128, + tables=[ + TableStats(table_name="workflow_runs", row_count=1, checksum="", size_bytes=64), + TableStats(table_name="workflow_app_logs", row_count=2, checksum="", size_bytes=32), + ], + ) WorkflowRunArchiver._merge_result_stats(summary, result) @@ -406,15 +426,10 @@ class TestDryRunArchive: class TestArchiveDbRetry: - def test_archive_bundle_groups_retries_with_fresh_session(self): + def test_archive_bundle_groups_retries_with_fresh_session(self, sqlite_session_factory: sessionmaker[Session]): archiver = WorkflowRunArchiver(days=90) run = _run() - session_maker = MagicMock( - side_effect=[ - _session_context(MagicMock(name="session-1")), - _session_context(MagicMock(name="session-2")), - ] - ) + session_maker = MagicMock(wraps=sqlite_session_factory) success = ArchiveResult( bundle_id=archiver._build_bundle_identity([run]).bundle_id, tenant_id=run.tenant_id, @@ -434,16 +449,12 @@ class TestArchiveDbRetry: assert session_maker.call_count == 2 sleep.assert_called_once_with(1.0) - def test_archive_bundle_groups_returns_failed_result_after_retry_exhaustion(self): + def test_archive_bundle_groups_returns_failed_result_after_retry_exhaustion( + self, sqlite_session_factory: sessionmaker[Session] + ): archiver = WorkflowRunArchiver(days=90) run = _run() - session_maker = MagicMock( - side_effect=[ - _session_context(MagicMock(name="session-1")), - _session_context(MagicMock(name="session-2")), - _session_context(MagicMock(name="session-3")), - ] - ) + session_maker = MagicMock(wraps=sqlite_session_factory) with ( patch.object(archiver, "_archive_bundle", side_effect=[_db_disconnect_error()] * 3) as archive_bundle, @@ -458,17 +469,18 @@ class TestArchiveDbRetry: assert session_maker.call_count == archiver.DB_RETRY_ATTEMPTS assert sleep.call_count == archiver.DB_RETRY_ATTEMPTS - 1 - def test_archive_bundle_uses_safe_rollback_when_failure_rolls_back_badly(self): + def test_archive_bundle_uses_safe_rollback_when_failure_rolls_back_badly(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90, dry_run=True) - session = MagicMock() - session.rollback.side_effect = RuntimeError("rollback failed") - with patch.object(archiver, "_extract_bundle_data", side_effect=RuntimeError("extract failed")): - result = archiver._archive_bundle(session, None, [_run()]) + with ( + patch.object(sqlite_session, "rollback", side_effect=RuntimeError("rollback failed")) as rollback, + patch.object(archiver, "_extract_bundle_data", side_effect=RuntimeError("extract failed")), + ): + result = archiver._archive_bundle(sqlite_session, None, [_run()]) assert result.success is False assert result.error == "extract failed" - session.rollback.assert_called_once() + rollback.assert_called_once() class TestArchiveRunIdempotency: @@ -487,56 +499,52 @@ class TestArchiveRunIdempotency: ).encode() return index_key, payload - def test_locked_bundle_is_skipped(self): + def test_locked_bundle_is_skipped(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) - run = MagicMock() - run.id = "run-1" - run.tenant_id = "tenant-1" - run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) + run = _run() with ( patch.object(archiver, "_lock_runs_for_archive", return_value=[]), ): storage = MagicMock() storage.object_exists.return_value = False - result = archiver._archive_bundle(MagicMock(), storage, [run]) + result = archiver._archive_bundle(sqlite_session, storage, [run]) assert result.success is True assert result.skipped is True assert result.error == "one or more runs locked or deleted by another archiver" - def test_already_archived_bundle_is_skipped(self): + def test_already_archived_bundle_is_skipped(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) - run = MagicMock() - run.id = "run-1" - run.tenant_id = "tenant-1" - run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) + run = _run() storage = MagicMock() storage.object_exists.return_value = True with patch.object(archiver, "_sync_existing_bundle_index") as sync_existing_bundle_index: - result = archiver._archive_bundle(MagicMock(), storage, [run]) + result = archiver._archive_bundle(sqlite_session, storage, [run]) assert result.success is True assert result.skipped is True assert result.error == "bundle already archived" sync_existing_bundle_index.assert_called_once() - def test_existing_bundle_catalog_publication_failure_is_not_success(self): + def test_existing_bundle_catalog_publication_failure_is_not_success(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) run = _run() - session = MagicMock() storage = MagicMock() storage.object_exists.return_value = True - with patch.object(archiver, "_sync_existing_bundle_index", side_effect=RuntimeError("catalog unavailable")): - result = archiver._archive_bundle(session, storage, [run]) + with ( + patch.object(sqlite_session, "rollback", wraps=sqlite_session.rollback) as rollback, + patch.object(archiver, "_sync_existing_bundle_index", side_effect=RuntimeError("catalog unavailable")), + ): + result = archiver._archive_bundle(sqlite_session, storage, [run]) assert result.success is False assert result.error == "catalog unavailable" - session.rollback.assert_called_once() + rollback.assert_called_once() - def test_retry_repairs_index_after_catalog_commit_then_index_write_failure(self): + def test_retry_repairs_index_after_catalog_commit_then_index_write_failure(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) run = _run() identity = archiver._build_bundle_identity([run]) @@ -555,15 +563,13 @@ class TestArchiveRunIdempotency: return original_put_object(key, data) storage.put_object = MagicMock(side_effect=put_object) - first_session = MagicMock() - first_session.scalar.return_value = None table_data = {"workflow_runs": [{"id": run.id, "tenant_id": run.tenant_id}]} with ( patch.object(archiver, "_lock_runs_for_archive", return_value=[run]), patch.object(archiver, "_extract_bundle_data", return_value=table_data), ): - first_result = archiver._archive_bundle(first_session, storage, [run]) + first_result = archiver._archive_bundle(sqlite_session, storage, [run]) assert first_result.success is False assert first_result.error == "index write failed" @@ -571,10 +577,7 @@ class TestArchiveRunIdempotency: assert json.loads(storage.objects[index_key])["run_ids"] == [] storage.list_objects = MagicMock(wraps=storage.list_objects) - retry_session = MagicMock() - retry_session.scalar.return_value = None - - retry_result = archiver._archive_bundle(retry_session, storage, [run]) + retry_result = archiver._archive_bundle(sqlite_session, storage, [run]) assert retry_result.success is True assert retry_result.skipped is True @@ -582,7 +585,7 @@ class TestArchiveRunIdempotency: assert json.loads(storage.objects[index_key])["run_ids"] == [run.id] storage.list_objects.assert_not_called() - def test_existing_manifest_with_missing_index_fails_without_partial_rebuild(self): + def test_existing_manifest_with_missing_index_fails_without_partial_rebuild(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) run = _run() identity = archiver._build_bundle_identity([run]) @@ -596,21 +599,17 @@ class TestArchiveRunIdempotency: storage = FakeArchiveStorage({manifest_key: manifest_data}) storage.list_objects = MagicMock(wraps=storage.list_objects) - result = archiver._archive_bundle(MagicMock(), storage, [run]) + result = archiver._archive_bundle(sqlite_session, storage, [run]) assert result.success is False assert "archive shard index missing" in (result.error or "") assert index_key not in storage.objects storage.list_objects.assert_not_called() - def test_successful_bundle_persists_archive_index(self): + def test_successful_bundle_persists_archive_index(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) - run = MagicMock() - run.id = str(uuid.uuid4()) + run = _run(str(uuid.uuid4())) run.tenant_id = str(uuid.uuid4()) - run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) - session = MagicMock() - session.scalar.return_value = None storage = MagicMock() storage.object_exists.return_value = False table_data = { @@ -622,9 +621,11 @@ class TestArchiveRunIdempotency: patch.object(archiver, "_lock_runs_for_archive", return_value=[run]), patch.object(archiver, "_extract_bundle_data", return_value=table_data), ): - result = archiver._archive_bundle(session, storage, [run]) + result = archiver._archive_bundle(sqlite_session, storage, [run]) - archived_bundle = session.add.call_args.args[0] + archived_bundle = sqlite_session.scalar( + select(WorkflowRunArchiveBundle).where(WorkflowRunArchiveBundle.tenant_id == run.tenant_id) + ) assert result.success is True assert isinstance(archived_bundle, WorkflowRunArchiveBundle) assert archived_bundle.tenant_id == run.tenant_id @@ -632,40 +633,35 @@ class TestArchiveRunIdempotency: assert archived_bundle.month == 3 assert archived_bundle.workflow_run_count == 1 assert archived_bundle.row_count == 2 - session.commit.assert_called_once() - def test_new_bundle_catalog_commit_failure_is_not_success(self): + def test_new_bundle_catalog_commit_failure_is_not_success(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) run = _run(str(uuid.uuid4())) run.tenant_id = str(uuid.uuid4()) - session = MagicMock() - session.scalar.return_value = None - session.commit.side_effect = RuntimeError("catalog commit failed") storage = MagicMock() storage.object_exists.return_value = False storage.list_objects.return_value = [] table_data = {"workflow_runs": [{"id": run.id, "tenant_id": run.tenant_id}]} with ( + patch.object(sqlite_session, "commit", side_effect=RuntimeError("catalog commit failed")), + patch.object(sqlite_session, "rollback", wraps=sqlite_session.rollback) as rollback, patch.object(archiver, "_lock_runs_for_archive", return_value=[run]), patch.object(archiver, "_extract_bundle_data", return_value=table_data), ): - result = archiver._archive_bundle(session, storage, [run]) + result = archiver._archive_bundle(sqlite_session, storage, [run]) assert result.success is False assert result.error == "catalog commit failed" - session.rollback.assert_called_once() + rollback.assert_called_once() - def test_index_skips_all_already_archived_runs(self): + def test_index_skips_all_already_archived_runs(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) - run = MagicMock() - run.id = "run-1" - run.tenant_id = "tenant-1" - run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) + run = _run() index_key, index_payload = self._index_payload(archiver, ["run-1"], run) storage = FakeArchiveStorage({index_key: index_payload}) - result = archiver._archive_bundle(MagicMock(), storage, [run]) + result = archiver._archive_bundle(sqlite_session, storage, [run]) assert result.success is True assert result.skipped is True @@ -673,15 +669,10 @@ class TestArchiveRunIdempotency: assert result.skipped_run_count == 1 assert result.error == "all runs already archived in shard index" - def test_index_filters_duplicate_runs_before_archive(self): + def test_index_filters_duplicate_runs_before_archive(self, sqlite_session: Session): archiver = WorkflowRunArchiver(days=90) - archived_run = MagicMock() - archived_run.id = "run-1" - archived_run.tenant_id = "tenant-1" - archived_run.created_at = datetime.datetime(2025, 3, 15, 10, 0, 0) - new_run = MagicMock() - new_run.id = "run-2" - new_run.tenant_id = "tenant-1" + archived_run = _run() + new_run = _run("run-2") new_run.created_at = datetime.datetime(2025, 3, 15, 11, 0, 0) index_key, index_payload = self._index_payload(archiver, ["run-1"], archived_run) storage = FakeArchiveStorage({index_key: index_payload}) @@ -690,7 +681,7 @@ class TestArchiveRunIdempotency: patch.object(archiver, "_lock_runs_for_archive", return_value=[new_run]) as lock_runs, patch.object(archiver, "_extract_bundle_data", return_value={"workflow_runs": [{"id": "run-2"}]}), ): - result = archiver._archive_bundle(MagicMock(), storage, [archived_run, new_run]) + result = archiver._archive_bundle(sqlite_session, storage, [archived_run, new_run]) assert result.success is True assert result.skipped is False