diff --git a/ace/application/__init__.py b/ace/application/__init__.py index cc34279..21396e8 100644 --- a/ace/application/__init__.py +++ b/ace/application/__init__.py @@ -310,9 +310,11 @@ IntelligenceBuildEffect, IntelligenceBuildExecutor, IntelligenceBuildHostServices, + IntelligenceBuildRecordedSourcePort, IntelligenceBuildResourcePagePort, IntelligenceBuildStartV1, ProductScopedImmutableRecordStore, + RecordedSourceReferenceV1, ) from ace.application.intelligence_builder import ( ConnectionAgent, @@ -478,6 +480,13 @@ PersonalIntelligenceOwnershipService, PersonalIntelligenceOwnershipStore, ) +from ace.application.recorded_source_admission import ( + CoreRecordedSourceAdmissionService, + RecordedSourceAcquisitionReceiptV1Alpha1, + RecordedSourceAdmission, + RecordedSourceAdmissionError, + RecordedSourceMaterialV1Alpha1, +) from ace.application.supersession_impact import ( SupersessionImpactAdmission, SupersessionImpactAdmissionError, @@ -785,10 +794,17 @@ "IntelligenceBuildExecutor", "IntelligenceBuildHostServices", "IntelligenceBuildResourcePagePort", + "IntelligenceBuildRecordedSourcePort", "IntelligenceBuildStartV1", "AuthorizedIntelligenceBuild", "ProductScopedImmutableRecordStore", + "CoreRecordedSourceAdmissionService", "REQUIRED_INTELLIGENCE_BUILD_EFFECTS", + "RecordedSourceAcquisitionReceiptV1Alpha1", + "RecordedSourceAdmission", + "RecordedSourceAdmissionError", + "RecordedSourceMaterialV1Alpha1", + "RecordedSourceReferenceV1", "IntelligenceActivationPlanV1Alpha2", "IntelligenceAgent", "IntelligenceAgentAttributionError", diff --git a/ace/application/intelligence_build_execution.py b/ace/application/intelligence_build_execution.py index 3b1c13e..6123a1b 100644 --- a/ace/application/intelligence_build_execution.py +++ b/ace/application/intelligence_build_execution.py @@ -9,9 +9,9 @@ from dataclasses import dataclass from datetime import datetime -from typing import Literal, Protocol +from typing import TYPE_CHECKING, Literal, Protocol -from pydantic import BaseModel, ConfigDict, Field, field_validator +from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator from ace.application.intelligence_resource_plane import ( IntelligenceResourceKind, @@ -26,6 +26,10 @@ ) from ace.core.runtime_use import AuthorityUseReceiptV1Alpha1 from ace.core.state import CoreAuthorityResolver, ResolvedApprovalReceiptV1 +from ace.intelligence.contracts.common import validate_digest, validate_reference, validate_slug + +if TYPE_CHECKING: + from ace.application.recorded_source_admission import RecordedSourceAdmission, RecordedSourceMaterialV1Alpha1 IntelligenceBuildEffect = Literal[ "connect_sources", @@ -41,6 +45,31 @@ ) +class RecordedSourceReferenceV1(BaseModel): + """Exact recorded material selected in the reviewed Builder request.""" + + model_config = ConfigDict(extra="forbid", frozen=True, strict=True, revalidate_instances="always") + + source_group_id: str + material_id: str + material_digest: str + + @field_validator("source_group_id") + @classmethod + def _validate_group(cls, value: str) -> str: + return validate_slug(value, name="source_group_id") + + @field_validator("material_id") + @classmethod + def _validate_id(cls, value: str) -> str: + return validate_reference(value, name="material_id") + + @field_validator("material_digest") + @classmethod + def _validate_digest(cls, value: str) -> str: + return validate_digest(value) + + class IntelligenceBuildStartV1(BaseModel): """One reviewed Atrium plan submitted for governed execution.""" @@ -55,6 +84,7 @@ class IntelligenceBuildStartV1(BaseModel): subject: str = Field(min_length=8, max_length=2_000) outcome_id: str = Field(min_length=1, max_length=240) source_group_ids: tuple[str, ...] = Field(default_factory=tuple, max_length=64) + recorded_source_refs: tuple[RecordedSourceReferenceV1, ...] = Field(default_factory=tuple, max_length=64) cadence_id: str = Field(min_length=1, max_length=240) approved_effects: tuple[IntelligenceBuildEffect, ...] requested_at: datetime @@ -66,6 +96,17 @@ def _unique_source_groups(cls, value: tuple[str, ...]) -> tuple[str, ...]: raise ValueError("source_group_ids must be unique") return value + @field_validator("recorded_source_refs") + @classmethod + def _exact_recorded_sources( + cls, + value: tuple[RecordedSourceReferenceV1, ...], + ) -> tuple[RecordedSourceReferenceV1, ...]: + keys = [(item.source_group_id, item.material_id) for item in value] + if len(keys) != len(set(keys)): + raise ValueError("recorded_source_refs must name each exact recorded material once") + return tuple(sorted(value, key=lambda item: (item.source_group_id, item.material_id))) + @field_validator("approved_effects") @classmethod def _exact_bounded_effects(cls, value: tuple[IntelligenceBuildEffect, ...]) -> tuple[IntelligenceBuildEffect, ...]: @@ -73,6 +114,13 @@ def _exact_bounded_effects(cls, value: tuple[IntelligenceBuildEffect, ...]) -> t raise ValueError("approved_effects must preserve the exact bounded onboarding effect sequence") return value + @model_validator(mode="after") + def _recorded_sources_are_in_reviewed_groups(self): + selected = set(self.source_group_ids) + if any(item.source_group_id not in selected for item in self.recorded_source_refs): + raise ValueError("every recorded source reference must belong to a reviewed source group") + return self + @dataclass(frozen=True, slots=True) class AuthorizedIntelligenceBuild: @@ -159,6 +207,15 @@ async def query( ) -> IntelligenceResourcePageV1Alpha1: ... +class IntelligenceBuildRecordedSourcePort(Protocol): + """Narrow host capability for the exact recorded material set in one build.""" + + async def admit( + self, + materials: tuple["RecordedSourceMaterialV1Alpha1", ...], + ) -> "RecordedSourceAdmission": ... + + @dataclass(frozen=True, slots=True) class IntelligenceBuildHostServices: """Invocation-scoped capabilities Core grants to one trusted executor.""" @@ -166,6 +223,7 @@ class IntelligenceBuildHostServices: records: ImmutableRecordStore resources: IntelligenceBuildResourcePagePort activation_authority: CoreAuthorityResolver + recorded_sources: IntelligenceBuildRecordedSourcePort | None = None class IntelligenceBuildExecutor(Protocol): @@ -182,7 +240,9 @@ async def start( "IntelligenceBuildExecutor", "IntelligenceBuildHostServices", "IntelligenceBuildResourcePagePort", + "IntelligenceBuildRecordedSourcePort", "IntelligenceBuildStartV1", "ProductScopedImmutableRecordStore", + "RecordedSourceReferenceV1", "REQUIRED_INTELLIGENCE_BUILD_EFFECTS", ] diff --git a/ace/application/recorded_source_admission.py b/ace/application/recorded_source_admission.py new file mode 100644 index 0000000..0269064 --- /dev/null +++ b/ace/application/recorded_source_admission.py @@ -0,0 +1,635 @@ +"""Governed admission of explicitly reviewed recorded source material. + +This boundary is intentionally different from LIVE source ingress. It performs +no network request and makes no freshness claim. Core binds the exact recorded +bytes to the already-authorized ``connect_sources`` build effect, the current +authority-grant head, and one exact committed Domain Activation head before it +atomically persists a recorded-replay receipt, canonical source snapshot, and +canonical PREPARED Observation. +""" + +from __future__ import annotations + +import hashlib +import json +from dataclasses import dataclass +from datetime import UTC, datetime +from typing import Any, Literal, Self + +from pydantic import ConfigDict, Field, field_validator, model_validator + +from ace.application.domain_activation import ( + DOMAIN_ACTIVATION_STATE_KIND, + CommittedActivationBinding, + CommittedDomainActivation, + bind_committed_activation, +) +from ace.application.intelligence_build_execution import ( + AuthorizedIntelligenceBuild, + IntelligenceBuildRecordedSourcePort, +) +from ace.application.intelligence_ledger import PREPARED_RECORD_SPACE +from ace.core.contracts import FrozenContract, canonical_hash, canonical_json +from ace.core.records import ( + AppendOnlyTransactionReceiptV1, + AppendOnlyTransactionRequestV1, + ImmutableRecordStore, + ImmutableRecordV1, +) +from ace.core.runtime_use import AuthorityUseReceiptV1Alpha1 +from ace.core.source import CanonicalSourceSnapshotV1Alpha1, SourceAcquisitionMode +from ace.core.state import GovernedStateHeadPreconditionV1Alpha1 +from ace.intelligence.contracts.common import ( + validate_digest, + validate_reference, + validate_slug, +) +from ace.intelligence.contracts.ledger import IntelligenceRecordKind +from ace.intelligence.contracts.resources import ( + ActivationRevisionReferenceV1Alpha1, + EvidenceAcquisitionMode, + IntelligenceResourceMode, + ObservationV1Alpha1, +) +from ace.intelligence.contracts.source_mapping import ResolvedSubjectBindingV1Alpha1 +from ace.intelligence.source_mapping import interpret_prepared_source_mapping + +RECORDED_SOURCE_MATERIAL_VERSION = "ace.application.recorded-source-material/v1alpha1" +RECORDED_SOURCE_ACQUISITION_RECEIPT_VERSION = "ace.application.recorded-source-acquisition-receipt/v1alpha1" +RECORDED_SOURCE_RECORD_KIND = "recorded_source_acquisition" +SOURCE_SNAPSHOT_RECORD_KIND = "source_snapshot" +CONNECT_SOURCES_EFFECT = "connect_sources" +INTELLIGENCE_BUILD_OPERATION = "start_intelligence_build" +INTELLIGENCE_BUILD_AUTHORITY = "intelligence_build" + + +class RecordedSourceAdmissionError(RuntimeError): + """Recorded material failed exact governed admission or replay.""" + + +class _StrictFrozenContract(FrozenContract): + model_config = ConfigDict( + extra="forbid", + frozen=True, + strict=True, + revalidate_instances="always", + validate_default=True, + allow_inf_nan=False, + ) + + +def _aware(value: datetime, *, name: str) -> datetime: + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError(f"{name} must include a timezone") + return value.astimezone(UTC) + + +def _derive(instance: _StrictFrozenContract, *, id_field: str, digest_field: str, prefix: str) -> None: + material = instance.model_dump(mode="json", exclude={id_field, digest_field}) + digest = canonical_hash(material) + expected_id = f"{prefix}:{digest[:32]}" + expected_digest = f"sha256:{digest}" + supplied_id = getattr(instance, id_field) + supplied_digest = getattr(instance, digest_field) + if supplied_id is not None and supplied_id != expected_id: + raise ValueError(f"{id_field} does not match exact contract material") + if supplied_digest is not None and supplied_digest != expected_digest: + raise ValueError(f"{digest_field} does not match exact contract material") + object.__setattr__(instance, id_field, expected_id) + object.__setattr__(instance, digest_field, expected_digest) + + +class RecordedSourceMaterialV1Alpha1(_StrictFrozenContract): + """Exact reviewed bytes and metadata; never proof of a fresh acquisition.""" + + contract: Literal["ace.application.recorded-source-material/v1alpha1"] = RECORDED_SOURCE_MATERIAL_VERSION + source_group_id: str + mapping_id: str + subject_binding: ResolvedSubjectBindingV1Alpha1 + source_definition_ref: str + source_type_ref: str + source_uri: str = Field(min_length=3, max_length=2_048) + captured_payload_json: str = Field(min_length=1, max_length=1_000_000) + captured_payload_digest: str + source_published_at: datetime | None = None + event_effective_at: datetime | None = None + observed_at: datetime + locator: str | None = Field(default=None, min_length=1, max_length=1_000) + material_id: str | None = None + material_digest: str | None = None + + @field_validator("source_group_id", "mapping_id") + @classmethod + def validate_ids(cls, value: str, info) -> str: + return validate_slug(value, name=info.field_name) + + @field_validator("source_definition_ref", "material_id") + @classmethod + def validate_refs(cls, value: str | None, info) -> str | None: + return validate_reference(value, name=info.field_name) if value is not None else None + + @field_validator("captured_payload_digest", "material_digest") + @classmethod + def validate_digests(cls, value: str | None) -> str | None: + return validate_digest(value) if value is not None else None + + @field_validator("captured_payload_json") + @classmethod + def require_canonical_payload(cls, value: str) -> str: + try: + parsed = json.loads(value, object_pairs_hook=lambda pairs: _unique_object(pairs)) + normalized = canonical_json(parsed) + except (TypeError, UnicodeError, ValueError, RecursionError) as exc: + raise ValueError("captured_payload_json must be bounded canonical JSON") from exc + if normalized != value: + raise ValueError("captured_payload_json must already use exact canonical JSON bytes") + return value + + @field_validator("source_published_at", "event_effective_at", "observed_at") + @classmethod + def normalize_times(cls, value: datetime | None, info) -> datetime | None: + return _aware(value, name=info.field_name) if value is not None else None + + @model_validator(mode="after") + def validate_exact_recorded_material(self) -> Self: + # CanonicalSourceSnapshot performs the strict JSON parse and + # canonicalization. Reject an unreviewed byte digest here before Core + # assigns any acquisition or Observation identity. + expected = "sha256:" + hashlib.sha256(self.captured_payload_json.encode("utf-8")).hexdigest() + if self.captured_payload_digest != expected: + raise ValueError("captured_payload_digest does not match exact submitted bytes") + if self.source_published_at is not None and self.source_published_at > self.observed_at: + raise ValueError("source_published_at cannot follow observed_at") + if self.event_effective_at is not None and self.event_effective_at > self.observed_at: + raise ValueError("event_effective_at cannot follow observed_at") + _derive(self, id_field="material_id", digest_field="material_digest", prefix="recorded_source_material") + return self + + +class RecordedSourceAcquisitionReceiptV1Alpha1(_StrictFrozenContract): + """Proof of governed recorded-replay admission, never a network-capture claim.""" + + contract: Literal["ace.application.recorded-source-acquisition-receipt/v1alpha1"] = ( + RECORDED_SOURCE_ACQUISITION_RECEIPT_VERSION + ) + disposition: Literal["recorded_material_admitted"] = "recorded_material_admitted" + acquisition_mode: Literal[EvidenceAcquisitionMode.RECORDED_REPLAY] = EvidenceAcquisitionMode.RECORDED_REPLAY + product_id: str + actor_ref: str + build_id: str + build_request_digest: str + effect: Literal["connect_sources"] = CONNECT_SOURCES_EFFECT + build_authority_use: AuthorityUseReceiptV1Alpha1 + activation_revision: ActivationRevisionReferenceV1Alpha1 + activation_head_precondition: GovernedStateHeadPreconditionV1Alpha1 + recorded_material_id: str + recorded_material_digest: str + source_group_id: str + source_definition_ref: str + source_type_ref: str + source_uri: str + captured_payload_digest: str + source_published_at: datetime | None = None + event_effective_at: datetime | None = None + observed_at: datetime + admitted_at: datetime + network_capture_performed: Literal[False] = False + freshness_verified: Literal[False] = False + receipt_id: str | None = None + receipt_digest: str | None = None + + @field_validator( + "product_id", + "actor_ref", + "build_id", + "recorded_material_id", + "source_definition_ref", + "receipt_id", + ) + @classmethod + def validate_refs(cls, value: str | None, info) -> str | None: + return validate_reference(value, name=info.field_name) if value is not None else None + + @field_validator("source_group_id") + @classmethod + def validate_group(cls, value: str) -> str: + return validate_slug(value, name="source_group_id") + + @field_validator( + "build_request_digest", + "recorded_material_digest", + "captured_payload_digest", + "receipt_digest", + ) + @classmethod + def validate_digests(cls, value: str | None) -> str | None: + return validate_digest(value) if value is not None else None + + @field_validator("source_published_at", "event_effective_at", "observed_at", "admitted_at") + @classmethod + def normalize_times(cls, value: datetime | None, info) -> datetime | None: + return _aware(value, name=info.field_name) if value is not None else None + + @model_validator(mode="after") + def validate_governed_recorded_admission(self) -> Self: + authority = self.build_authority_use + if ( + authority.product_id != self.product_id + or authority.actor_ref != self.actor_ref + or authority.use_subject_ref != self.build_id + or authority.use_subject_digest != self.build_request_digest + or authority.operation != INTELLIGENCE_BUILD_OPERATION + or authority.authority != INTELLIGENCE_BUILD_AUTHORITY + or authority.evaluated_at != self.admitted_at + ): + raise ValueError("build authority use does not bind the exact recorded admission") + if ( + self.activation_head_precondition.state_kind != DOMAIN_ACTIVATION_STATE_KIND + or self.activation_head_precondition.product_id != self.product_id + or self.activation_head_precondition.state_id != self.activation_revision.activation_id + or self.activation_head_precondition.sequence != self.activation_revision.revision + or self.activation_head_precondition.revision_id != self.activation_revision.revision_id + ): + raise ValueError("recorded admission does not bind the exact activation head") + if self.activation_revision.product_id != self.product_id: + raise ValueError("activation revision crossed recorded admission product scope") + if self.observed_at > self.admitted_at: + raise ValueError("recorded material cannot be admitted before its stated observation time") + _derive(self, id_field="receipt_id", digest_field="receipt_digest", prefix="recorded_source_acquisition") + return self + + @property + def live_acquisition(self) -> Literal[False]: + return False + + @property + def reusable_authority(self) -> Literal[False]: + return False + + +@dataclass(frozen=True, slots=True) +class RecordedSourceAdmission: + """Exact reopened recorded receipts, snapshots, Observations, and append receipt.""" + + acquisition_receipts: tuple[RecordedSourceAcquisitionReceiptV1Alpha1, ...] + source_snapshots: tuple[CanonicalSourceSnapshotV1Alpha1, ...] + observations: tuple[ObservationV1Alpha1, ...] + transaction_receipt: AppendOnlyTransactionReceiptV1 + replayed: bool + + @property + def live_acquisition(self) -> Literal[False]: + return False + + +def _activation_head(binding: CommittedActivationBinding) -> GovernedStateHeadPreconditionV1Alpha1: + revision = binding.prepared_binding.revision + receipt = binding.commit_receipt + if revision.activation_id is None or revision.revision_id is None or receipt.receipt_id is None: + raise RecordedSourceAdmissionError("committed activation is missing exact head coordinates") + return GovernedStateHeadPreconditionV1Alpha1( + state_kind=DOMAIN_ACTIVATION_STATE_KIND, + product_id=revision.spec.product_id, + state_id=revision.activation_id, + sequence=revision.revision, + revision_id=revision.revision_id, + commit_receipt_id=receipt.receipt_id, + ) + + +def _unique_object(pairs: list[tuple[str, Any]]) -> dict[str, Any]: + result: dict[str, Any] = {} + for key, value in pairs: + if key in result: + raise ValueError(f"duplicate JSON object key: {key}") + result[key] = value + return result + + +def _record( + payload, + *, + product_id: str, + kind: str, + key: str, + as_of: datetime, + available_at: datetime, + order: int, +) -> ImmutableRecordV1: + return ImmutableRecordV1( + product_id=product_id, + record_space=PREPARED_RECORD_SPACE, + record_kind=kind, + record_key=key, + payload_contract=payload.contract, + payload=payload.model_dump(mode="python"), + as_of=as_of, + available_at=available_at, + processing_order=order, + ) + + +class CoreRecordedSourceAdmissionService(IntelligenceBuildRecordedSourcePort): + """Bind reviewed replay bytes to one authorized build and committed activation.""" + + def __init__( + self, + *, + build: AuthorizedIntelligenceBuild, + binding: CommittedActivationBinding, + store: ImmutableRecordStore, + ) -> None: + self.build = build + self.binding = self._validate_binding(binding) + self.store = store + self._validate_build() + + @staticmethod + def _validate_binding(binding: CommittedActivationBinding) -> CommittedActivationBinding: + try: + exact = bind_committed_activation( + pack=binding.prepared_binding.pack, + committed=CommittedDomainActivation( + revision=binding.prepared_binding.revision, + commit_receipt=binding.commit_receipt, + ), + ) + except Exception: + raise RecordedSourceAdmissionError("committed activation binding failed exact revalidation") from None + if exact != binding: + raise RecordedSourceAdmissionError("committed activation binding changed during revalidation") + return exact + + def _validate_build(self) -> None: + authority = self.build.authority_use + if ( + authority.product_id != self.build.product_id + or authority.actor_ref != self.build.actor_ref + or authority.use_subject_ref != self.build.build_id + or authority.use_subject_digest != self.build.request_digest + or authority.operation != INTELLIGENCE_BUILD_OPERATION + or authority.authority != INTELLIGENCE_BUILD_AUTHORITY + or authority.grant_ref != self.build.request.authority_grant_ref + or CONNECT_SOURCES_EFFECT not in self.build.request.approved_effects + ): + raise RecordedSourceAdmissionError("authorized build does not cover exact recorded source admission") + if self.binding.prepared_binding.reference.product_id != self.build.product_id: + raise RecordedSourceAdmissionError("committed activation crossed the authorized build product") + + def _materials( + self, + materials: tuple[RecordedSourceMaterialV1Alpha1, ...], + ) -> tuple[RecordedSourceMaterialV1Alpha1, ...]: + if not materials: + raise RecordedSourceAdmissionError("recorded source admission requires exact reviewed material") + try: + exact = tuple( + RecordedSourceMaterialV1Alpha1.model_validate(material.model_dump(mode="python")) + for material in materials + ) + except (AttributeError, TypeError, ValueError) as exc: + raise RecordedSourceAdmissionError("recorded source material failed exact revalidation") from exc + ordered = tuple(sorted(exact, key=lambda item: (item.source_group_id, str(item.material_id)))) + actual_refs = tuple( + (item.source_group_id, str(item.material_id), str(item.material_digest)) for item in ordered + ) + authorized_refs = tuple( + (item.source_group_id, item.material_id, item.material_digest) + for item in self.build.request.recorded_source_refs + ) + if actual_refs != authorized_refs: + raise RecordedSourceAdmissionError( + "recorded source set does not exactly match the reviewed group, material IDs, and digests" + ) + if any(item.subject_binding.product_id != self.build.product_id for item in ordered): + raise RecordedSourceAdmissionError("recorded material crossed the authorized build product") + if any(item.subject_binding.activation_revision != self.binding.prepared_binding.reference for item in ordered): + raise RecordedSourceAdmissionError("recorded material does not bind the exact committed activation") + return ordered + + def _transaction_key(self) -> str: + coordinates = tuple( + (item.source_group_id, item.material_id, item.material_digest) + for item in self.build.request.recorded_source_refs + ) + return f"recorded_source_admission:{canonical_hash([self.build.build_id, coordinates])[:32]}" + + async def admit(self, materials: tuple[RecordedSourceMaterialV1Alpha1, ...]) -> RecordedSourceAdmission: + exact = self._materials(materials) + transaction_key = self._transaction_key() + replay = await self._replay(transaction_key=transaction_key, expected=exact) + if replay is not None: + return replay + + admitted_at = self.build.authority_use.evaluated_at + activation_head = _activation_head(self.binding) + acquisitions: list[RecordedSourceAcquisitionReceiptV1Alpha1] = [] + snapshots: list[CanonicalSourceSnapshotV1Alpha1] = [] + observations: list[ObservationV1Alpha1] = [] + for material in exact: + acquisition = RecordedSourceAcquisitionReceiptV1Alpha1( + product_id=self.build.product_id, + actor_ref=self.build.actor_ref, + build_id=self.build.build_id, + build_request_digest=self.build.request_digest, + build_authority_use=self.build.authority_use, + activation_revision=self.binding.prepared_binding.reference, + activation_head_precondition=activation_head, + recorded_material_id=str(material.material_id), + recorded_material_digest=str(material.material_digest), + source_group_id=material.source_group_id, + source_definition_ref=material.source_definition_ref, + source_type_ref=material.source_type_ref, + source_uri=material.source_uri, + captured_payload_digest=material.captured_payload_digest, + source_published_at=material.source_published_at, + event_effective_at=material.event_effective_at, + observed_at=material.observed_at, + admitted_at=admitted_at, + ) + snapshot = CanonicalSourceSnapshotV1Alpha1( + source_definition_ref=material.source_definition_ref, + source_type_ref=material.source_type_ref, + source_uri=material.source_uri, + captured_payload_json=material.captured_payload_json, + captured_payload_digest=material.captured_payload_digest, + source_published_at=material.source_published_at, + event_effective_at=material.event_effective_at, + observed_at=material.observed_at, + ingested_at=admitted_at, + locator=material.locator, + acquisition_mode=SourceAcquisitionMode.RECORDED_REPLAY, + acquisition_receipt_ref=str(acquisition.receipt_id), + acquisition_receipt_digest=str(acquisition.receipt_digest), + ) + try: + mapped = interpret_prepared_source_mapping( + binding=self.binding.prepared_binding, + mapping_id=material.mapping_id, + source_snapshot=snapshot, + subject_binding=material.subject_binding, + ) + except Exception: + raise RecordedSourceAdmissionError("recorded material failed activation-bound source mapping") from None + observation = mapped.observation + if ( + observation.mode is not IntelligenceResourceMode.PREPARED + or observation.acquisition_mode is not EvidenceAcquisitionMode.RECORDED_REPLAY + or observation.acquisition_receipt_ref != acquisition.receipt_id + or observation.acquisition_receipt_digest != acquisition.receipt_digest + ): + raise RecordedSourceAdmissionError("prepared mapping changed recorded acquisition truth") + acquisitions.append(acquisition) + snapshots.append(snapshot) + observations.append(observation) + + records_list: list[ImmutableRecordV1] = [] + for acquisition, snapshot, observation in zip(acquisitions, snapshots, observations, strict=True): + base = len(records_list) + records_list.extend( + ( + _record( + acquisition, + product_id=self.build.product_id, + kind=RECORDED_SOURCE_RECORD_KIND, + key=str(acquisition.receipt_id), + as_of=acquisition.observed_at, + available_at=admitted_at, + order=base, + ), + _record( + snapshot, + product_id=self.build.product_id, + kind=SOURCE_SNAPSHOT_RECORD_KIND, + key=str(snapshot.source_snapshot_ref), + as_of=snapshot.ingested_at, + available_at=admitted_at, + order=base + 1, + ), + _record( + observation, + product_id=self.build.product_id, + kind=IntelligenceRecordKind.OBSERVATION.value, + key=str(observation.resource_id), + as_of=observation.as_of, + available_at=admitted_at, + order=base + 2, + ), + ) + ) + request = AppendOnlyTransactionRequestV1( + product_id=self.build.product_id, + record_space=PREPARED_RECORD_SPACE, + transaction_key=transaction_key, + records=tuple(records_list), + submitted_at=admitted_at, + governed_state_preconditions=( + activation_head, + self.build.authority_use.state_head_precondition, + ), + ) + receipt = await self.store.append(request) + if receipt != request.receipt(): + raise RecordedSourceAdmissionError("Core append receipt does not bind exact recorded admission") + reopened = await self._replay(transaction_key=transaction_key, expected=exact, replayed=False) + if reopened is None: + raise RecordedSourceAdmissionError("recorded source admission did not reopen") + return reopened + + async def _replay( + self, + *, + transaction_key: str, + expected: tuple[RecordedSourceMaterialV1Alpha1, ...], + replayed: bool = True, + ) -> RecordedSourceAdmission | None: + receipt = await self.store.load_transaction_receipt( + product_id=self.build.product_id, + record_space=PREPARED_RECORD_SPACE, + transaction_key=transaction_key, + ) + if receipt is None: + return None + if len(receipt.records) != len(expected) * 3: + raise RecordedSourceAdmissionError("recorded admission receipt lost exact material-set shape") + loaded: list[ImmutableRecordV1] = [] + for reference in receipt.records: + record = await self.store.load_record( + reference.storage_id, + product_id=self.build.product_id, + record_space=PREPARED_RECORD_SPACE, + record_kind=reference.record_kind, + ) + if record is None or record.reference() != reference: + raise RecordedSourceAdmissionError("recorded admission has missing or changed immutable material") + loaded.append(record) + expected_kinds = ( + RECORDED_SOURCE_RECORD_KIND, + SOURCE_SNAPSHOT_RECORD_KIND, + IntelligenceRecordKind.OBSERVATION.value, + ) * len(expected) + if tuple(item.record_kind for item in loaded) != expected_kinds: + raise RecordedSourceAdmissionError("recorded admission record kinds changed") + acquisitions: list[RecordedSourceAcquisitionReceiptV1Alpha1] = [] + snapshots: list[CanonicalSourceSnapshotV1Alpha1] = [] + observations: list[ObservationV1Alpha1] = [] + for index, material in enumerate(expected): + offset = index * 3 + try: + acquisition = RecordedSourceAcquisitionReceiptV1Alpha1.model_validate(loaded[offset].payload) + snapshot = CanonicalSourceSnapshotV1Alpha1.model_validate(loaded[offset + 1].payload) + observation = ObservationV1Alpha1.model_validate(loaded[offset + 2].payload) + except (TypeError, ValueError) as exc: + raise RecordedSourceAdmissionError("recorded admission payload failed exact replay") from exc + if ( + acquisition.recorded_material_id != material.material_id + or acquisition.recorded_material_digest != material.material_digest + or acquisition.build_id != self.build.build_id + or acquisition.build_request_digest != self.build.request_digest + or acquisition.activation_revision != self.binding.prepared_binding.reference + or snapshot.acquisition_receipt_ref != acquisition.receipt_id + or snapshot.acquisition_receipt_digest != acquisition.receipt_digest + or observation.source_ref != snapshot.source_snapshot_ref + or observation.source_digest != snapshot.source_snapshot_digest + or observation.acquisition_receipt_ref != acquisition.receipt_id + or observation.acquisition_receipt_digest != acquisition.receipt_digest + or observation.activation_revision != self.binding.prepared_binding.reference + or observation.mode is not IntelligenceResourceMode.PREPARED + or observation.acquisition_mode is not EvidenceAcquisitionMode.RECORDED_REPLAY + ): + raise RecordedSourceAdmissionError("recorded admission chain crossed exact governed material") + acquisitions.append(acquisition) + snapshots.append(snapshot) + observations.append(observation) + first_acquisition = acquisitions[0] + expected_preconditions = tuple( + sorted( + ( + first_acquisition.activation_head_precondition, + first_acquisition.build_authority_use.state_head_precondition, + ), + key=lambda item: (item.state_kind, item.product_id, item.state_id), + ) + ) + if receipt.governed_state_preconditions != expected_preconditions: + raise RecordedSourceAdmissionError("recorded admission lost activation or authority precondition") + return RecordedSourceAdmission( + acquisition_receipts=tuple(acquisitions), + source_snapshots=tuple(snapshots), + observations=tuple(observations), + transaction_receipt=receipt, + replayed=replayed, + ) + + +__all__ = [ + "CONNECT_SOURCES_EFFECT", + "CoreRecordedSourceAdmissionService", + "IntelligenceBuildRecordedSourcePort", + "RECORDED_SOURCE_ACQUISITION_RECEIPT_VERSION", + "RECORDED_SOURCE_MATERIAL_VERSION", + "RecordedSourceAcquisitionReceiptV1Alpha1", + "RecordedSourceAdmission", + "RecordedSourceAdmissionError", + "RecordedSourceMaterialV1Alpha1", +] diff --git a/tests/intelligence/test_recorded_source_admission.py b/tests/intelligence/test_recorded_source_admission.py new file mode 100644 index 0000000..97e26c0 --- /dev/null +++ b/tests/intelligence/test_recorded_source_admission.py @@ -0,0 +1,249 @@ +from __future__ import annotations + +import hashlib +from datetime import UTC, datetime, timedelta + +import pytest + +from ace.application.domain_activation import DomainActivationAdmissionService, bind_committed_activation +from ace.application.intelligence_build_execution import ( + REQUIRED_INTELLIGENCE_BUILD_EFFECTS, + AuthorizedIntelligenceBuild, + IntelligenceBuildStartV1, + RecordedSourceReferenceV1, +) +from ace.application.intelligence_resource_projection import IntelligenceLedgerResourceProjectionReader +from ace.application.recorded_source_admission import ( + CoreRecordedSourceAdmissionService, + RecordedSourceAdmissionError, + RecordedSourceMaterialV1Alpha1, +) +from ace.core import ( + AuthenticatedRuntimeContextV1Alpha1, + AuthorityUseReceiptV1Alpha1, + GovernedStateHeadPreconditionV1Alpha1, + GovernedStateHeadV1, + canonical_json, +) +from ace.core.runtime_use import AUTHORITY_GRANT_STATE_KIND +from ace.intelligence import EvidenceAcquisitionMode, IntelligenceResourceMode +from ace.intelligence.contracts.resource_plane import ( + IntelligenceResourceKind, + IntelligenceResourceQueryV1Alpha1, +) +from ace.testing import InMemoryImmutableRecordStore +from tests.intelligence.test_domain_activation_admission import _Authority, _MemoryStore +from tests.intelligence.test_source_mapping import _binding, _compiled, _fixture_documents, _subject + +pytestmark = pytest.mark.unit + +PRODUCT = "product:recorded-source" +ACTOR = "principal:personal-operator" +OBSERVED_AT = datetime(2026, 8, 12, 18, tzinfo=UTC) +ADMITTED_AT = datetime(2026, 8, 13, 18, tzinfo=UTC) + + +async def _stack(): + pack = _compiled("numeric") + prepared = _binding(pack, product_id=PRODUCT) + activation_store = _MemoryStore() + committed = await DomainActivationAdmissionService( + store=activation_store, + authority=_Authority(), + ).admit( + prepared.revision, + expected_head_revision_id=None, + committed_at=prepared.revision.occurred_at + timedelta(seconds=1), + ) + binding = bind_committed_activation(pack=pack, committed=committed) + + grant_head = GovernedStateHeadV1( + state_kind=AUTHORITY_GRANT_STATE_KIND, + product_id=PRODUCT, + state_id="authority_grant:atrium-intelligence-build", + sequence=1, + revision_id="authority_grant_revision:recorded-source", + commit_receipt_id="governed_state_commit:recorded-source", + updated_at=ADMITTED_AT - timedelta(minutes=1), + ) + activation_head = activation_store.heads[ + ("domain_activation", PRODUCT, str(binding.prepared_binding.revision.activation_id)) + ] + records = InMemoryImmutableRecordStore( + governed_state_heads={ + (activation_head.state_kind, activation_head.product_id, activation_head.state_id): activation_head, + (grant_head.state_kind, grant_head.product_id, grant_head.state_id): grant_head, + } + ) + context = AuthenticatedRuntimeContextV1Alpha1( + product_id=PRODUCT, + actor_ref=ACTOR, + authentication_receipt_ref="authentication_receipt:recorded-source", + authentication_receipt_digest="sha256:" + "1" * 64, + authenticated_at=ADMITTED_AT - timedelta(minutes=2), + expires_at=ADMITTED_AT + timedelta(hours=1), + ) + _, _, payload = _fixture_documents("numeric") + payload_json = canonical_json(payload) + material = RecordedSourceMaterialV1Alpha1( + source_group_id="official_records", + mapping_id="reading_snapshot", + subject_binding=_subject(binding.prepared_binding, "numeric"), + source_definition_ref="source_definition:numeric", + source_type_ref="source:reading/v1", + source_uri="https://example.invalid/recorded/reading-1", + captured_payload_json=payload_json, + captured_payload_digest="sha256:" + hashlib.sha256(payload_json.encode()).hexdigest(), + source_published_at=OBSERVED_AT - timedelta(hours=1), + event_effective_at=OBSERVED_AT - timedelta(minutes=30), + observed_at=OBSERVED_AT, + locator="record:1", + ) + request = IntelligenceBuildStartV1( + authority_grant_ref=grant_head.state_id, + resource_authority_grant_ref="authority_grant:atrium-observe-read", + activation_approval_receipt_ref=str(binding.commit_receipt.approval.receipt_ref), + activation_approval_subject_ref=str(binding.prepared_binding.revision.spec.spec_id), + client_request_id="atrium_request:recorded-source", + profile_id="intelligence_onboarding_profile:recorded-source-fixture", + subject="Track the reviewed recorded source for material changes.", + outcome_id="decision_readiness", + source_group_ids=("official_records",), + recorded_source_refs=( + RecordedSourceReferenceV1( + source_group_id=material.source_group_id, + material_id=str(material.material_id), + material_digest=str(material.material_digest), + ), + ), + cadence_id="daily_pulse", + approved_effects=REQUIRED_INTELLIGENCE_BUILD_EFFECTS, + requested_at=ADMITTED_AT - timedelta(minutes=1), + ) + build_id = "intelligence_build:recorded-source" + request_digest = "sha256:" + "2" * 64 + authority_use = AuthorityUseReceiptV1Alpha1( + product_id=PRODUCT, + actor_ref=ACTOR, + authenticated_context=context, + use_subject_ref=build_id, + use_subject_digest=request_digest, + operation="start_intelligence_build", + authority="intelligence_build", + grant_ref=grant_head.state_id, + grant_hash="3" * 64, + evaluated_at=ADMITTED_AT, + expires_at=ADMITTED_AT + timedelta(hours=1), + state_head_precondition=GovernedStateHeadPreconditionV1Alpha1.from_head(grant_head), + ) + build = AuthorizedIntelligenceBuild( + build_id=build_id, + request_digest=request_digest, + product_id=PRODUCT, + actor_ref=ACTOR, + request=request, + authority_use=authority_use, + activation_approval=binding.commit_receipt.approval, + ) + return binding, build, records, material + + +@pytest.mark.asyncio +async def test_recorded_replay_admits_canonical_observation_and_reopens_exactly() -> None: + binding, build, records, material = await _stack() + service = CoreRecordedSourceAdmissionService(build=build, binding=binding, store=records) + + first = await service.admit((material,)) + replay = await CoreRecordedSourceAdmissionService( + build=build, + binding=binding, + store=records, + ).admit((material,)) + + acquisition = first.acquisition_receipts[0] + assert acquisition.network_capture_performed is False + assert acquisition.freshness_verified is False + assert acquisition.live_acquisition is False + assert acquisition.reusable_authority is False + assert first.source_snapshots[0].acquisition_mode.value == "recorded_replay" + assert first.observations[0].mode is IntelligenceResourceMode.PREPARED + assert first.observations[0].acquisition_mode is EvidenceAcquisitionMode.RECORDED_REPLAY + assert first.observations[0].acquisition_receipt_ref == acquisition.receipt_id + assert len(first.transaction_receipt.governed_state_preconditions) == 2 + assert replay.replayed is True + assert replay.acquisition_receipts == first.acquisition_receipts + assert replay.source_snapshots == first.source_snapshots + assert replay.observations == first.observations + assert replay.transaction_receipt == first.transaction_receipt + + +@pytest.mark.asyncio +async def test_recorded_observation_is_visible_from_fresh_canonical_resource_projection() -> None: + binding, build, records, material = await _stack() + admitted = await CoreRecordedSourceAdmissionService( + build=build, + binding=binding, + store=records, + ).admit((material,)) + query = IntelligenceResourceQueryV1Alpha1( + authenticated_context=build.authority_use.authenticated_context, + product_id=PRODUCT, + authority_grant_ref="authority_grant:atrium-observe-read", + resource_kinds=(IntelligenceResourceKind.OBSERVATION,), + subject_refs=(material.subject_binding.entity_ref,), + as_of=ADMITTED_AT, + available_at=ADMITTED_AT, + page_size=10, + ) + page = await IntelligenceLedgerResourceProjectionReader(store=records).read( + query=query, + after=None, + limit=10, + ) + + assert len(page.records) == 1 + assert page.records[0].reference.resource_id == admitted.observations[0].resource_id + assert page.records[0].reference.resource_kind is IntelligenceResourceKind.OBSERVATION + + +@pytest.mark.asyncio +async def test_substituted_recorded_material_fails_before_any_write() -> None: + binding, build, records, material = await _stack() + changed_payload = canonical_json({"reading": {"value": "999.000"}, "subject": {"code": "AX"}}) + changed = RecordedSourceMaterialV1Alpha1( + **material.model_dump( + mode="python", + exclude={"captured_payload_json", "captured_payload_digest", "material_id", "material_digest"}, + ), + captured_payload_json=changed_payload, + captured_payload_digest="sha256:" + hashlib.sha256(changed_payload.encode()).hexdigest(), + ) + with pytest.raises(RecordedSourceAdmissionError, match="exactly match"): + await CoreRecordedSourceAdmissionService(build=build, binding=binding, store=records).admit((changed,)) + assert await records.scan_product_records(product_id=PRODUCT) == () + + +@pytest.mark.asyncio +async def test_missing_or_extra_recorded_material_fails_before_any_write() -> None: + binding, build, records, material = await _stack() + service = CoreRecordedSourceAdmissionService(build=build, binding=binding, store=records) + with pytest.raises(RecordedSourceAdmissionError, match="requires exact reviewed material"): + await service.admit(()) + + extra = RecordedSourceMaterialV1Alpha1( + **material.model_dump(mode="python", exclude={"source_uri", "material_id", "material_digest"}), + source_uri="https://example.invalid/recorded/reading-2", + ) + with pytest.raises(RecordedSourceAdmissionError, match="exactly match"): + await service.admit((material, extra)) + assert await records.scan_product_records(product_id=PRODUCT) == () + + +@pytest.mark.asyncio +async def test_stale_activation_or_authority_head_fails_atomic_append() -> None: + binding, build, records, material = await _stack() + records.governed_state_heads.clear() + + with pytest.raises(Exception, match="precondition"): + await CoreRecordedSourceAdmissionService(build=build, binding=binding, store=records).admit((material,)) + assert await records.scan_product_records(product_id=PRODUCT) == () diff --git a/tests/test_api_intelligence_builds.py b/tests/test_api_intelligence_builds.py index 18f8711..0c27429 100644 --- a/tests/test_api_intelligence_builds.py +++ b/tests/test_api_intelligence_builds.py @@ -181,6 +181,7 @@ async def _request( authority: _Authority, executor: _Executor, activation_authority: _ActivationAuthority | None = None, + body: dict | None = None, ): app = FastAPI() app.include_router(router) @@ -193,7 +194,7 @@ async def _request( executor=executor, ) async with AsyncClient(transport=ASGITransport(app=app), base_url="http://test") as client: - response = await client.post("/v1/intelligence/builds/start", json=_body()) + response = await client.post("/v1/intelligence/builds/start", json=body or _body()) return response, records @@ -266,6 +267,34 @@ async def test_default_runtime_has_no_implicit_activation_approval_authority() - await runtime.activation_authority.resolve_approval() +@pytest.mark.asyncio +async def test_exact_recorded_material_coordinates_change_the_authorized_build_digest() -> None: + first_body = _body() + first_body["source_group_ids"] = ["official_records"] + first_body["recorded_source_refs"] = [ + { + "source_group_id": "official_records", + "material_id": "recorded_source_material:one", + "material_digest": "sha256:" + "1" * 64, + } + ] + second_body = {**first_body} + second_body["recorded_source_refs"] = [ + { + **first_body["recorded_source_refs"][0], + "material_digest": "sha256:" + "2" * 64, + } + ] + first_authority = _Authority() + second_authority = _Authority() + + first, _ = await _request(claims=_claims(), authority=first_authority, executor=_Executor(), body=first_body) + second, _ = await _request(claims=_claims(), authority=second_authority, executor=_Executor(), body=second_body) + + assert first.status_code == second.status_code == 200 + assert first_authority.calls[0]["use_subject_digest"] != second_authority.calls[0]["use_subject_digest"] + + def test_start_build_openapi_exposes_stable_request_and_result_contracts() -> None: app = FastAPI() app.include_router(router)