diff --git a/ace/application/__init__.py b/ace/application/__init__.py index de32c51..283f698 100644 --- a/ace/application/__init__.py +++ b/ace/application/__init__.py @@ -350,6 +350,7 @@ IntelligenceLedgerProjectionError, IntelligenceLedgerResourceProjectionReader, IntelligenceResourceProjectionContributor, + LiveSourceResourceProjectionReader, MonitoringResourceProjectionReader, ) from ace.application.live_intelligence_bridge import ( @@ -461,6 +462,7 @@ "IntelligenceLedgerProjectionError", "IntelligenceLedgerResourceProjectionReader", "IntelligenceResourceProjectionContributor", + "LiveSourceResourceProjectionReader", "CompositeIntelligenceResourceProjectionReader", "DecisionOutcomeFeedbackResourceProjectionReader", "MonitoringResourceProjectionReader", diff --git a/ace/application/intelligence_resource_projection.py b/ace/application/intelligence_resource_projection.py index 6af93e4..61ff308 100644 --- a/ace/application/intelligence_resource_projection.py +++ b/ace/application/intelligence_resource_projection.py @@ -14,9 +14,10 @@ IntelligenceResourceProjectionReader, ) from ace.application.monitoring import LIVE_MONITORING_RECORD_SPACE -from ace.core.contracts import canonical_json +from ace.core.contracts import canonical_hash, canonical_json from ace.core.decisions import DecisionV1Alpha1, OutcomeV1Alpha1 from ace.core.records import ImmutableRecordStore, ImmutableRecordV1 +from ace.core.source import CanonicalSourceSnapshotV1Alpha1 from ace.intelligence.contracts.feedback import FeedbackProposalV1Alpha1 from ace.intelligence.contracts.ledger import IntelligenceRecordKind from ace.intelligence.contracts.monitoring import ( @@ -46,6 +47,11 @@ ShiftV1Alpha1, SignalV1Alpha1, ) +from ace.intelligence.contracts.source_acquisition import ( + LiveSourceAdmissionReceiptV1Alpha1, + LiveSourceIngressRecordKind, + SourceAcquisitionReceiptV1Alpha1, +) _JSON_OBJECT = TypeAdapter(dict[str, Any]) @@ -75,6 +81,12 @@ IntelligenceResourceKind.FEEDBACK, } ) +LIVE_SOURCE_RESOURCE_KINDS = frozenset( + { + IntelligenceResourceKind.CONNECTION, + IntelligenceResourceKind.SOURCE, + } +) _LINEAGE_TO_PUBLIC: dict[LineageResourceKind, IntelligenceResourceKind] = { LineageResourceKind.OBSERVATION: IntelligenceResourceKind.OBSERVATION, LineageResourceKind.ENTITY_SNAPSHOT: IntelligenceResourceKind.ENTITY, @@ -729,6 +741,324 @@ async def read( ) +def _source_reference( + *, + product_id: str, + source_definition_ref: str, + snapshot: CanonicalSourceSnapshotV1Alpha1, + admission: LiveSourceAdmissionReceiptV1Alpha1, + revision: int, +) -> IntelligenceResourceReferenceV1Alpha1: + return IntelligenceResourceReferenceV1Alpha1( + product_id=product_id, + resource_kind=IntelligenceResourceKind.SOURCE, + resource_id=source_definition_ref, + resource_digest=str(snapshot.source_snapshot_digest), + resource_contract=snapshot.contract, + revision=revision, + as_of=snapshot.as_of, + available_at=admission.admitted_at, + ) + + +def _connection_reference( + *, + product_id: str, + source_definition_ref: str, + acquisition: SourceAcquisitionReceiptV1Alpha1, + admission: LiveSourceAdmissionReceiptV1Alpha1, + revision: int, +) -> IntelligenceResourceReferenceV1Alpha1: + return IntelligenceResourceReferenceV1Alpha1( + product_id=product_id, + resource_kind=IntelligenceResourceKind.CONNECTION, + resource_id=f"connection:{canonical_hash([product_id, source_definition_ref])[:32]}", + resource_digest=str(admission.receipt_digest), + resource_contract=admission.contract, + revision=revision, + as_of=acquisition.captured_at, + available_at=admission.admitted_at, + ) + + +def _live_source_records( + *, + acquisition: SourceAcquisitionReceiptV1Alpha1, + snapshot: CanonicalSourceSnapshotV1Alpha1, + admission: LiveSourceAdmissionReceiptV1Alpha1, + revision: int, + prior_connection: IntelligenceResourceReferenceV1Alpha1 | None = None, + prior_source: IntelligenceResourceReferenceV1Alpha1 | None = None, +) -> tuple[IntelligenceResourceRecordV1Alpha1, IntelligenceResourceRecordV1Alpha1]: + source_definition_ref = acquisition.source_definition_ref + source_reference = _source_reference( + product_id=acquisition.product_id, + source_definition_ref=source_definition_ref, + snapshot=snapshot, + admission=admission, + revision=revision, + ) + connection_reference = _connection_reference( + product_id=acquisition.product_id, + source_definition_ref=source_definition_ref, + acquisition=acquisition, + admission=admission, + revision=revision, + ) + if revision == 1 and (prior_connection is not None or prior_source is not None): + raise ValueError("first live source revision cannot supersede prior material") + if revision > 1 and ( + prior_connection is None + or prior_source is None + or prior_connection.resource_id != connection_reference.resource_id + or prior_source.resource_id != source_reference.resource_id + or prior_connection.revision != revision - 1 + or prior_source.revision != revision - 1 + ): + raise ValueError("later live source revision requires both exact prior references") + connection_payload = { + "actor_ref": acquisition.actor_ref, + "source_definition_ref": source_definition_ref, + "source_type_ref": acquisition.source_type_ref, + "configuration_ref": acquisition.configuration_ref, + "configuration_digest": acquisition.configuration_digest, + "adapter_artifact": acquisition.adapter_artifact.model_dump(mode="json"), + "capability_use_receipt_ref": admission.capability_use_receipt_ref, + "capability_use_receipt_digest": admission.capability_use_receipt_digest, + "authority_use_receipt_ref": admission.authority_use_receipt_ref, + "authority_use_receipt_digest": admission.authority_use_receipt_digest, + "acquisition_receipt_ref": str(acquisition.receipt_id), + "acquisition_receipt_digest": str(acquisition.receipt_digest), + "admission_receipt_ref": str(admission.receipt_id), + "admission_receipt_digest": str(admission.receipt_digest), + "captured_at": acquisition.captured_at.isoformat(), + "admitted_at": admission.admitted_at.isoformat(), + } + source_payload = { + "source_definition_ref": source_definition_ref, + "source_type_ref": snapshot.source_type_ref, + "source_snapshot_ref": str(snapshot.source_snapshot_ref), + "source_snapshot_digest": str(snapshot.source_snapshot_digest), + "source_published_at": ( + None if snapshot.source_published_at is None else snapshot.source_published_at.isoformat() + ), + "event_effective_at": ( + None if snapshot.event_effective_at is None else snapshot.event_effective_at.isoformat() + ), + "observed_at": snapshot.observed_at.isoformat(), + "ingested_at": snapshot.ingested_at.isoformat(), + "acquisition_receipt_ref": str(acquisition.receipt_id), + "acquisition_receipt_digest": str(acquisition.receipt_digest), + "admission_receipt_ref": str(admission.receipt_id), + "admission_receipt_digest": str(admission.receipt_digest), + "captured_payload_redacted": True, + } + subjects = tuple( + sorted( + { + acquisition.actor_ref, + source_definition_ref, + acquisition.source_type_ref, + } + ) + ) + connection = IntelligenceResourceRecordV1Alpha1( + reference=connection_reference, + availability=IntelligenceResourceAvailability.AVAILABLE, + title=f"Connection: {source_definition_ref}", + summary=f"Successful governed capture through {acquisition.adapter_artifact.implementation_id}.", + subject_refs=subjects, + supersedes=prior_connection, + payload=CanonicalJsonValueV1Alpha1(value_json=canonical_json(connection_payload)), + ) + source = IntelligenceResourceRecordV1Alpha1( + reference=source_reference, + availability=IntelligenceResourceAvailability.AVAILABLE, + title=f"Source: {source_definition_ref}", + summary=f"Latest admitted {snapshot.source_type_ref} capture metadata; captured payload is redacted.", + subject_refs=subjects, + provenance=(connection_reference,), + supersedes=prior_source, + payload=CanonicalJsonValueV1Alpha1(value_json=canonical_json(source_payload)), + ) + return connection, source + + +def _decode_live_source_chain( + *, + acquisition_record: ImmutableRecordV1, + snapshot_record: ImmutableRecordV1, + admission_record: ImmutableRecordV1, +) -> tuple[SourceAcquisitionReceiptV1Alpha1, CanonicalSourceSnapshotV1Alpha1, LiveSourceAdmissionReceiptV1Alpha1]: + acquisition = SourceAcquisitionReceiptV1Alpha1.model_validate(acquisition_record.payload) + snapshot = CanonicalSourceSnapshotV1Alpha1.model_validate(snapshot_record.payload) + admission = LiveSourceAdmissionReceiptV1Alpha1.model_validate(admission_record.payload) + if ( + len({acquisition_record.product_id, snapshot_record.product_id, admission_record.product_id}) != 1 + or acquisition_record.record_space != IntelligenceResourceMode.LIVE.value + or snapshot_record.record_space != IntelligenceResourceMode.LIVE.value + or admission_record.record_space != IntelligenceResourceMode.LIVE.value + or acquisition_record.record_kind != LiveSourceIngressRecordKind.SOURCE_ACQUISITION.value + or snapshot_record.record_kind != LiveSourceIngressRecordKind.SOURCE_SNAPSHOT.value + or admission_record.record_kind != LiveSourceIngressRecordKind.SOURCE_ADMISSION.value + or acquisition_record.record_key != acquisition.receipt_id + or snapshot_record.record_key != snapshot.source_snapshot_ref + or admission_record.record_key != admission.receipt_id + or acquisition_record.payload_contract != acquisition.contract + or snapshot_record.payload_contract != snapshot.contract + or admission_record.payload_contract != admission.contract + or acquisition_record.as_of != acquisition.captured_at + or snapshot_record.as_of != snapshot.as_of + or admission_record.as_of != admission.admitted_at + or acquisition_record.available_at != admission.admitted_at + or snapshot_record.available_at != admission.admitted_at + or admission_record.available_at != admission.admitted_at + or snapshot.acquisition_receipt_ref != acquisition.receipt_id + or snapshot.acquisition_receipt_digest != acquisition.receipt_digest + or admission.acquisition_receipt_ref != acquisition.receipt_id + or admission.acquisition_receipt_digest != acquisition.receipt_digest + or admission.source_snapshot_ref != snapshot.source_snapshot_ref + or admission.source_snapshot_digest != snapshot.source_snapshot_digest + or acquisition.product_id != admission.product_id + or acquisition.actor_ref != admission.actor_ref + or acquisition.use_subject_ref != admission.use_subject_ref + or acquisition.use_subject_digest != admission.use_subject_digest + or acquisition.operation != admission.operation + or acquisition.source_definition_ref != snapshot.source_definition_ref + or acquisition.source_type_ref != snapshot.source_type_ref + or acquisition.source_definition_head_precondition != admission.source_definition_head_precondition + or acquisition.receipt_id != snapshot.acquisition_receipt_ref + or acquisition.receipt_digest != snapshot.acquisition_receipt_digest + or acquisition.captured_payload_digest != snapshot.captured_payload_digest + or acquisition.source_published_at != snapshot.source_published_at + or acquisition.event_effective_at != snapshot.event_effective_at + or acquisition.observed_at != snapshot.observed_at + or acquisition.capability_use.receipt_id != admission.capability_use_receipt_ref + or acquisition.capability_use.receipt_digest != admission.capability_use_receipt_digest + or acquisition.authority_use.receipt_id != admission.authority_use_receipt_ref + or acquisition.authority_use.receipt_digest != admission.authority_use_receipt_digest + or admission.source_definition_head_precondition.state_id != acquisition.source_definition_ref + ): + raise ValueError("live source records do not form one exact admitted chain") + return acquisition, snapshot, admission + + +class LiveSourceResourceProjectionReader(IntelligenceResourceProjectionReader): + """Project successful governed captures as Connection and redacted Source revisions.""" + + def __init__(self, *, store: ImmutableRecordStore, degrade_unsupported: bool = True) -> None: + self.store = store + self.degrade_unsupported = degrade_unsupported + + @property + def supported_kinds(self) -> frozenset[IntelligenceResourceKind]: + return LIVE_SOURCE_RESOURCE_KINDS + + async def read( + self, + *, + query: IntelligenceResourceQueryV1Alpha1, + after: IntelligenceResourceCursorV1Alpha1 | None, + limit: int, + ) -> IntelligenceResourceProjectionBatch: + requested = set(query.resource_kinds) + relevant = requested & LIVE_SOURCE_RESOURCE_KINDS + degraded = { + f"degraded_reason:unsupported-{kind.value}" + for kind in requested - LIVE_SOURCE_RESOURCE_KINDS + if self.degrade_unsupported + } + if not relevant: + return IntelligenceResourceProjectionBatch( + records=(), + state=(IntelligenceResourcePageState.DEGRADED if degraded else IntelligenceResourcePageState.COMPLETE), + degraded_reason_refs=tuple(sorted(degraded)), + ) + buckets: dict[str, tuple[ImmutableRecordV1, ...]] = {} + for kind in ( + LiveSourceIngressRecordKind.SOURCE_ACQUISITION, + LiveSourceIngressRecordKind.SOURCE_SNAPSHOT, + LiveSourceIngressRecordKind.SOURCE_ADMISSION, + ): + try: + buckets[kind.value] = await self.store.read_as_of( + product_id=query.product_id, + record_space=IntelligenceResourceMode.LIVE.value, + record_kind=kind.value, + available_at=query.available_at, + ) + except Exception: + degraded.add(f"degraded_reason:read-live-{kind.value}") + buckets[kind.value] = () + acquisition_records = buckets[LiveSourceIngressRecordKind.SOURCE_ACQUISITION] + snapshot_records = buckets[LiveSourceIngressRecordKind.SOURCE_SNAPSHOT] + acquisitions = {record.record_key: record for record in acquisition_records} + snapshots = {record.record_key: record for record in snapshot_records} + if len(acquisitions) != len(acquisition_records) or len(snapshots) != len(snapshot_records): + degraded.add("degraded_reason:duplicate-live-source-record") + chains: dict[ + str, + list[ + tuple[ + SourceAcquisitionReceiptV1Alpha1, + CanonicalSourceSnapshotV1Alpha1, + LiveSourceAdmissionReceiptV1Alpha1, + ] + ], + ] = defaultdict(list) + used_acquisition_refs: set[str] = set() + used_snapshot_refs: set[str] = set() + for admission_record in buckets[LiveSourceIngressRecordKind.SOURCE_ADMISSION]: + try: + admission = LiveSourceAdmissionReceiptV1Alpha1.model_validate(admission_record.payload) + acquisition_record = acquisitions[str(admission.acquisition_receipt_ref)] + snapshot_record = snapshots[str(admission.source_snapshot_ref)] + chain = _decode_live_source_chain( + acquisition_record=acquisition_record, + snapshot_record=snapshot_record, + admission_record=admission_record, + ) + chains[chain[0].source_definition_ref].append(chain) + used_acquisition_refs.add(str(chain[0].receipt_id)) + used_snapshot_refs.add(str(chain[1].source_snapshot_ref)) + except Exception: + degraded.add("degraded_reason:invalid-live-source-chain") + if set(acquisitions) - used_acquisition_refs or set(snapshots) - used_snapshot_refs: + degraded.add("degraded_reason:orphan-live-source-record") + projected: list[IntelligenceResourceRecordV1Alpha1] = [] + for source_definition_ref, source_chains in chains.items(): + ordered = sorted( + source_chains, + key=lambda chain: (chain[2].admitted_at, str(chain[2].receipt_digest)), + ) + prior_connection: IntelligenceResourceReferenceV1Alpha1 | None = None + prior_source: IntelligenceResourceReferenceV1Alpha1 | None = None + for revision, (acquisition, snapshot, admission) in enumerate(ordered, start=1): + connection, source = _live_source_records( + acquisition=acquisition, + snapshot=snapshot, + admission=admission, + revision=revision, + prior_connection=prior_connection, + prior_source=prior_source, + ) + prior_connection = connection.reference + prior_source = source.reference + for item in (connection, source): + if item.reference.resource_kind not in relevant or item.reference.as_of > query.as_of: + continue + if query.subject_refs and set(query.subject_refs).isdisjoint(item.subject_refs): + continue + projected.append(item) + visible = _after_cursor(projected, after)[:limit] + reasons = tuple(sorted(degraded)) + return IntelligenceResourceProjectionBatch( + records=tuple(visible), + state=(IntelligenceResourcePageState.DEGRADED if reasons else IntelligenceResourcePageState.COMPLETE), + degraded_reason_refs=reasons, + ) + + class CompositeIntelligenceResourceProjectionReader(IntelligenceResourceProjectionReader): """Merge disjoint rebuildable projection contributors into one stable page.""" @@ -786,5 +1116,6 @@ async def read( "IntelligenceLedgerProjectionError", "IntelligenceLedgerResourceProjectionReader", "IntelligenceResourceProjectionContributor", + "LiveSourceResourceProjectionReader", "MonitoringResourceProjectionReader", ] diff --git a/core/engine/core/intelligence_resource_plane.py b/core/engine/core/intelligence_resource_plane.py index 74d992d..2de0c12 100644 --- a/core/engine/core/intelligence_resource_plane.py +++ b/core/engine/core/intelligence_resource_plane.py @@ -25,6 +25,7 @@ IntelligenceResourcePlaneService, IntelligenceResourceProjectionReader, IntelligenceResourceQueryV1Alpha1, + LiveSourceResourceProjectionReader, MonitoringResourceProjectionReader, ) from ace.core import ImmutableRecordPersistenceError, ImmutableRecordStore @@ -99,6 +100,10 @@ def intelligence_resource_projection_reader(records: ImmutableRecordStore) -> In store=records, degrade_unsupported=False, ), + LiveSourceResourceProjectionReader( + store=records, + degrade_unsupported=False, + ), ) diff --git a/docs/design/intelligence-resource-plane-v0.8.0-work-packet-v1.md b/docs/design/intelligence-resource-plane-v0.8.0-work-packet-v1.md index 5d1e756..cfc6613 100644 --- a/docs/design/intelligence-resource-plane-v0.8.0-work-packet-v1.md +++ b/docs/design/intelligence-resource-plane-v0.8.0-work-packet-v1.md @@ -1,10 +1,10 @@ # ACE 0.8.0 unified Intelligence resource plane work packet Status: **active 0.8C packet; facade, intelligence/monitoring projections, governed HTTP query, -and Decision → Outcome → Feedback closure implemented** +Decision → Outcome → Feedback closure, and live Connection/Source projections implemented** Public milestone: [issue #40](https://github.com/augmented-cognition-engine/core/issues/40) -Accepted base: `main@cf53360` (0.8A architecture, AM4 lifecycle, completed 0.8B, facade, -ledger/monitoring projection, and governed HTTP query) +Accepted base: `main@5003cfb` (0.8A architecture, AM4 lifecycle, completed 0.8B, facade, +ledger/monitoring projection, governed HTTP query, and Decision → Outcome → Feedback closure) ## Outcome @@ -97,6 +97,16 @@ types remain visible only as explicitly degraded truth. The supported host compo contributors through one named factory so future resource families cannot silently bypass the same API path. +The live-source contributor projects successful governed admissions as versioned Connections and +Sources. It requires the exact acquisition → snapshot → admission chain, preserves the exact prior +revision rather than reconstructing a partial reference, and exposes only redacted source metadata; +captured payloads, URIs, locators, resolved addresses, and credentials never enter the public read +model. Rebuilding the contributor over the same immutable store reproduces the same revisions. +Partial, orphaned, duplicate, or inconsistent chains degrade without exposing material. Source +Health remains explicitly unsupported because the current immutable records prove successful +admissions but do not yet provide failure and health telemetry; 0.8C will not infer health from +success-only history. + 0.8C must still add governed-state projections for the remaining canonical families and complete packaged schema/import integrity. 0.8D must prove Atrium consumes this interface rather than privileged internal state. diff --git a/tests/intelligence/test_live_source_resource_projection.py b/tests/intelligence/test_live_source_resource_projection.py new file mode 100644 index 0000000..6404126 --- /dev/null +++ b/tests/intelligence/test_live_source_resource_projection.py @@ -0,0 +1,160 @@ +from __future__ import annotations + +from datetime import timedelta + +import pytest + +from ace.application import LiveSourceResourceProjectionReader +from ace.core import AuthenticatedRuntimeContextV1Alpha1 +from ace.intelligence import ( + IntelligenceResourceKind, + IntelligenceResourcePageState, + IntelligenceResourceQueryV1Alpha1, + LiveSourceIngressRecordKind, +) +from tests.intelligence.test_live_source_ingress import BASE, PRODUCT, _Clock, _environment + +pytestmark = pytest.mark.unit + + +def _query(*kinds: IntelligenceResourceKind) -> IntelligenceResourceQueryV1Alpha1: + return IntelligenceResourceQueryV1Alpha1( + authenticated_context=AuthenticatedRuntimeContextV1Alpha1( + product_id=PRODUCT, + actor_ref="principal:analyst", + authentication_receipt_ref="authentication_receipt:source-projection", + authentication_receipt_digest="sha256:" + "d" * 64, + authenticated_at=BASE, + expires_at=BASE + timedelta(minutes=5), + ), + product_id=PRODUCT, + authority_grant_ref="authority_grant:resource-read", + resource_kinds=kinds, + as_of=BASE + timedelta(minutes=1), + available_at=BASE + timedelta(minutes=1), + page_size=20, + ) + + +@pytest.mark.asyncio +async def test_successful_governed_capture_projects_redacted_connection_and_source() -> None: + env = await _environment() + await env.service.admit(request=env.request, pack=env.pack) + + batch = await LiveSourceResourceProjectionReader(store=env.record_store).read( + query=_query(IntelligenceResourceKind.CONNECTION, IntelligenceResourceKind.SOURCE), + after=None, + limit=20, + ) + + assert batch.state is IntelligenceResourcePageState.COMPLETE + assert [item.reference.resource_kind for item in batch.records] == [ + IntelligenceResourceKind.CONNECTION, + IntelligenceResourceKind.SOURCE, + ] + connection, source = batch.records + assert connection.reference.revision == source.reference.revision == 1 + assert source.provenance == (connection.reference,) + assert connection.payload is not None + assert source.payload is not None + connection_payload = connection.payload.parsed_value() + source_payload = source.payload.parsed_value() + assert connection_payload["source_definition_ref"] == env.request.source_definition_ref + assert source_payload["captured_payload_redacted"] is True + serialized = f"{connection_payload}{source_payload}" + assert "captured_payload_json" not in serialized + assert "requested_uri" not in serialized + assert "effective_uri" not in serialized + assert "resolved_ip_addresses" not in serialized + assert "locator" not in serialized + + restarted = await LiveSourceResourceProjectionReader(store=env.record_store).read( + query=_query(IntelligenceResourceKind.CONNECTION, IntelligenceResourceKind.SOURCE), + after=None, + limit=20, + ) + assert restarted == batch + + +@pytest.mark.asyncio +async def test_repeated_capture_preserves_exact_revision_lineage() -> None: + first = await _environment() + await first.service.admit(request=first.request, pack=first.pack) + second = await _environment( + record_store=first.record_store, + clock=_Clock( + BASE + timedelta(seconds=20), + BASE + timedelta(seconds=22), + BASE + timedelta(seconds=23), + ), + ) + second_request = second.request.model_copy( + update={ + "idempotency_key": "live-ingress:two", + "requested_at": BASE + timedelta(seconds=10), + "request_id": None, + "request_digest": None, + } + ) + second_request = type(second.request).model_validate(second_request.model_dump(mode="python")) + await second.service.admit(request=second_request, pack=second.pack) + + batch = await LiveSourceResourceProjectionReader(store=second.record_store).read( + query=_query(IntelligenceResourceKind.CONNECTION, IntelligenceResourceKind.SOURCE), + after=None, + limit=20, + ) + + assert batch.state is IntelligenceResourcePageState.COMPLETE + by_kind = { + kind: sorted( + (item for item in batch.records if item.reference.resource_kind is kind), + key=lambda item: item.reference.revision, + ) + for kind in (IntelligenceResourceKind.CONNECTION, IntelligenceResourceKind.SOURCE) + } + for revisions in by_kind.values(): + assert [item.reference.revision for item in revisions] == [1, 2] + assert revisions[1].supersedes == revisions[0].reference + assert revisions[1].reference.resource_digest != revisions[0].reference.resource_digest + assert by_kind[IntelligenceResourceKind.SOURCE][1].provenance == ( + by_kind[IntelligenceResourceKind.CONNECTION][1].reference, + ) + + +@pytest.mark.asyncio +async def test_partial_or_invalid_admission_degrades_without_exposing_source_material() -> None: + env = await _environment() + await env.service.admit(request=env.request, pack=env.pack) + admission_storage_id = next( + key + for key, record in env.record_store.records.items() + if record.record_kind == LiveSourceIngressRecordKind.SOURCE_ADMISSION.value + ) + del env.record_store.records[admission_storage_id] + + batch = await LiveSourceResourceProjectionReader(store=env.record_store).read( + query=_query(IntelligenceResourceKind.CONNECTION, IntelligenceResourceKind.SOURCE), + after=None, + limit=20, + ) + + assert batch.records == () + assert batch.state is IntelligenceResourcePageState.DEGRADED + assert batch.degraded_reason_refs == ("degraded_reason:orphan-live-source-record",) + + +@pytest.mark.asyncio +async def test_source_health_is_not_fabricated_from_success_only_admission_records() -> None: + env = await _environment() + await env.service.admit(request=env.request, pack=env.pack) + + batch = await LiveSourceResourceProjectionReader(store=env.record_store).read( + query=_query(IntelligenceResourceKind.SOURCE_HEALTH), + after=None, + limit=20, + ) + + assert batch.records == () + assert batch.state is IntelligenceResourcePageState.DEGRADED + assert batch.degraded_reason_refs == ("degraded_reason:unsupported-source_health",)