From 7a14237853ee1a354ed24fccb4ddd07a9e1eb37e Mon Sep 17 00:00:00 2001 From: eshulman2 Date: Thu, 27 Aug 2026 21:48:27 +0300 Subject: [PATCH 1/3] feat: add convergent observation ledger --- src/forge/domain/observations.py | 18 +- src/forge/reconciliation/__init__.py | 27 +++ src/forge/reconciliation/ledger.py | 280 +++++++++++++++++++++++ src/forge/reconciliation/models.py | 38 +++ tests/unit/reconciliation/test_ledger.py | 97 ++++++++ 5 files changed, 459 insertions(+), 1 deletion(-) create mode 100644 src/forge/reconciliation/__init__.py create mode 100644 src/forge/reconciliation/ledger.py create mode 100644 src/forge/reconciliation/models.py create mode 100644 tests/unit/reconciliation/test_ledger.py diff --git a/src/forge/domain/observations.py b/src/forge/domain/observations.py index f668d3a9..8b4419f8 100644 --- a/src/forge/domain/observations.py +++ b/src/forge/domain/observations.py @@ -7,7 +7,7 @@ from pydantic import Field -from forge.domain.identity import ResourceIdentity +from forge.domain.identity import ResourceIdentity, stable_identity from forge.domain.schema import JsonValue, VersionedDomainModel @@ -23,8 +23,24 @@ class Observation(VersionedDomainModel): source_system: str = Field(min_length=1) resource: ResourceIdentity resource_revision: str | None = None + revision_order: int | None = Field(default=None, ge=0) observed_at: datetime received_at: datetime facts: dict[str, JsonValue] = Field(default_factory=dict) correlation: dict[str, JsonValue] = Field(default_factory=dict) evidence_reference: str | None = None + + @property + def delivery_identity(self) -> str: + """Identity shared by poller and webhook deliveries of the same revision.""" + return stable_identity( + "observation-delivery", + { + "source_system": self.source_system, + "resource_type": self.resource.resource_type, + "external_id": self.resource.external_id, + "namespace": self.resource.namespace, + "resource_revision": self.resource_revision, + "revision_order": self.revision_order, + }, + ) diff --git a/src/forge/reconciliation/__init__.py b/src/forge/reconciliation/__init__.py new file mode 100644 index 00000000..acbf1620 --- /dev/null +++ b/src/forge/reconciliation/__init__.py @@ -0,0 +1,27 @@ +"""Convergent webhook and poller observation handling.""" + +from forge.reconciliation.ledger import ( + InMemoryObservationLedger, + ObservationLedger, + RedisObservationLedger, + classify_observation, + resource_identity, +) +from forge.reconciliation.models import ( + DriftClass, + ObservationDecision, + ObservationDisposition, + ReconciledResource, +) + +__all__ = [ + "DriftClass", + "InMemoryObservationLedger", + "ObservationDecision", + "ObservationDisposition", + "ObservationLedger", + "RedisObservationLedger", + "ReconciledResource", + "classify_observation", + "resource_identity", +] diff --git a/src/forge/reconciliation/ledger.py b/src/forge/reconciliation/ledger.py new file mode 100644 index 00000000..e5e31ecb --- /dev/null +++ b/src/forge/reconciliation/ledger.py @@ -0,0 +1,280 @@ +"""Observation ledger with source-independent deduplication and monotonic revisions.""" + +from __future__ import annotations + +import asyncio +from collections.abc import Sequence +from datetime import UTC, datetime +from typing import Any, Protocol + +from redis.exceptions import WatchError + +from forge.domain import Observation, stable_identity +from forge.orchestrator.checkpointer import get_redis_client +from forge.reconciliation.models import ( + DriftClass, + ObservationDecision, + ObservationDisposition, + ReconciledResource, +) + +PROTECTED_WORKFLOW_FACTS = { + "current_node", + "workflow_name", + "workflow_revision", + "workflow_digest", + "workflow_transition_count", +} +_RESOURCE_PREFIX = "forge:observations:resource:" +_DELIVERY_PREFIX = "forge:observations:delivery:" +_HISTORY_PREFIX = "forge:observations:history:" + + +class ObservationLedger(Protocol): + async def record(self, observation: Observation) -> ObservationDecision: ... + + async def latest(self, observation: Observation) -> ReconciledResource | None: ... + + async def history(self, observation: Observation) -> Sequence[ObservationDecision]: ... + + +def resource_identity(observation: Observation) -> str: + return stable_identity( + "observed-resource", + { + "source_system": observation.source_system, + "resource_type": observation.resource.resource_type, + "external_id": observation.resource.external_id, + "namespace": observation.resource.namespace, + }, + ) + + +class InMemoryObservationLedger: + """Reference implementation used by ingress conformance fixtures.""" + + def __init__(self) -> None: + self._resources: dict[str, ReconciledResource] = {} + self._history: dict[str, list[ObservationDecision]] = {} + self._deliveries: dict[str, ObservationDecision] = {} + self._lock = asyncio.Lock() + + async def record(self, observation: Observation) -> ObservationDecision: + async with self._lock: + protected = sorted(PROTECTED_WORKFLOW_FACTS & observation.facts.keys()) + if protected: + decision = _decision( + observation, + ObservationDisposition.CONFLICT, + DriftClass.POLICY_BLOCKING, + f"external observation attempted to set workflow-owned facts: {protected}", + ) + self._append(observation, decision) + return decision + delivery = observation.delivery_identity + duplicate = self._deliveries.get(delivery) + if duplicate is not None: + same_facts = duplicate.observation.facts == observation.facts + decision = _decision( + observation, + ObservationDisposition.DUPLICATE + if same_facts + else ObservationDisposition.CONFLICT, + DriftClass.EXPECTED if same_facts else DriftClass.OPERATOR_REQUIRED, + "provider revision was already observed through an ingress source" + if same_facts + else "same provider revision contains different facts", + ) + self._append(observation, decision) + return decision + + key = resource_identity(observation) + current = self._resources.get(key) + disposition, drift, reason = classify_observation( + current.latest if current else None, observation + ) + decision = _decision( + observation, + disposition, + drift, + reason, + supersedes=current.latest_delivery_identity + if current and disposition is ObservationDisposition.ACCEPTED + else None, + ) + self._deliveries[delivery] = decision + self._append(observation, decision) + if disposition is ObservationDisposition.ACCEPTED: + self._resources[key] = ReconciledResource( + latest=observation, + latest_delivery_identity=delivery, + updated_at=decision.decided_at, + ) + return decision + + def _append(self, observation: Observation, decision: ObservationDecision) -> None: + self._history.setdefault(resource_identity(observation), []).append(decision) + + async def latest(self, observation: Observation) -> ReconciledResource | None: + return self._resources.get(resource_identity(observation)) + + async def history(self, observation: Observation) -> Sequence[ObservationDecision]: + return tuple(self._history.get(resource_identity(observation), ())) + + +class RedisObservationLedger: + """Production ledger using optimistic transactions for monotonic acceptance.""" + + def __init__(self, redis_client: Any = None) -> None: + self._redis = redis_client + + async def _client(self) -> Any: + if self._redis is None: + self._redis = await get_redis_client() + return self._redis + + async def record(self, observation: Observation) -> ObservationDecision: + protected = sorted(PROTECTED_WORKFLOW_FACTS & observation.facts.keys()) + if protected: + decision = _decision( + observation, + ObservationDisposition.CONFLICT, + DriftClass.POLICY_BLOCKING, + f"external observation attempted to set workflow-owned facts: {protected}", + ) + await (await self._client()).rpush( + self._history_key(observation), decision.model_dump_json() + ) + return decision + + redis = await self._client() + resource_key = self._resource_key(observation) + delivery_key = f"{_DELIVERY_PREFIX}{observation.delivery_identity}" + while True: + async with redis.pipeline(transaction=True) as pipeline: + try: + await pipeline.watch(resource_key, delivery_key) + delivery_raw = await pipeline.get(delivery_key) + current_raw = await pipeline.get(resource_key) + prior = ( + ObservationDecision.model_validate_json(delivery_raw) + if delivery_raw + else None + ) + current = ( + ReconciledResource.model_validate_json(current_raw) if current_raw else None + ) + if prior is not None: + same_facts = prior.observation.facts == observation.facts + decision = _decision( + observation, + ObservationDisposition.DUPLICATE + if same_facts + else ObservationDisposition.CONFLICT, + DriftClass.EXPECTED if same_facts else DriftClass.OPERATOR_REQUIRED, + "provider revision was already observed through an ingress source" + if same_facts + else "same provider revision contains different facts", + ) + else: + disposition, drift, reason = classify_observation( + current.latest if current else None, observation + ) + decision = _decision( + observation, + disposition, + drift, + reason, + current.latest_delivery_identity + if current and disposition is ObservationDisposition.ACCEPTED + else None, + ) + pipeline.multi() + pipeline.rpush(self._history_key(observation), decision.model_dump_json()) + if prior is None: + pipeline.set(delivery_key, decision.model_dump_json()) + if decision.disposition is ObservationDisposition.ACCEPTED: + reconciled = ReconciledResource( + latest=observation, + latest_delivery_identity=observation.delivery_identity, + updated_at=decision.decided_at, + ) + pipeline.set(resource_key, reconciled.model_dump_json()) + await pipeline.execute() + return decision + except WatchError: + continue + + async def latest(self, observation: Observation) -> ReconciledResource | None: + value = await (await self._client()).get(self._resource_key(observation)) + return ReconciledResource.model_validate_json(value) if value else None + + async def history(self, observation: Observation) -> Sequence[ObservationDecision]: + values = await (await self._client()).lrange(self._history_key(observation), 0, -1) + return tuple(ObservationDecision.model_validate_json(value) for value in values) + + @staticmethod + def _resource_key(observation: Observation) -> str: + return f"{_RESOURCE_PREFIX}{resource_identity(observation)}" + + @staticmethod + def _history_key(observation: Observation) -> str: + return f"{_HISTORY_PREFIX}{resource_identity(observation)}" + + +def classify_observation( + current: Observation | None, incoming: Observation +) -> tuple[ObservationDisposition, DriftClass, str]: + if current is None: + return ObservationDisposition.ACCEPTED, DriftClass.EXPECTED, "first observed revision" + if incoming.revision_order is not None and current.revision_order is not None: + if incoming.revision_order < current.revision_order: + return ObservationDisposition.STALE, DriftClass.EXPECTED, "older provider revision" + if incoming.revision_order == current.revision_order: + if incoming.facts == current.facts: + return ( + ObservationDisposition.DUPLICATE, + DriftClass.EXPECTED, + "same revision and facts", + ) + return ( + ObservationDisposition.CONFLICT, + DriftClass.OPERATOR_REQUIRED, + "same provider revision contains different facts", + ) + return ( + ObservationDisposition.ACCEPTED, + DriftClass.AUTO_RECONCILABLE, + "newer provider revision updates the external projection", + ) + if incoming.resource_revision == current.resource_revision: + if incoming.facts == current.facts: + return ObservationDisposition.DUPLICATE, DriftClass.EXPECTED, "same revision and facts" + return ( + ObservationDisposition.CONFLICT, + DriftClass.OPERATOR_REQUIRED, + "opaque provider revision contains different facts", + ) + return ( + ObservationDisposition.CONFLICT, + DriftClass.OPERATOR_REQUIRED, + "opaque revisions cannot be ordered safely", + ) + + +def _decision( + observation: Observation, + disposition: ObservationDisposition, + drift: DriftClass, + reason: str, + supersedes: str | None = None, +) -> ObservationDecision: + return ObservationDecision( + observation=observation, + delivery_identity=observation.delivery_identity, + disposition=disposition, + drift=drift, + reason=reason, + decided_at=datetime.now(UTC), + supersedes_delivery_identity=supersedes, + ) diff --git a/src/forge/reconciliation/models.py b/src/forge/reconciliation/models.py new file mode 100644 index 00000000..86e4d798 --- /dev/null +++ b/src/forge/reconciliation/models.py @@ -0,0 +1,38 @@ +"""Durable decisions made while reconciling external observations.""" + +from __future__ import annotations + +from datetime import datetime +from enum import StrEnum + +from forge.domain import DomainModel, Observation + + +class ObservationDisposition(StrEnum): + ACCEPTED = "accepted" + DUPLICATE = "duplicate" + STALE = "stale" + CONFLICT = "conflict" + + +class DriftClass(StrEnum): + EXPECTED = "expected" + AUTO_RECONCILABLE = "auto_reconcilable" + POLICY_BLOCKING = "policy_blocking" + OPERATOR_REQUIRED = "operator_required" + + +class ObservationDecision(DomainModel): + observation: Observation + delivery_identity: str + disposition: ObservationDisposition + drift: DriftClass + reason: str + decided_at: datetime + supersedes_delivery_identity: str | None = None + + +class ReconciledResource(DomainModel): + latest: Observation + latest_delivery_identity: str + updated_at: datetime diff --git a/tests/unit/reconciliation/test_ledger.py b/tests/unit/reconciliation/test_ledger.py new file mode 100644 index 00000000..01304535 --- /dev/null +++ b/tests/unit/reconciliation/test_ledger.py @@ -0,0 +1,97 @@ +from datetime import UTC, datetime + +import pytest + +from forge.domain import Observation, ObservationSource, ResourceIdentity +from forge.reconciliation import ( + DriftClass, + InMemoryObservationLedger, + ObservationDisposition, +) + + +def observation( + source: ObservationSource, + order: int, + *, + status: str = "open", +) -> Observation: + now = datetime.now(UTC) + return Observation( + observation_id=f"{source}-{order}", + source=source, + source_system="github", + resource=ResourceIdentity( + resource_type="change_request", external_id="17", namespace="org/repo" + ), + resource_revision=f"revision-{order}", + revision_order=order, + observed_at=now, + received_at=now, + facts={"status": status}, + ) + + +@pytest.mark.asyncio +async def test_webhook_and_poller_delivery_share_identity_and_deduplicate() -> None: + ledger = InMemoryObservationLedger() + webhook = observation(ObservationSource.WEBHOOK, 4) + polled = observation(ObservationSource.POLLER, 4) + + first = await ledger.record(webhook) + duplicate = await ledger.record(polled) + + assert webhook.delivery_identity == polled.delivery_identity + assert first.disposition is ObservationDisposition.ACCEPTED + assert duplicate.disposition is ObservationDisposition.DUPLICATE + + +@pytest.mark.asyncio +async def test_stale_delivery_cannot_overwrite_latest_projection() -> None: + ledger = InMemoryObservationLedger() + newest = observation(ObservationSource.WEBHOOK, 5, status="merged") + stale = observation(ObservationSource.POLLER, 3) + await ledger.record(newest) + + decision = await ledger.record(stale) + + assert decision.disposition is ObservationDisposition.STALE + assert (await ledger.latest(stale)).latest.facts == {"status": "merged"} + + +@pytest.mark.asyncio +async def test_same_revision_with_different_facts_requires_operator() -> None: + ledger = InMemoryObservationLedger() + await ledger.record(observation(ObservationSource.WEBHOOK, 5)) + + decision = await ledger.record(observation(ObservationSource.POLLER, 5, status="closed")) + + assert decision.disposition is ObservationDisposition.CONFLICT + assert decision.drift is DriftClass.OPERATOR_REQUIRED + + +@pytest.mark.asyncio +async def test_newer_observation_updates_projection_without_workflow_position() -> None: + ledger = InMemoryObservationLedger() + first = observation(ObservationSource.WEBHOOK, 1) + await ledger.record(first) + + decision = await ledger.record(observation(ObservationSource.POLLER, 2, status="merged")) + + assert decision.disposition is ObservationDisposition.ACCEPTED + assert decision.drift is DriftClass.AUTO_RECONCILABLE + assert "current_node" not in decision.observation.facts + + +@pytest.mark.asyncio +async def test_external_observation_cannot_overwrite_workflow_position() -> None: + ledger = InMemoryObservationLedger() + incoming = observation(ObservationSource.POLLER, 1).model_copy( + update={"facts": {"status": "merged", "current_node": "complete"}} + ) + + decision = await ledger.record(incoming) + + assert decision.disposition is ObservationDisposition.CONFLICT + assert decision.drift is DriftClass.POLICY_BLOCKING + assert await ledger.latest(incoming) is None From 91a482ef95ca895c1b0f4a437b3c90e3ddd03737 Mon Sep 17 00:00:00 2001 From: eshulman2 Date: Thu, 27 Aug 2026 21:52:39 +0300 Subject: [PATCH 2/3] feat: reconcile observations before command interpretation --- src/forge/orchestrator/worker.py | 27 ++++++++++++++++++++++++++ tests/unit/orchestrator/test_worker.py | 9 +++++++++ 2 files changed, 36 insertions(+) diff --git a/src/forge/orchestrator/worker.py b/src/forge/orchestrator/worker.py index a15d3431..c1694fdc 100644 --- a/src/forge/orchestrator/worker.py +++ b/src/forge/orchestrator/worker.py @@ -67,6 +67,11 @@ from forge.orchestrator.review_enrichment import ReviewEnrichmentService from forge.queue.consumer import QueueConsumer from forge.queue.models import QueueMessage +from forge.reconciliation import ( + ObservationDisposition, + ObservationLedger, + RedisObservationLedger, +) from forge.skills.orchestrator import ensure_skills from forge.skills.utils import extract_project_key from forge.utils.redaction import redact_secrets @@ -204,6 +209,7 @@ def __init__( command_handlers: CommandHandlerRegistry | None = None, review_enrichment: ReviewEnrichmentService | None = None, effect_service: EffectService | None = None, + observation_ledger: ObservationLedger | None = None, ) -> None: """Initialize the worker. @@ -222,6 +228,7 @@ def __init__( self.command_handlers = command_handlers or create_default_command_handler_registry() self.review_enrichment = review_enrichment self.effect_service = effect_service or create_default_effect_service() + self.observation_ledger = observation_ledger or RedisObservationLedger() self._shutdown_event = asyncio.Event() self._checkpointer = None self._compiled_workflows: dict[str, Any] = {} # Cache compiled workflows by name @@ -274,6 +281,14 @@ async def _invoke_workflow( with bind_effect_runtime(self._durable_effect_service(), identity): return await compiled_workflow.ainvoke(invocation_input, config=config) + def _observation_ledger(self) -> ObservationLedger: + """Lazily restore reconciliation for fixtures that bypass ``__init__``.""" + ledger = getattr(self, "observation_ledger", None) + if ledger is None: + ledger = RedisObservationLedger() + self.observation_ledger = ledger + return ledger + async def _execute_required_jira_effect( self, *, @@ -544,6 +559,18 @@ async def _process_workflow(self, message: QueueMessage) -> None: try: ingress = self.event_adapters.adapt(message) + observation_decision = await self._observation_ledger().record(ingress.observation) + if observation_decision.disposition in { + ObservationDisposition.STALE, + ObservationDisposition.CONFLICT, + }: + logger.info( + "Ignoring %s observation %s: %s", + observation_decision.disposition.value, + ingress.observation.observation_id, + observation_decision.reason, + ) + return # Determine ticket type early to select workflow ticket_type = self._extract_ticket_type(message) diff --git a/tests/unit/orchestrator/test_worker.py b/tests/unit/orchestrator/test_worker.py index 6148dc8c..a05074c4 100644 --- a/tests/unit/orchestrator/test_worker.py +++ b/tests/unit/orchestrator/test_worker.py @@ -30,6 +30,7 @@ QueueMessage, normalized_event_to_dict, ) +from forge.reconciliation import InMemoryObservationLedger from forge.workflow.utils.source_control import identity_for @@ -49,6 +50,14 @@ def durable_effect_service_mock(): yield service +@pytest.fixture(autouse=True) +def observation_ledger_mock(monkeypatch): + """Keep unit workers isolated from the production Redis observation ledger.""" + ledger = InMemoryObservationLedger() + monkeypatch.setattr("forge.orchestrator.worker.RedisObservationLedger", lambda: ledger) + return ledger + + @pytest.mark.parametrize( ("result", "error_before_invoke", "expected"), [ From 028fe32e2a66e1e9014646bfaa24dd00840eac16 Mon Sep 17 00:00:00 2001 From: eshulman2 Date: Fri, 28 Aug 2026 00:06:29 +0300 Subject: [PATCH 3/3] Complete webhook and poller reconciliation --- docs/architecture/internals.md | 6 +- docs/architecture/option-b-completion-plan.md | 15 +- .../phase-6-observation-contract.md | 54 ++++ .../phase-6-reconciliation-contract.md | 56 +++++ docs/architecture/reference.md | 4 +- src/forge/domain/__init__.py | 9 +- src/forge/domain/observations.py | 74 +++++- .../source_control/observations.py | 62 +++-- .../orchestrator/event_adapters/commands.py | 6 +- src/forge/orchestrator/event_adapters/jira.py | 137 +++++++++- src/forge/orchestrator/worker.py | 1 + src/forge/reconciliation/ledger.py | 112 ++++++++- .../github_pull_request_revision.json | 53 ++++ .../fixtures/reconciliation/README.md | 10 + .../source_control_sequence.json | 34 +++ tests/contracts/reconciliation/__init__.py | 1 + .../reconciliation/test_convergence.py | 235 ++++++++++++++++++ tests/contracts/test_observation_contract.py | 27 ++ .../source_control/test_observations.py | 30 +++ .../event_adapters/test_registry.py | 111 +++++++++ .../test_reconciliation_worker.py | 47 ++++ tests/unit/reconciliation/test_ledger.py | 79 ++++++ zensical.toml | 2 + 23 files changed, 1104 insertions(+), 61 deletions(-) create mode 100644 docs/architecture/phase-6-observation-contract.md create mode 100644 docs/architecture/phase-6-reconciliation-contract.md create mode 100644 tests/contracts/fixtures/observations/github_pull_request_revision.json create mode 100644 tests/contracts/fixtures/reconciliation/README.md create mode 100644 tests/contracts/fixtures/reconciliation/source_control_sequence.json create mode 100644 tests/contracts/reconciliation/__init__.py create mode 100644 tests/contracts/reconciliation/test_convergence.py create mode 100644 tests/contracts/test_observation_contract.py create mode 100644 tests/unit/orchestrator/test_reconciliation_worker.py diff --git a/docs/architecture/internals.md b/docs/architecture/internals.md index cae45a50..836329a9 100644 --- a/docs/architecture/internals.md +++ b/docs/architecture/internals.md @@ -16,7 +16,11 @@ Gateway and Worker communicate only through Redis and can be deployed on separat **Checkpointing:** LangGraph workflow state is persisted via `AsyncRedisSaver`, keyed by Jira ticket key (e.g., `AISOS-123`). Checkpoints are written after each graph node completes. When a new event arrives for an existing ticket, the workflow resumes from its last checkpoint. -**Idempotency:** A `DeduplicationService` exists but is not yet wired into the webhook routes. Branch creation and label operations are naturally idempotent; Jira comment posting is not. +**Idempotency:** The worker records normalized webhook/poller observations in +the reconciliation ledger before command interpretation. Duplicate, stale, +and conflicting observations do not re-enter the workflow. External writes +still use durable effect identities; provider operations such as Jira comment +posting must remain idempotent across crash recovery. **Consistency boundary:** Workflow mutations are persisted as stable effect intents before provider execution. Provider-specific recovery evidence and idempotent ref updates close diff --git a/docs/architecture/option-b-completion-plan.md b/docs/architecture/option-b-completion-plan.md index 0ce7dbae..06aacba7 100644 --- a/docs/architecture/option-b-completion-plan.md +++ b/docs/architecture/option-b-completion-plan.md @@ -125,7 +125,7 @@ Exit gate: the supported process can be rendered from a versioned definition, an ## Phase 6 — Reconciliation and poller convergence -PR: new stacked PR based on Phase 5. Status: missing. +PR: #331 (`phase6/reconciliation-contract`). Status: complete. Purpose: make webhooks and the existing forge-poller equivalent observation sources while workflow state remains authoritative for transition progress. @@ -141,6 +141,17 @@ Work: Exit gate: lost, duplicated, and reordered delivery converges to the same workflow/effect state without duplicate effects. +Evidence: the versioned Observation fixtures in +`tests/contracts/fixtures/observations/` and +`tests/contracts/fixtures/reconciliation/`, adapter and identity tests in +`tests/contracts/` and `tests/unit/integrations/source_control/`, ledger +classification tests in `tests/unit/reconciliation/`, and replay/worker +convergence tests in `tests/contracts/reconciliation/` and +`tests/unit/orchestrator/test_reconciliation_worker.py`. The explicit +no-native-revision limitation is documented in +`docs/architecture/phase-6-observation-contract.md`; such input is retained +for operator review and never used to infer workflow position. + ## Phase 7 — Execution read models and operations PR: #329, rebased onto Phase 6. Status: partial. @@ -181,7 +192,7 @@ Exit gate: all supported flows use normalized observations, authoritative comman 2. Finish #326, then rebase all descendants. 3. Finish #327, then rebase all descendants. 4. Finish #328. -5. Create `phase6/reconciliation-contract` from the Phase 5 tip and open its stacked PR. +5. Phase 6 is implemented by Forge PR #331 and the companion forge-poller PR #10. 6. Rebase `phase7/execution-read-models` from the old Phase 5 base onto Phase 6; change #329's base to the Phase 6 branch. 7. Rebase `phase8-compatibility-removal` onto the rebased Phase 7 branch; retain #330's Phase 7 base. 8. Finish Phase 6, then Phase 7, then Phase 8, rebasing descendants after each phase changes. diff --git a/docs/architecture/phase-6-observation-contract.md b/docs/architecture/phase-6-observation-contract.md new file mode 100644 index 00000000..26f765f4 --- /dev/null +++ b/docs/architecture/phase-6-observation-contract.md @@ -0,0 +1,54 @@ +# Phase 6 Observation contract + +Forge accepts webhook and `forge-poller` deliveries through the same ingress +adapters. Both paths are normalized to the versioned `Observation` record +(`schema_version: "1.0"`). The strict contract contains source, +provider/resource identity, provider revision, observation times, normalized +facts, correlation metadata, and an optional evidence reference. Unknown +fields and future schema versions are rejected at the boundary. + +## Identity + +`observation_id` identifies the provider event record. It is generated from +the provider event ID plus provider/resource identity and revision context; it +does not include `source`. Replays that preserve the provider event ID retain +the same observation ID. When a poller must use a different transport ID, +`delivery_identity` still remains stable whenever the provider revision is +present. + +`delivery_identity` is the deduplication identity used by Forge's observation +ledger. It is generated from source system, resource identity, and +`resource_revision`. Therefore webhook and poller deliveries of one provider +revision have the same delivery identity even when their transport event IDs, +received times, or source values differ. `revision_order` is ordering metadata +and is not included when a provider revision is available. For resources with +no revision, the provider event identity is used; callers must provide +`correlation.provider_event_id` (or `transport_event_id`) in that case. + +The Jira adapter derives revisions from the immutable comment ID, the issue's +`updated` timestamp, or a changelog fingerprint. Source-control adapters use +the change-request head SHA, check state scoped to its commit, or immutable +comment/review IDs. The poller can therefore forward its existing Jira and +GitHub payloads; Forge assigns the source-independent identity before ledger +processing. Command IDs likewise use `delivery_identity`, so a poller retry +cannot create a second command/effect for a webhook-delivered revision. + +An observation without a native provider revision is deliberately limited: it +is deduplicated only when its provider event ID is stable, and a later +revision cannot be ordered safely. Forge records such input as an +operator-visible conflict rather than guessing which external state is newer. +This is the explicit no-native-revision limitation, not a workflow checkpoint. + +The shared fixtures at +[`github_pull_request_revision.json`](../../tests/contracts/fixtures/observations/github_pull_request_revision.json) +and +[`source_control_sequence.json`](../../tests/contracts/fixtures/reconciliation/source_control_sequence.json) +show equivalent source-control revisions and replay behavior. Contract tests +cover both ingress markers, cross-source deduplication, monotonic stale/reorder +handling, command identity, and the Jira revision derivation cases. + +The companion poller changes preserve GitHub review/comment IDs and head SHAs, +include Jira issue `updated` values and comment IDs/timestamps in forwarded +payloads, and make synthetic delivery IDs replay-stable. The poller contract +tests cover those payload guarantees; Forge's adapter tests then verify their +identity and reconciliation semantics. diff --git a/docs/architecture/phase-6-reconciliation-contract.md b/docs/architecture/phase-6-reconciliation-contract.md new file mode 100644 index 00000000..86ef2d5e --- /dev/null +++ b/docs/architecture/phase-6-reconciliation-contract.md @@ -0,0 +1,56 @@ +# Phase 6: reconciliation contract + +Phase 6 makes webhook and polling delivery interchangeable inputs to Forge. +Workflow position remains owned by the workflow instance; an external +observation may update an external-state projection but cannot set +`current_node`, workflow identity, or transition counters. + +## Conformance contract + +Every ingress source must provide a versioned `Observation` with: + +- `source`: `webhook` or `poller` (transport metadata only); +- provider/resource identity (`source_system`, `resource`, and + `resource_revision`); +- a stable provider event identity in `correlation.provider_event_id`; +- monotonic `revision_order` whenever a provider revision cannot be ordered + from its native identifier; and +- provider facts that are identical for equivalent revisions. + +`Observation.delivery_identity` intentionally excludes `source`, so equivalent +webhook and poller deliveries deduplicate. A command may be evaluated only for +an accepted observation. Duplicate, stale, or conflicting deliveries are +recorded for inspection and cannot create another external effect. + +The shared fixture is +[`source_control_sequence.json`](../../tests/contracts/fixtures/reconciliation/source_control_sequence.json). +It is JSON rather than a Python fixture so Forge and `forge-poller` can replay +the same provider revisions. Forge's contract tests cover equivalent envelopes +and replay the sequence with a lost first revision, duplicate cross-source +delivery, stale reordered delivery, and a duplicate replay. The resulting +latest revision, accepted command, and effect list must match the clean replay. + +Run the Forge side with: + +```shell +pytest -q tests/contracts/reconciliation tests/unit/reconciliation +``` + +## Poller companion behavior and limitation + +`forge-poller` remains an independently deployable delivery source. Its +webhook-shaped Jira and GitHub payloads preserve the provider IDs and revision +metadata needed by Forge's normalizing adapters, which assign the +source-independent revision and delivery identity before workflow evaluation. +Poller cursors are delivery optimizations only; they are not read or written +as workflow checkpoints. +Consequently a lost, repeated, or reordered delivery cannot move workflow +position or create a duplicate command/effect. + +The provider limitation is explicit: if a payload contains neither a native +revision nor a stable provider event ID, Forge cannot establish ordering. It +records the observation but classifies a competing update as +`operator_required`; it does not infer chronology from receipt time or a +poller cursor. Jira payloads lacking `issue.fields.updated`, changelog items, +or a comment ID are in this category and require the provider/poller to add a +stable revision before automatic convergence is possible. diff --git a/docs/architecture/reference.md b/docs/architecture/reference.md index 983fb39e..5ee66e3f 100644 --- a/docs/architecture/reference.md +++ b/docs/architecture/reference.md @@ -26,7 +26,9 @@ Workflows pause at defined gates and wait indefinitely for human approval. The ` - **No PEL reclaim**: Unacknowledged messages from crashed workers remain in Redis PEL indefinitely. Recovery requires manual `XCLAIM`. - **No distributed per-ticket lock**: Multiple workers can process events for the same ticket concurrently, causing potential checkpoint conflicts. -- **Webhook deduplication not wired**: `DeduplicationService` exists but is not connected to webhook routes. +- **Ingress delivery is at-least-once**: webhook and poller observations are + deduplicated and classified by the reconciliation ledger, but the gateway + may still enqueue a retried transport message before the worker records it. - **Webhook signature validation is optional**: Endpoints accept unsigned payloads when secrets are not configured. - **No approval gate timeout**: Paused workflows wait indefinitely with no escalation. - **Single Redis dependency**: No Sentinel, Cluster, or HA. Redis is a single point of failure. diff --git a/src/forge/domain/__init__.py b/src/forge/domain/__init__.py index ee6835b4..81b54988 100644 --- a/src/forge/domain/__init__.py +++ b/src/forge/domain/__init__.py @@ -9,7 +9,12 @@ stable_identity, ) from forge.domain.interactions import CommentType, classify_comment -from forge.domain.observations import Observation, ObservationSource +from forge.domain.observations import ( + Observation, + ObservationSource, + observation_delivery_identity, + observation_identity, +) from forge.domain.schema import DomainModel, JsonValue, VersionedDomainModel from forge.domain.stations import ( StationFailure, @@ -27,6 +32,8 @@ "JsonValue", "Observation", "ObservationSource", + "observation_delivery_identity", + "observation_identity", "ResourceIdentity", "StationFailure", "StationInvocationIdentity", diff --git a/src/forge/domain/observations.py b/src/forge/domain/observations.py index 8b4419f8..e1f17c56 100644 --- a/src/forge/domain/observations.py +++ b/src/forge/domain/observations.py @@ -33,14 +33,66 @@ class Observation(VersionedDomainModel): @property def delivery_identity(self) -> str: """Identity shared by poller and webhook deliveries of the same revision.""" - return stable_identity( - "observation-delivery", - { - "source_system": self.source_system, - "resource_type": self.resource.resource_type, - "external_id": self.resource.external_id, - "namespace": self.resource.namespace, - "resource_revision": self.resource_revision, - "revision_order": self.revision_order, - }, - ) + return observation_delivery_identity(self) + + +def observation_delivery_identity(observation: Observation) -> str: + """Return the source-independent identity of an observation delivery. + + A provider revision is the strongest identity available: delivery IDs are + transport metadata and differ when the same state arrives from a webhook + and from the poller. For event-shaped resources that do not expose a + revision, the provider event ID (stored in correlation metadata) keeps + distinct events from collapsing into one delivery. ``observation_id`` is + the final fallback for callers constructing an observation without either + kind of provider identity. + + ``revision_order`` is deliberately not included when ``resource_revision`` + is present. It is ordering metadata, not part of the provider revision; + including it would make equivalent deliveries deduplicate differently. + """ + parts: dict[str, JsonValue] = { + "source_system": observation.source_system, + "resource_type": observation.resource.resource_type, + "external_id": observation.resource.external_id, + "namespace": observation.resource.namespace, + } + if observation.resource_revision is not None: + parts["resource_revision"] = observation.resource_revision + elif observation.revision_order is not None: + parts["revision_order"] = observation.revision_order + else: + provider_event_id = observation.correlation.get("provider_event_id") + if not isinstance(provider_event_id, str): + provider_event_id = observation.correlation.get("transport_event_id") + if isinstance(provider_event_id, str) and provider_event_id: + parts["provider_event_id"] = provider_event_id + else: + parts["observation_id"] = observation.observation_id + return stable_identity("observation-delivery", parts) + + +def observation_identity( + *, + source_system: str, + provider_event_id: str, + resource: ResourceIdentity, + resource_revision: str | None = None, +) -> str: + """Build the deterministic identity assigned to a provider observation. + + This identity remains stable when the delivery source changes. The event + ID distinguishes separate provider events, while the resource revision is + included as context for providers that reuse event IDs across resources. + """ + return stable_identity( + "observation", + { + "source_system": source_system, + "provider_event_id": provider_event_id, + "resource_type": resource.resource_type, + "external_id": resource.external_id, + "namespace": resource.namespace, + "resource_revision": resource_revision, + }, + ) diff --git a/src/forge/integrations/source_control/observations.py b/src/forge/integrations/source_control/observations.py index 32f8b269..2c6cf126 100644 --- a/src/forge/integrations/source_control/observations.py +++ b/src/forge/integrations/source_control/observations.py @@ -12,7 +12,7 @@ Observation, ObservationSource, ResourceIdentity, - stable_identity, + observation_identity, ) from forge.integrations.source_control.contracts import NormalizedEvent @@ -44,20 +44,25 @@ def normalized_event_to_observation( external_id = event.repo_ref.id resource_type = "repository" revision: str | None = None - if change_request: - resource_type = "change_request" - external_id = f"{event.repo_ref.id}#{native_id}" - revision = change_request.head_sha or None - elif event.check: + # Event-specific resources must win over the optional change-request + # context attached by GitHub's check/comment/review payloads. Otherwise + # every check or comment on one head SHA is misidentified as that PR state. + if event.check: resource_type = "check" external_id = f"{event.repo_ref.id}:{event.check.name}" - revision = f"{event.check.status.value}:{event.check.conclusion.value}" + revision = _check_revision(event) elif event.comment: resource_type = "comment" external_id = f"{event.repo_ref.id}:{event.comment.id}" + revision = event.comment.id elif event.review: resource_type = "review" external_id = f"{event.repo_ref.id}:{event.review.id}" + revision = event.review.id + elif change_request: + resource_type = "change_request" + external_id = f"{event.repo_ref.id}#{native_id}" + revision = change_request.head_sha or None facts = _json_value( { @@ -72,27 +77,44 @@ def normalized_event_to_observation( } ) assert isinstance(facts, dict) - observation_id = stable_identity( - "observation", - { - "source_system": event.repo_ref.provider.value, - "event_id": event.id, - "resource_revision": revision, - }, + resource = ResourceIdentity( + resource_type=resource_type, + external_id=external_id, + namespace=event.repo_ref.connection, + ) + observation_id = observation_identity( + source_system=event.repo_ref.provider.value, + provider_event_id=event.id, + resource=resource, + resource_revision=revision, ) return Observation( observation_id=observation_id, source=source, source_system=event.repo_ref.provider.value, - resource=ResourceIdentity( - resource_type=resource_type, - external_id=external_id, - namespace=event.repo_ref.connection, - ), + resource=resource, resource_revision=revision, observed_at=event.received_at, received_at=event.received_at, facts=facts, - correlation={"transport_event_id": event.id, "repository_id": event.repo_ref.id}, + correlation={ + "provider_event_id": event.id, + "transport_event_id": event.id, + "repository_id": event.repo_ref.id, + }, evidence_reference=f"source-control-event:{event.id}" if event.raw else None, ) + + +def _check_revision(event: NormalizedEvent) -> str: + """Return a revision for a check state, scoped to the checked commit. + + ``CheckRun`` intentionally keeps only provider-neutral fields. Combining + those fields with the associated change-request head prevents a successful + check on two commits from being treated as one revision. For standalone + checks, the event ID is the provider's only available immutable identity. + """ + assert event.check is not None + head_sha = event.change_request.head_sha if event.change_request else "" + scope = head_sha or event.id + return f"{scope}:{event.check.status.value}:{event.check.conclusion.value}" diff --git a/src/forge/orchestrator/event_adapters/commands.py b/src/forge/orchestrator/event_adapters/commands.py index 67711855..20bbcd2a 100644 --- a/src/forge/orchestrator/event_adapters/commands.py +++ b/src/forge/orchestrator/event_adapters/commands.py @@ -141,7 +141,11 @@ def interpret_event( command_id = stable_identity( "workflow-command", { - "event_id": message.event_id, + # The transport event ID changes when one provider revision is + # delivered by both webhook and poller. Command identity must be + # tied to the source-independent observation so either delivery + # selects the same durable command/effect. + "observation_delivery_identity": adapted.observation.delivery_identity, "run_id": workflow.run_id, "command_type": command_type.value, }, diff --git a/src/forge/orchestrator/event_adapters/jira.py b/src/forge/orchestrator/event_adapters/jira.py index f3dd574c..8f3a583a 100644 --- a/src/forge/orchestrator/event_adapters/jira.py +++ b/src/forge/orchestrator/event_adapters/jira.py @@ -3,9 +3,16 @@ from __future__ import annotations import logging +from datetime import datetime from typing import Any -from forge.domain import Observation, ObservationSource, ResourceIdentity, stable_identity +from forge.domain import ( + Observation, + ObservationSource, + ResourceIdentity, + observation_identity, + stable_identity, +) from forge.models.events import EventSource from forge.models.workflow import TicketType from forge.orchestrator.event_adapters.contracts import AdaptedEvent, IngressMessage @@ -33,8 +40,9 @@ class JiraEventAdapter: def adapt(self, message: IngressMessage) -> AdaptedEvent: issue = message.payload.get("issue", {}) + issue_fields = issue.get("fields", {}) if isinstance(issue, dict) else {} ticket_key = str(issue.get("key") or message.ticket_key) - ticket_type_name = str(issue.get("fields", {}).get("issuetype", {}).get("name", "Unknown")) + ticket_type_name = str(issue_fields.get("issuetype", {}).get("name", "Unknown")) if ticket_type_name in {"Epic", "Task", "Sub-task"} and message.payload.get( "source_ticket_key" ): @@ -45,24 +53,33 @@ def adapt(self, message: IngressMessage) -> AdaptedEvent: except ValueError: logger.warning("Unknown ticket type '%s' for %s", ticket_type_name, ticket_key) ticket_type = TicketType.UNKNOWN + resource = ResourceIdentity(resource_type="issue", external_id=ticket_key) + comment = message.payload.get("comment") observation = Observation( - observation_id=stable_identity( - "observation", {"source_system": "jira", "event_id": message.event_id} + observation_id=observation_identity( + source_system="jira", + provider_event_id=message.event_id, + resource=resource, ), source=ObservationSource.WEBHOOK, source_system="jira", - resource=ResourceIdentity(resource_type="issue", external_id=ticket_key), + resource=resource, + resource_revision=_jira_revision(message.payload), + revision_order=_jira_revision_order(message.payload), observed_at=message.timestamp, received_at=message.timestamp, - facts={ - "event_type": message.event_type, - "issue": issue, - "changelog": message.payload.get("changelog", {}), - "comment": message.payload.get("comment"), - "comment_text": _comment_text(message.payload.get("comment", {}).get("body", "")), - "source_ticket_key": message.payload.get("source_ticket_key"), + facts=_canonical_facts( + message.event_type, + ticket_key=ticket_key, + issue_fields=issue_fields, + comment=comment, + source_ticket_key=message.payload.get("source_ticket_key"), + ), + correlation={ + "provider_event_id": message.event_id, + "transport_event_id": message.event_id, + "workflow_ticket_key": message.ticket_key, }, - correlation={"workflow_ticket_key": message.ticket_key}, ) return AdaptedEvent( source=message.source, @@ -71,3 +88,97 @@ def adapt(self, message: IngressMessage) -> AdaptedEvent: ticket_type=ticket_type, observation=observation, ) + + +def _canonical_facts( + event_type: str, + *, + ticket_key: str, + issue_fields: dict[str, Any], + comment: Any, + source_ticket_key: Any, +) -> dict[str, Any]: + """Build provider-neutral Jira facts shared by webhook and poller paths. + + Jira webhooks commonly contain a full issue, author metadata, changelog + history, and ADF comment objects. The poller intentionally forwards a + smaller webhook-shaped payload. None of those provider details are + needed for command selection: only ticket identity/type/status/labels and + normalized comment text are. Keeping that small stable projection makes + equivalent revisions compare equal in the observation ledger. + """ + issue_type = issue_fields.get("issuetype", {}) + status = issue_fields.get("status", {}) + labels = issue_fields.get("labels", []) + if not isinstance(labels, list | tuple | set): + labels = [] + canonical_fields: dict[str, Any] = {} + if isinstance(issue_type, dict) and issue_type.get("name") is not None: + canonical_fields["issuetype"] = {"name": str(issue_type["name"])} + if isinstance(status, dict) and status.get("name") is not None: + canonical_fields["status"] = {"name": str(status["name"])} + if isinstance(labels, (list, tuple, set)): + canonical_fields["labels"] = sorted({str(label) for label in labels}) + return { + "event_type": event_type, + "issue": { + "key": ticket_key, + "fields": canonical_fields, + }, + # Changelog and comment objects contain transport/provider-specific + # metadata. Keep their historical empty/null compatibility shape; + # command interpretation uses the original ingress payload for + # changelog routing and only needs normalized comment text here. + "changelog": {}, + "comment": None, + "comment_text": _comment_text(comment.get("body", "")) + if isinstance(comment, dict) + else "", + "source_ticket_key": str(source_ticket_key) if source_ticket_key else None, + } + + +def _jira_revision(payload: dict[str, Any]) -> str | None: + """Return a stable revision for a Jira issue observation. + + Comment IDs are immutable provider identities and take precedence over the + issue update timestamp. For issue/label changes, Jira's ``updated`` field + is the only native revision exposed by the issue endpoint. A changelog + fingerprint is used when a webhook has changelog data but no ``updated`` + field. The final ``None`` fallback keeps malformed/legacy payloads + observable without pretending that their UUID delivery ID is orderable. + """ + comment = payload.get("comment") + if isinstance(comment, dict) and comment.get("id") is not None: + return f"comment:{comment['id']}" + + issue = payload.get("issue", {}) + fields = issue.get("fields", {}) if isinstance(issue, dict) else {} + updated = fields.get("updated") + if isinstance(updated, str) and updated: + return f"updated:{updated}" + + changelog = payload.get("changelog") + if isinstance(changelog, dict) and changelog.get("items"): + return stable_identity("jira-changelog", {"items": changelog["items"]}) + return None + + +def _jira_revision_order(payload: dict[str, Any]) -> int | None: + """Convert Jira's update timestamp to comparable monotonic metadata.""" + issue = payload.get("issue", {}) + fields = issue.get("fields", {}) if isinstance(issue, dict) else {} + value = fields.get("updated") + if not isinstance(value, str) or not value: + comment = payload.get("comment") + if isinstance(comment, dict): + value = comment.get("created") or comment.get("updated") + if not isinstance(value, str) or not value: + return None + try: + parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) + except ValueError: + return None + if parsed.tzinfo is None: + return None + return max(0, int(parsed.timestamp() * 1_000_000)) diff --git a/src/forge/orchestrator/worker.py b/src/forge/orchestrator/worker.py index c1694fdc..8f6fa6f6 100644 --- a/src/forge/orchestrator/worker.py +++ b/src/forge/orchestrator/worker.py @@ -561,6 +561,7 @@ async def _process_workflow(self, message: QueueMessage) -> None: ingress = self.event_adapters.adapt(message) observation_decision = await self._observation_ledger().record(ingress.observation) if observation_decision.disposition in { + ObservationDisposition.DUPLICATE, ObservationDisposition.STALE, ObservationDisposition.CONFLICT, }: diff --git a/src/forge/reconciliation/ledger.py b/src/forge/reconciliation/ledger.py index e5e31ecb..c07bb3f7 100644 --- a/src/forge/reconciliation/ledger.py +++ b/src/forge/reconciliation/ledger.py @@ -19,11 +19,28 @@ ) PROTECTED_WORKFLOW_FACTS = { + # Execution position and immutable process identity. "current_node", "workflow_name", "workflow_revision", "workflow_digest", + "workflow_definition_revision", + "workflow_definition_digest", + "workflow_definition", + "workflow_pin_status", + "workflow_state_profile", + "workflow_position", "workflow_transition_count", + "workflow_node_attempts", + # Checkpoint control fields. External providers may report a status, but + # they cannot directly pause, block, retry, or move a workflow checkpoint. + "is_paused", + "is_blocked", + "retry_count", + "last_error", + "node_outcome", + "pending_effects", + "effect_journal", } _RESOURCE_PREFIX = "forge:observations:resource:" _DELIVERY_PREFIX = "forge:observations:delivery:" @@ -74,6 +91,15 @@ async def record(self, observation: Observation) -> ObservationDecision: delivery = observation.delivery_identity duplicate = self._deliveries.get(delivery) if duplicate is not None: + if _revision_metadata_conflicts(duplicate.observation, observation): + decision = _decision( + observation, + ObservationDisposition.CONFLICT, + DriftClass.OPERATOR_REQUIRED, + "provider revision has inconsistent ordering metadata", + ) + self._append(observation, decision) + return decision same_facts = duplicate.observation.facts == observation.facts decision = _decision( observation, @@ -165,17 +191,25 @@ async def record(self, observation: Observation) -> ObservationDecision: ReconciledResource.model_validate_json(current_raw) if current_raw else None ) if prior is not None: - same_facts = prior.observation.facts == observation.facts - decision = _decision( - observation, - ObservationDisposition.DUPLICATE - if same_facts - else ObservationDisposition.CONFLICT, - DriftClass.EXPECTED if same_facts else DriftClass.OPERATOR_REQUIRED, - "provider revision was already observed through an ingress source" - if same_facts - else "same provider revision contains different facts", - ) + if _revision_metadata_conflicts(prior.observation, observation): + decision = _decision( + observation, + ObservationDisposition.CONFLICT, + DriftClass.OPERATOR_REQUIRED, + "provider revision has inconsistent ordering metadata", + ) + else: + same_facts = prior.observation.facts == observation.facts + decision = _decision( + observation, + ObservationDisposition.DUPLICATE + if same_facts + else ObservationDisposition.CONFLICT, + DriftClass.EXPECTED if same_facts else DriftClass.OPERATOR_REQUIRED, + "provider revision was already observed through an ingress source" + if same_facts + else "same provider revision contains different facts", + ) else: disposition, drift, reason = classify_observation( current.latest if current else None, observation @@ -228,6 +262,31 @@ def classify_observation( if current is None: return ObservationDisposition.ACCEPTED, DriftClass.EXPECTED, "first observed revision" if incoming.revision_order is not None and current.revision_order is not None: + # The numeric order and provider token describe the same revision. A + # mismatch is not safely orderable: choosing either value could apply + # facts to the wrong version or move a projection backwards. + if ( + incoming.revision_order == current.revision_order + and incoming.resource_revision is not None + and current.resource_revision is not None + and incoming.resource_revision != current.resource_revision + ): + return ( + ObservationDisposition.CONFLICT, + DriftClass.OPERATOR_REQUIRED, + "same revision order contains different provider revisions", + ) + if ( + incoming.resource_revision is not None + and current.resource_revision is not None + and incoming.resource_revision == current.resource_revision + and incoming.revision_order != current.revision_order + ): + return ( + ObservationDisposition.CONFLICT, + DriftClass.OPERATOR_REQUIRED, + "provider revision has inconsistent ordering metadata", + ) if incoming.revision_order < current.revision_order: return ObservationDisposition.STALE, DriftClass.EXPECTED, "older provider revision" if incoming.revision_order == current.revision_order: @@ -247,6 +306,12 @@ def classify_observation( DriftClass.AUTO_RECONCILABLE, "newer provider revision updates the external projection", ) + if incoming.resource_revision is None and current.resource_revision is None: + return ( + ObservationDisposition.CONFLICT, + DriftClass.OPERATOR_REQUIRED, + "unversioned observations cannot be ordered safely", + ) if incoming.resource_revision == current.resource_revision: if incoming.facts == current.facts: return ObservationDisposition.DUPLICATE, DriftClass.EXPECTED, "same revision and facts" @@ -262,6 +327,31 @@ def classify_observation( ) +def _revision_metadata_conflicts(left: Observation, right: Observation) -> bool: + """Return whether revision token and order contradict one another. + + This check is intentionally separate from delivery identity. Ordering is + optional metadata, so a webhook and poller may legitimately provide only + one representation of the same revision; however, two representations + that assert the same token at different orders (or different tokens at + one order) cannot both be true. + """ + if ( + left.resource_revision is None + or right.resource_revision is None + or left.revision_order is None + or right.revision_order is None + ): + return False + return ( + left.resource_revision == right.resource_revision + and left.revision_order != right.revision_order + ) or ( + left.resource_revision != right.resource_revision + and left.revision_order == right.revision_order + ) + + def _decision( observation: Observation, disposition: ObservationDisposition, diff --git a/tests/contracts/fixtures/observations/github_pull_request_revision.json b/tests/contracts/fixtures/observations/github_pull_request_revision.json new file mode 100644 index 00000000..c82ca6b4 --- /dev/null +++ b/tests/contracts/fixtures/observations/github_pull_request_revision.json @@ -0,0 +1,53 @@ +{ + "webhook": { + "schema_version": "1.0", + "observation_id": "observation:webhook-delivery-7", + "source": "webhook", + "source_system": "github", + "resource": { + "resource_type": "change_request", + "external_id": "acme/api#42", + "namespace": "public" + }, + "resource_revision": "abc123", + "revision_order": 17, + "observed_at": "2026-08-27T10:00:00Z", + "received_at": "2026-08-27T10:00:01Z", + "facts": { + "kind": "cr_updated", + "state": "open", + "head_sha": "abc123" + }, + "correlation": { + "provider_event_id": "delivery-7", + "transport_event_id": "delivery-7", + "repository_id": "acme/api" + }, + "evidence_reference": "source-control-event:delivery-7" + }, + "poller": { + "schema_version": "1.0", + "observation_id": "observation:poller-observation-99", + "source": "poller", + "source_system": "github", + "resource": { + "resource_type": "change_request", + "external_id": "acme/api#42", + "namespace": "public" + }, + "resource_revision": "abc123", + "revision_order": 17, + "observed_at": "2026-08-27T10:00:00Z", + "received_at": "2026-08-27T10:05:00Z", + "facts": { + "kind": "cr_updated", + "state": "open", + "head_sha": "abc123" + }, + "correlation": { + "provider_event_id": "poller-observation-99", + "transport_event_id": "poller-observation-99", + "repository_id": "acme/api" + } + } +} diff --git a/tests/contracts/fixtures/reconciliation/README.md b/tests/contracts/fixtures/reconciliation/README.md new file mode 100644 index 00000000..14fe2ed3 --- /dev/null +++ b/tests/contracts/fixtures/reconciliation/README.md @@ -0,0 +1,10 @@ +# Reconciliation conformance fixtures + +`source_control_sequence.json` is a provider-independent sequence of two +observations for one change request. It is intentionally JSON so the Forge +tests and `forge-poller` tests can consume the exact same revisions without +importing one project's Python package. + +The `provider_event_id`, `resource_revision`, and `revision_order` values are +part of the observation contract. Transport delivery metadata (webhook versus +poller) is not part of the fixture and must not change the resulting command. diff --git a/tests/contracts/fixtures/reconciliation/source_control_sequence.json b/tests/contracts/fixtures/reconciliation/source_control_sequence.json new file mode 100644 index 00000000..6b5c13f6 --- /dev/null +++ b/tests/contracts/fixtures/reconciliation/source_control_sequence.json @@ -0,0 +1,34 @@ +{ + "schema_version": "1.0", + "description": "Provider revisions used by Forge/poller observation conformance tests.", + "source_system": "github", + "resource": { + "resource_type": "change_request", + "external_id": "acme/widgets#17", + "namespace": "default" + }, + "revisions": [ + { + "name": "opened", + "provider_event_id": "github-pr-17-opened", + "resource_revision": "sha-open", + "revision_order": 1, + "facts": { + "kind": "cr_opened", + "change_request_state": "open" + }, + "expected_command": null + }, + { + "name": "merged", + "provider_event_id": "github-pr-17-merged", + "resource_revision": "sha-merge", + "revision_order": 2, + "facts": { + "kind": "cr_merged", + "change_request_state": "merged" + }, + "expected_command": "approve" + } + ] +} diff --git a/tests/contracts/reconciliation/__init__.py b/tests/contracts/reconciliation/__init__.py new file mode 100644 index 00000000..0549b753 --- /dev/null +++ b/tests/contracts/reconciliation/__init__.py @@ -0,0 +1 @@ +"""Cross-ingress reconciliation contract tests.""" diff --git a/tests/contracts/reconciliation/test_convergence.py b/tests/contracts/reconciliation/test_convergence.py new file mode 100644 index 00000000..9af950f1 --- /dev/null +++ b/tests/contracts/reconciliation/test_convergence.py @@ -0,0 +1,235 @@ +"""Conformance tests for webhook/poller observation convergence. + +The fixture contains provider revisions, while this module supplies the two +transport paths. Keeping the provider revision data transport-neutral makes +it possible for forge-poller to run the same fixture once it emits the +versioned Observation envelope. +""" + +from __future__ import annotations + +import json +from dataclasses import replace +from datetime import UTC, datetime +from pathlib import Path +from typing import Any + +import pytest + +from forge.domain import Observation, ObservationSource, WorkflowCommandType +from forge.integrations.source_control.contracts import ( + Actor, + ChangeRequest, + ChangeRequestIdentity, + ChangeRequestState, + EventKind, + NormalizedEvent, + Provider, + RepositoryRef, +) +from forge.integrations.source_control.observations import normalized_event_to_observation +from forge.models.events import EventSource +from forge.orchestrator.event_adapters import ( + CommandDecisionStatus, + create_default_event_adapter_registry, + interpret_event, +) +from forge.queue.models import QueueMessage, normalized_event_to_dict +from forge.reconciliation import InMemoryObservationLedger, ObservationDisposition + +FIXTURE = Path(__file__).parents[1] / "fixtures" / "reconciliation" / "source_control_sequence.json" +NOW = datetime(2026, 8, 27, 12, 0, tzinfo=UTC) +WORKFLOW_STATE: dict[str, Any] = { + "thread_id": "FORGE-42", + "ticket_key": "FORGE-42", + "workflow_name": "feature", + "workflow_definition_revision": 3, + "current_node": "human_review_gate", +} + + +def _fixture() -> dict[str, Any]: + return json.loads(FIXTURE.read_text()) + + +def _event(revision: dict[str, Any]) -> NormalizedEvent: + resource = _fixture()["resource"] + repo_id, number = resource["external_id"].split("#") + return NormalizedEvent( + id=revision["provider_event_id"], + kind=EventKind(revision["facts"]["kind"]), + repo_ref=RepositoryRef( + id=repo_id, + provider=Provider.GITHUB, + connection=resource["namespace"], + namespace=repo_id, + default_branch="main", + change_request_mode="direct", + ), + actor=Actor(login="alice", is_bot=False), + received_at=NOW, + change_request=ChangeRequest( + identity=ChangeRequestIdentity( + connection=resource["namespace"], repository_id=repo_id, native_id=number + ), + url=f"https://github.com/{repo_id}/pull/{number}", + title="Conformance fixture", + body="", + state=ChangeRequestState(revision["facts"]["change_request_state"]), + source_branch="feature", + target_branch="main", + head_sha=revision["resource_revision"], + ), + ) + + +def _observation( + revision: dict[str, Any], source: ObservationSource +) -> tuple[Observation, NormalizedEvent]: + event = _event(revision) + observation = normalized_event_to_observation(event, source=source).model_copy( + update={"revision_order": revision["revision_order"]} + ) + return observation, event + + +def _command( + observation: Observation, + event: NormalizedEvent, + *, + transport_event_id: str | None = None, +) -> tuple[str, str | None, str | None, dict[str, Any]]: + transport_event_id = transport_event_id or event.id + message = QueueMessage( + message_id=f"message-{transport_event_id}", + event_id=transport_event_id, + source=EventSource.SOURCE_CONTROL, + event_type=event.kind.value, + ticket_key="FORGE-42", + normalized_event=normalized_event_to_dict(event), + timestamp=NOW, + ) + adapted = create_default_event_adapter_registry().adapt(message) + # Adaptation is deliberately shared; only the ingress source marker differs. + adapted = replace(adapted, observation=observation) + decision = interpret_event(message, adapted, WORKFLOW_STATE) + return ( + decision.status.value, + decision.command.command_id if decision.command else None, + decision.command.command_type.value if decision.command else None, + decision.command.arguments if decision.command else {}, + ) + + +async def _replay( + revisions: list[dict[str, Any]], sources: list[ObservationSource] +) -> dict[str, Any]: + ledger = InMemoryObservationLedger() + command_decisions: list[tuple[str, str | None, str | None, dict[str, Any]]] = [] + accepted_effects: list[str] = [] + delivery_dispositions: list[str] = [] + for revision, source in zip(revisions, sources, strict=True): + observation, event = _observation(revision, source) + reconciliation = await ledger.record(observation) + delivery_dispositions.append(reconciliation.disposition.value) + if reconciliation.disposition is ObservationDisposition.ACCEPTED: + decision = _command(observation, event) + command_decisions.append(decision) + if decision[0] is CommandDecisionStatus.ACCEPTED.value and decision[2] is not None: + accepted_effects.append(decision[2]) + + latest = await ledger.latest(_observation(revisions[-1], sources[-1])[0]) + assert latest is not None + return { + "latest_revision": latest.latest.resource_revision, + # Delivery history is intentionally omitted: a lost event or a + # duplicate delivery may change that history while the workflow and + # externally visible effect state must converge. + "command_decisions": [ + decision + for decision in command_decisions + if decision[0] == CommandDecisionStatus.ACCEPTED.value + ], + "effects": accepted_effects, + "delivery_dispositions": delivery_dispositions, + } + + +def test_shared_fixture_is_versioned_and_has_provider_revision_identity() -> None: + fixture = _fixture() + assert fixture["schema_version"] == "1.0" + assert fixture["resource"]["resource_type"] == "change_request" + assert all( + revision["provider_event_id"] + and revision["resource_revision"] + and revision["revision_order"] >= 0 + for revision in fixture["revisions"] + ) + + +@pytest.mark.asyncio +async def test_webhook_and_poller_paths_emit_equivalent_observations() -> None: + revisions = _fixture()["revisions"] + webhook, _ = _observation(revisions[1], ObservationSource.WEBHOOK) + poller, _ = _observation(revisions[1], ObservationSource.POLLER) + + assert webhook.source is ObservationSource.WEBHOOK + assert poller.source is ObservationSource.POLLER + assert webhook.delivery_identity == poller.delivery_identity + assert webhook.model_copy(update={"source": poller.source}) == poller + + +def test_transport_delivery_id_does_not_change_command_identity() -> None: + """A poller retry must select the same command as its webhook counterpart.""" + revision = _fixture()["revisions"][1] + webhook, event = _observation(revision, ObservationSource.WEBHOOK) + poller, _ = _observation(revision, ObservationSource.POLLER) + + webhook_command = _command(webhook, event, transport_event_id="github-delivery-17") + poller_command = _command(poller, event, transport_event_id="poller-delivery-17") + + assert webhook.delivery_identity == poller.delivery_identity + assert webhook_command == poller_command + + +@pytest.mark.asyncio +async def test_lost_duplicate_stale_and_reordered_delivery_converges() -> None: + revisions = _fixture()["revisions"] + opened, merged = revisions + expected = await _replay( + revisions, + [ObservationSource.WEBHOOK, ObservationSource.WEBHOOK], + ) + + # The opened event is lost; merge arrives through both paths, is replayed, + # and the old opened revision is delivered after it. + degraded = await _replay( + [merged, merged, opened, merged], + [ + ObservationSource.WEBHOOK, + ObservationSource.POLLER, + ObservationSource.POLLER, + ObservationSource.WEBHOOK, + ], + ) + + assert {key: expected[key] for key in ("latest_revision", "command_decisions", "effects")} == { + key: degraded[key] for key in ("latest_revision", "command_decisions", "effects") + } + assert degraded["delivery_dispositions"] == ["accepted", "duplicate", "stale", "duplicate"] + assert expected["effects"] == [WorkflowCommandType.APPROVE.value] + + +@pytest.mark.asyncio +async def test_fixture_expected_commands_match_both_ingress_sources() -> None: + revisions = _fixture()["revisions"] + for source in (ObservationSource.WEBHOOK, ObservationSource.POLLER): + for revision in revisions: + observation, event = _observation(revision, source) + status, _command_id, command_type, _arguments = _command(observation, event) + expected = revision["expected_command"] + if expected is None: + assert status == CommandDecisionStatus.IGNORED.value + else: + assert status == CommandDecisionStatus.ACCEPTED.value + assert command_type == expected diff --git a/tests/contracts/test_observation_contract.py b/tests/contracts/test_observation_contract.py new file mode 100644 index 00000000..26b829a9 --- /dev/null +++ b/tests/contracts/test_observation_contract.py @@ -0,0 +1,27 @@ +"""Provider-facing conformance fixtures for the Observation v1 contract.""" + +import json +from pathlib import Path + +from forge.domain import Observation, ObservationSource + +FIXTURES = Path(__file__).parent / "fixtures" / "observations" + + +def test_shared_fixture_accepts_both_ingress_sources() -> None: + payload = json.loads((FIXTURES / "github_pull_request_revision.json").read_text()) + webhook = Observation.model_validate_json(json.dumps(payload["webhook"])) + poller = Observation.model_validate_json(json.dumps(payload["poller"])) + + assert webhook.source is ObservationSource.WEBHOOK + assert poller.source is ObservationSource.POLLER + assert webhook.resource == poller.resource + assert webhook.resource_revision == poller.resource_revision + assert webhook.delivery_identity == poller.delivery_identity + + +def test_shared_fixture_is_strict_and_json_round_trips() -> None: + payload = json.loads((FIXTURES / "github_pull_request_revision.json").read_text()) + observation = Observation.model_validate_json(json.dumps(payload["webhook"])) + + assert Observation.model_validate_json(observation.model_dump_json()) == observation diff --git a/tests/unit/integrations/source_control/test_observations.py b/tests/unit/integrations/source_control/test_observations.py index 9bc7f1bf..3095c043 100644 --- a/tests/unit/integrations/source_control/test_observations.py +++ b/tests/unit/integrations/source_control/test_observations.py @@ -10,6 +10,7 @@ NormalizedEvent, Provider, RepositoryRef, + ReviewComment, ) from forge.integrations.source_control.observations import normalized_event_to_observation @@ -63,3 +64,32 @@ def test_poller_and_webhook_use_same_identity_for_same_external_event() -> None: assert webhook.observation_id == polled.observation_id assert polled.source is ObservationSource.POLLER + + +def test_poller_and_webhook_deduplicate_revision_even_with_different_delivery_ids() -> None: + webhook_event = _event() + poller_event = _event() + poller_event.id = "poller-observation-99" + + webhook = normalized_event_to_observation(webhook_event) + polled = normalized_event_to_observation(poller_event, source=ObservationSource.POLLER) + + # The observation records retain their provider delivery identity, while + # the delivery key is derived from the external revision and is shared. + assert webhook.observation_id != polled.observation_id + assert webhook.delivery_identity == polled.delivery_identity + + +def test_event_resources_do_not_share_a_change_request_delivery_key() -> None: + first = _event() + first.kind = EventKind.COMMENT_CREATED + first.comment = ReviewComment(id="comment-1", body="one", author="octocat") + second = _event() + second.kind = EventKind.COMMENT_CREATED + second.comment = ReviewComment(id="comment-2", body="two", author="octocat") + + first_observation = normalized_event_to_observation(first) + second_observation = normalized_event_to_observation(second) + + assert first_observation.resource.resource_type == "comment" + assert first_observation.delivery_identity != second_observation.delivery_identity diff --git a/tests/unit/orchestrator/event_adapters/test_registry.py b/tests/unit/orchestrator/event_adapters/test_registry.py index e16b288e..2da14f8b 100644 --- a/tests/unit/orchestrator/event_adapters/test_registry.py +++ b/tests/unit/orchestrator/event_adapters/test_registry.py @@ -19,6 +19,7 @@ from forge.orchestrator.event_adapters.registry import EventAdapterRegistry from forge.orchestrator.event_adapters.source_control import extract_change_request_url from forge.queue.models import QueueMessage, normalized_event_to_dict +from forge.reconciliation import InMemoryObservationLedger, ObservationDisposition def _message( @@ -89,6 +90,116 @@ def test_default_registry_adapts_jira_without_provider_clients() -> None: assert adapted.observation.facts["event_type"] == "updated" +def test_jira_issue_revision_is_shared_when_delivery_ids_differ() -> None: + payload = { + "issue": { + "key": "FORGE-42", + "fields": { + "issuetype": {"name": "Feature"}, + "updated": "2026-08-27T10:00:00.000+0000", + }, + }, + "changelog": {"items": [{"field": "labels", "toString": "forge:managed"}]}, + } + webhook = _message(source=EventSource.JIRA, payload=payload) + poller = _message(source=EventSource.JIRA, payload=payload) + poller.event_id = "poller-delivery-42" + + adapter = JiraEventAdapter() + webhook_observation = adapter.adapt(webhook).observation + poller_observation = adapter.adapt(poller).observation + + assert webhook_observation.resource_revision == "updated:2026-08-27T10:00:00.000+0000" + assert webhook_observation.revision_order is not None + assert webhook_observation.delivery_identity == poller_observation.delivery_identity + + +def test_jira_comment_id_wins_over_issue_revision_for_cross_source_replay() -> None: + webhook_payload = { + "issue": { + "key": "FORGE-42", + "fields": { + "issuetype": {"name": "Feature"}, + "updated": "2026-08-27T10:00:00.000+0000", + }, + }, + "comment": {"id": "10042", "body": "Please revise"}, + } + poller_payload = { + **webhook_payload, + "issue": { + **webhook_payload["issue"], + "fields": { + **webhook_payload["issue"]["fields"], + "updated": "2026-08-27T10:01:00.000+0000", + }, + }, + } + adapter = JiraEventAdapter() + webhook = adapter.adapt(_message(source=EventSource.JIRA, payload=webhook_payload)).observation + poller = adapter.adapt(_message(source=EventSource.JIRA, payload=poller_payload)).observation + + assert webhook.resource_revision == "comment:10042" + assert webhook.delivery_identity == poller.delivery_identity + + +def test_jira_comment_created_timestamp_orders_comments_without_issue_updated() -> None: + payload = { + "issue": {"key": "FORGE-42", "fields": {"issuetype": {"name": "Feature"}}}, + "comment": {"id": "10042", "created": "2026-08-27T10:01:00.000+0000"}, + } + + adapted = JiraEventAdapter().adapt(_message(source=EventSource.JIRA, payload=payload)) + + assert adapted.observation.resource_revision == "comment:10042" + assert adapted.observation.revision_order is not None + + +@pytest.mark.asyncio +async def test_rich_webhook_and_minimal_poller_facts_deduplicate_same_issue_revision() -> None: + rich_payload = { + "webhookEvent": "jira:issue_updated", + "issue": { + "id": "10042", + "key": "FORGE-42", + "fields": { + "issuetype": {"name": "Feature", "id": "10001"}, + "status": {"name": "In Progress", "id": "3"}, + "labels": ["forge:managed", "forge:pending"], + "summary": "A richer provider issue", + "description": {"type": "doc", "content": []}, + "updated": "2026-08-27T10:00:00.000+0000", + }, + }, + "changelog": {"id": "history-1", "items": [{"field": "labels"}]}, + "user": {"accountId": "provider-user", "displayName": "Provider"}, + } + minimal_payload = { + "webhookEvent": "jira:issue_updated", + "issue": { + "key": "FORGE-42", + "fields": { + "issuetype": {"name": "Feature"}, + "status": {"name": "In Progress"}, + "labels": ["forge:managed", "forge:pending"], + "updated": "2026-08-27T10:00:00.000+0000", + }, + }, + } + adapter = JiraEventAdapter() + webhook_message = _message(source=EventSource.JIRA, payload=rich_payload) + poller_message = _message(source=EventSource.JIRA, payload=minimal_payload) + poller_message.event_id = "poller-delivery-42" + webhook = adapter.adapt(webhook_message).observation + poller = adapter.adapt(poller_message).observation + + assert webhook.facts == poller.facts + assert webhook.delivery_identity == poller.delivery_identity + ledger = InMemoryObservationLedger() + assert (await ledger.record(webhook)).disposition is ObservationDisposition.ACCEPTED + assert (await ledger.record(poller)).disposition is ObservationDisposition.DUPLICATE + + def test_child_jira_event_rerouted_to_parent_does_not_start_child_workflow() -> None: message = _message( source=EventSource.JIRA, diff --git a/tests/unit/orchestrator/test_reconciliation_worker.py b/tests/unit/orchestrator/test_reconciliation_worker.py new file mode 100644 index 00000000..0431ca3f --- /dev/null +++ b/tests/unit/orchestrator/test_reconciliation_worker.py @@ -0,0 +1,47 @@ +"""Worker ingress tests for source-independent reconciliation.""" + +from unittest.mock import AsyncMock, patch + +import pytest + +from forge.models.events import EventSource +from forge.orchestrator.event_adapters import create_default_event_adapter_registry +from forge.orchestrator.worker import OrchestratorWorker +from forge.queue.models import QueueMessage +from forge.reconciliation import InMemoryObservationLedger, ObservationDisposition + + +def _jira_message() -> QueueMessage: + return QueueMessage( + message_id="message-1", + event_id="provider-event-1", + source=EventSource.JIRA, + event_type="jira:issue_updated", + ticket_key="FORGE-42", + payload={ + "issue": { + "key": "FORGE-42", + "fields": {"issuetype": {"name": "Feature"}}, + } + }, + ) + + +@pytest.mark.asyncio +async def test_duplicate_observation_does_not_reinterpret_or_start_workflow() -> None: + message = _jira_message() + adapters = create_default_event_adapter_registry() + observation = adapters.adapt(message).observation + ledger = InMemoryObservationLedger() + assert (await ledger.record(observation)).disposition is ObservationDisposition.ACCEPTED + + worker = OrchestratorWorker(consumer_name="test-worker", observation_ledger=ledger) + with ( + patch("forge.orchestrator.worker.ensure_skills", new=AsyncMock()), + patch("forge.orchestrator.worker.interpret_event") as interpret, + patch.object(worker, "_invoke_workflow", new=AsyncMock()) as invoke, + ): + await worker._process_workflow(message) + + interpret.assert_not_called() + invoke.assert_not_awaited() diff --git a/tests/unit/reconciliation/test_ledger.py b/tests/unit/reconciliation/test_ledger.py index 01304535..7f71e899 100644 --- a/tests/unit/reconciliation/test_ledger.py +++ b/tests/unit/reconciliation/test_ledger.py @@ -32,6 +32,25 @@ def observation( ) +def unversioned_observation( + source: ObservationSource, + observation_id: str, + *, + event_id: str | None = None, +) -> Observation: + now = datetime.now(UTC) + return Observation( + observation_id=observation_id, + source=source, + source_system="jira", + resource=ResourceIdentity(resource_type="issue", external_id="FORGE-17"), + observed_at=now, + received_at=now, + facts={"event_type": "issue_updated"}, + correlation={"provider_event_id": event_id} if event_id else {}, + ) + + @pytest.mark.asyncio async def test_webhook_and_poller_delivery_share_identity_and_deduplicate() -> None: ledger = InMemoryObservationLedger() @@ -95,3 +114,63 @@ async def test_external_observation_cannot_overwrite_workflow_position() -> None assert decision.disposition is ObservationDisposition.CONFLICT assert decision.drift is DriftClass.POLICY_BLOCKING assert await ledger.latest(incoming) is None + + +@pytest.mark.asyncio +async def test_unversioned_events_do_not_collapse_into_one_delivery() -> None: + ledger = InMemoryObservationLedger() + first = unversioned_observation(ObservationSource.WEBHOOK, "event-1") + second = unversioned_observation(ObservationSource.WEBHOOK, "event-2") + + first_decision = await ledger.record(first) + second_decision = await ledger.record(second) + + assert first.delivery_identity != second.delivery_identity + assert first_decision.disposition is ObservationDisposition.ACCEPTED + # Without an order or revision token, a second event cannot safely replace + # the first projection; it is retained as an operator-visible conflict. + assert second_decision.disposition is ObservationDisposition.CONFLICT + assert second_decision.drift is DriftClass.OPERATOR_REQUIRED + + +@pytest.mark.asyncio +async def test_unversioned_provider_event_id_is_shared_across_sources() -> None: + ledger = InMemoryObservationLedger() + webhook = unversioned_observation( + ObservationSource.WEBHOOK, "webhook-delivery", event_id="provider-event-7" + ) + polled = unversioned_observation( + ObservationSource.POLLER, "poll-delivery", event_id="provider-event-7" + ) + + assert webhook.delivery_identity == polled.delivery_identity + assert (await ledger.record(webhook)).disposition is ObservationDisposition.ACCEPTED + assert (await ledger.record(polled)).disposition is ObservationDisposition.DUPLICATE + + +@pytest.mark.asyncio +async def test_revision_order_and_token_mismatch_is_operator_conflict() -> None: + ledger = InMemoryObservationLedger() + await ledger.record(observation(ObservationSource.WEBHOOK, 4)) + incoming = observation(ObservationSource.POLLER, 4).model_copy( + update={"resource_revision": "different-revision"} + ) + + decision = await ledger.record(incoming) + + assert decision.disposition is ObservationDisposition.CONFLICT + assert decision.drift is DriftClass.OPERATOR_REQUIRED + + +@pytest.mark.asyncio +async def test_same_token_with_different_order_is_operator_conflict() -> None: + ledger = InMemoryObservationLedger() + await ledger.record(observation(ObservationSource.WEBHOOK, 4)) + incoming = observation(ObservationSource.POLLER, 5).model_copy( + update={"resource_revision": "revision-4"} + ) + + decision = await ledger.record(incoming) + + assert decision.disposition is ObservationDisposition.CONFLICT + assert decision.drift is DriftClass.OPERATOR_REQUIRED diff --git a/zensical.toml b/zensical.toml index 39e20f4f..004fd0ee 100644 --- a/zensical.toml +++ b/zensical.toml @@ -28,6 +28,8 @@ nav = [ {"Overview" = "architecture/index.md"}, {"System & Components" = "architecture/overview.md"}, {"Internals" = "architecture/internals.md"}, + {"Phase 6 Observation Contract" = "architecture/phase-6-observation-contract.md"}, + {"Phase 6 Reconciliation" = "architecture/phase-6-reconciliation-contract.md"}, {"Reference" = "architecture/reference.md"}, ]}, {"Local Setup" = "dev/setup.md"},