feat: consolidate legacy workflows into gyxx-flow
This commit is contained in:
@@ -0,0 +1,242 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from gyxx_flow.core.context import RunContext
|
||||
from gyxx_flow.core.layout import DataLayout
|
||||
from gyxx_flow.core.records import RunJournal
|
||||
from gyxx_flow.ops import IdempotencyConflict, Outbox, RunIndex
|
||||
|
||||
|
||||
def _journal(
|
||||
root: Path,
|
||||
*,
|
||||
workflow_id: str = "shop.weekly",
|
||||
business_date: str = "2026-07-27",
|
||||
suffix: str = "abc123",
|
||||
) -> RunJournal:
|
||||
context = RunContext.create(
|
||||
workflow_id,
|
||||
business_date,
|
||||
now=datetime(2026, 7, 27, 4, 0, tzinfo=timezone.utc),
|
||||
random_suffix=suffix,
|
||||
)
|
||||
return RunJournal.create(DataLayout(root), context)
|
||||
|
||||
|
||||
def test_run_index_snapshots_journals_and_supports_structured_queries(tmp_path: Path) -> None:
|
||||
first = _journal(tmp_path, suffix="abc123")
|
||||
second = _journal(
|
||||
tmp_path,
|
||||
workflow_id="product.daily",
|
||||
business_date="2026-07-26",
|
||||
suffix="def456",
|
||||
)
|
||||
first.finalize("success")
|
||||
second.finalize("failed", error="collector failed")
|
||||
index = RunIndex(tmp_path)
|
||||
|
||||
first_record = index.index_journal(first)
|
||||
second_record = index.index_journal(second.path)
|
||||
|
||||
assert first_record.status == "success"
|
||||
assert second_record.error == "collector failed"
|
||||
assert index.get(first_record.run_id) == first_record
|
||||
assert index.query(status="failed") == (second_record,)
|
||||
assert index.query(workflow_id="shop.weekly") == (first_record,)
|
||||
assert index.query(business_date="2026-07-26") == (second_record,)
|
||||
assert index.query(run_id=first_record.run_id) == (first_record,)
|
||||
assert list((tmp_path / "state" / "ops" / "run-index").glob(".*.tmp")) == []
|
||||
|
||||
|
||||
def test_run_index_refreshes_same_run_after_journal_changes(tmp_path: Path) -> None:
|
||||
journal = _journal(tmp_path)
|
||||
index = RunIndex(tmp_path)
|
||||
running = index.index_journal(journal)
|
||||
|
||||
journal.start_step("collect.attempt-1", attempt=1)
|
||||
journal.finish_step("collect.attempt-1", status="success", exit_code=0)
|
||||
journal.finalize("success")
|
||||
complete = index.index_journal(journal)
|
||||
|
||||
assert running.run_id == complete.run_id
|
||||
assert complete.status == "success"
|
||||
assert complete.step_counts == {"success": 1, "failed": 0, "skipped": 0, "running": 0}
|
||||
assert len(tuple((tmp_path / "state" / "ops" / "run-index").glob("*.json"))) == 1
|
||||
|
||||
|
||||
def test_run_index_rejects_malformed_journal_without_writing_index(tmp_path: Path) -> None:
|
||||
malformed = tmp_path / "malformed.json"
|
||||
malformed.write_text('{"run_id": "unsafe/../id"}', encoding="utf-8")
|
||||
index = RunIndex(tmp_path)
|
||||
|
||||
with pytest.raises(ValueError, match="journal"):
|
||||
index.index_journal(malformed)
|
||||
|
||||
assert index.query() == ()
|
||||
|
||||
|
||||
def test_outbox_enqueue_is_atomic_and_idempotent(tmp_path: Path) -> None:
|
||||
outbox = Outbox(tmp_path)
|
||||
|
||||
first = outbox.enqueue(
|
||||
idempotency_key="run-001:feishu-report",
|
||||
topic="feishu.report",
|
||||
payload={"report_id": "report-001"},
|
||||
run_id="run-001",
|
||||
)
|
||||
duplicate = outbox.enqueue(
|
||||
idempotency_key="run-001:feishu-report",
|
||||
topic="feishu.report",
|
||||
payload={"report_id": "report-001"},
|
||||
run_id="run-001",
|
||||
)
|
||||
|
||||
assert duplicate == first
|
||||
assert first.status == "pending"
|
||||
assert outbox.get(first.message_id) == first
|
||||
assert outbox.list(status="pending") == (first,)
|
||||
assert len(tuple((tmp_path / "state" / "ops" / "outbox" / "messages").glob("*.json"))) == 1
|
||||
assert list((tmp_path / "state" / "ops" / "outbox" / "messages").glob(".*.tmp")) == []
|
||||
|
||||
|
||||
def test_outbox_rejects_reusing_key_for_different_semantics(tmp_path: Path) -> None:
|
||||
outbox = Outbox(tmp_path)
|
||||
outbox.enqueue(
|
||||
idempotency_key="run-001:db-upsert",
|
||||
topic="db.upsert",
|
||||
payload={"row_id": 1},
|
||||
run_id="run-001",
|
||||
)
|
||||
|
||||
with pytest.raises(IdempotencyConflict, match="idempotency key"):
|
||||
outbox.enqueue(
|
||||
idempotency_key="run-001:db-upsert",
|
||||
topic="db.upsert",
|
||||
payload={"row_id": 2},
|
||||
run_id="run-001",
|
||||
)
|
||||
|
||||
|
||||
def test_outbox_idempotency_compares_normalized_json_semantics(tmp_path: Path) -> None:
|
||||
outbox = Outbox(tmp_path)
|
||||
first = outbox.enqueue(
|
||||
idempotency_key="run-001:normalized",
|
||||
topic="sink.publish",
|
||||
payload={"row_ids": [1, 2]},
|
||||
run_id="run-001",
|
||||
)
|
||||
|
||||
duplicate = outbox.enqueue(
|
||||
idempotency_key="run-001:normalized",
|
||||
topic="sink.publish",
|
||||
payload={"row_ids": (1, 2)},
|
||||
run_id="run-001",
|
||||
)
|
||||
|
||||
assert duplicate == first
|
||||
|
||||
|
||||
def test_outbox_state_transitions_and_safe_replay(tmp_path: Path) -> None:
|
||||
outbox = Outbox(tmp_path)
|
||||
pending = outbox.enqueue(
|
||||
idempotency_key="run-001:notify",
|
||||
topic="notification.official",
|
||||
payload={"event": "complete"},
|
||||
run_id="run-001",
|
||||
)
|
||||
|
||||
failed = outbox.mark_failed(pending.message_id, error="temporary outage")
|
||||
replayed = outbox.replay(failed.message_id)
|
||||
sent = outbox.mark_sent(replayed.message_id)
|
||||
|
||||
assert failed.status == "failed"
|
||||
assert failed.attempts == 1
|
||||
assert replayed.status == "pending"
|
||||
assert replayed.message_id == pending.message_id
|
||||
assert replayed.idempotency_key == pending.idempotency_key
|
||||
assert sent.status == "sent"
|
||||
assert outbox.replay(sent.message_id) == sent
|
||||
with pytest.raises(ValueError, match="sent"):
|
||||
outbox.mark_failed(sent.message_id, error="must not regress")
|
||||
|
||||
|
||||
def test_outbox_dispatch_never_redelivers_sent_and_retries_with_same_key(tmp_path: Path) -> None:
|
||||
outbox = Outbox(tmp_path)
|
||||
pending = outbox.enqueue(
|
||||
idempotency_key="run-001:publish",
|
||||
topic="sink.publish",
|
||||
payload={"artifact_id": "artifact-001"},
|
||||
run_id="run-001",
|
||||
)
|
||||
received_keys: list[str] = []
|
||||
|
||||
def fail_once(message: object) -> None:
|
||||
key = getattr(message, "idempotency_key")
|
||||
received_keys.append(key)
|
||||
if len(received_keys) == 1:
|
||||
raise RuntimeError("temporary")
|
||||
|
||||
failed = outbox.dispatch(pending.message_id, fail_once)
|
||||
assert failed.status == "failed"
|
||||
assert "RuntimeError" in (failed.error or "")
|
||||
|
||||
outbox.replay(failed.message_id)
|
||||
sent = outbox.dispatch(failed.message_id, fail_once)
|
||||
sent_again = outbox.dispatch(sent.message_id, fail_once)
|
||||
|
||||
assert sent.status == "sent"
|
||||
assert sent_again == sent
|
||||
assert received_keys == [pending.idempotency_key, pending.idempotency_key]
|
||||
|
||||
|
||||
def test_outbox_rejects_tampered_state_invariants(tmp_path: Path) -> None:
|
||||
outbox = Outbox(tmp_path)
|
||||
message = outbox.enqueue(
|
||||
idempotency_key="run-001:tamper",
|
||||
topic="audit.write",
|
||||
payload={},
|
||||
run_id="run-001",
|
||||
)
|
||||
path = outbox.message_path(message.message_id)
|
||||
raw = json.loads(path.read_text(encoding="utf-8"))
|
||||
raw["status"] = "sent"
|
||||
raw["sent_at"] = None
|
||||
path.write_text(json.dumps(raw), encoding="utf-8")
|
||||
|
||||
with pytest.raises(ValueError, match="invalid outbox message"):
|
||||
outbox.get(message.message_id)
|
||||
|
||||
|
||||
def test_outbox_files_are_structured_and_do_not_store_callable_results(tmp_path: Path) -> None:
|
||||
outbox = Outbox(tmp_path)
|
||||
message = outbox.enqueue(
|
||||
idempotency_key="run-001:audit",
|
||||
topic="audit.write",
|
||||
payload={"value": 1},
|
||||
run_id="run-001",
|
||||
)
|
||||
|
||||
raw = json.loads(outbox.message_path(message.message_id).read_text(encoding="utf-8"))
|
||||
|
||||
assert raw["schema_version"] == 1
|
||||
assert raw["status"] == "pending"
|
||||
assert raw["payload"] == {"value": 1}
|
||||
assert set(raw) == {
|
||||
"schema_version",
|
||||
"message_id",
|
||||
"idempotency_key",
|
||||
"topic",
|
||||
"payload",
|
||||
"run_id",
|
||||
"status",
|
||||
"attempts",
|
||||
"created_at",
|
||||
"updated_at",
|
||||
"sent_at",
|
||||
"error",
|
||||
}
|
||||
Reference in New Issue
Block a user