Module adcp.reporting.receipts.handler

One authenticated composition path for direct MCP/A2A and hydrated contexts.

Global variables

var ReceiptAccountResolver

Resolve AND reauthorize the exact account reference for this consumer on every call.

Return the canonical storage account ID, never a RequestContext cache key. A natural-key account reference is resolved by this same application ACL boundary. Raises ReportingReceiptError('UNAUTHORIZED') for unknown or denied accounts.

Classes

class ReportingReceiptHandler (store: ReportingReceiptBatchStore,
*,
resolve_account: ReceiptAccountResolver,
buyer_agents: BuyerAgentRegistry | None = None,
consumer_status_enabled: bool = False,
adcp_version: str | None = None)
Expand source code
class ReportingReceiptHandler(ADCPHandler[ToolContext]):
    """Mount receipts and the optional frozen feed; tier activation is separately gated.

    ``resolve_account`` is an application ACL, called even for a completed
    batch. ``buyer_agents`` optionally re-resolves API/OAuth/signed commercial
    identity on every call. Neither transport tenancy nor body fields supply
    the consumer. Authentication middleware must populate trusted context.
    """

    advertised_tools = {TASK, "get_reporting_status"}

    def __init__(
        self,
        store: ReportingReceiptBatchStore,
        *,
        resolve_account: ReceiptAccountResolver,
        buyer_agents: BuyerAgentRegistry | None = None,
        consumer_status_enabled: bool = False,
        adcp_version: str | None = None,
    ) -> None:
        super().__init__()
        if not isinstance(store, ReportingReceiptBatchStore):
            raise TypeError("receipt ingress requires an atomic ReportingReceiptBatchStore")
        self.receipt_store = store
        self._receipt_account_resolver = resolve_account
        self._receipt_registry = buyer_agents
        from adcp.reporting.feed.store import ReportingFeedStore

        self.reporting_feed_store: ReportingFeedStore | None = (
            store if isinstance(store, ReportingFeedStore) else None
        )
        self._feed_consumer_status_enabled = consumer_status_enabled
        self._adcp_version = resolve_adcp_version(adcp_version)
        if not is_adcp_version_at_least(self._adcp_version, "3.2-rc.6"):
            raise ConfigurationError(
                "reporting mounts require a supported AdCP 3.2 reporting contract; "
                "use 3.2 or omit the pin for the packaged default"
            )

    def get_adcp_version(self) -> str:
        """Select rendering and advertised MCP/A2A schemas for this mount."""
        return self._adcp_version

    def advertised_tools_for_instance(self) -> set[str]:
        return {TASK, "get_reporting_status"} if self.reporting_feed_store is not None else {TASK}

    async def get_reporting_status(
        self,
        params: GetReportingStatusRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any] | NotImplementedResponse:
        """One authenticated mount; the optional store freezes periods walks."""
        from adcp.reporting.feed.errors import ReportingFeedError
        from adcp.reporting.feed.request import FeedRequest
        from adcp.reporting.ledger.status import ReportingStatusCaller, ReportingStatusHandler
        from adcp.reporting.ledger.store import LedgerConflictError, ReportingLedgerStore

        if self.reporting_feed_store is None:
            return self._not_supported("get_reporting_status")
        request = (
            dict(params)
            if isinstance(params, dict)
            else params.model_dump(mode="json", exclude_unset=True)
        )
        if context is not None and context.resolved_adcp_version is not None:
            request["adcp_version"] = context.resolved_adcp_version
        else:
            request.setdefault("adcp_version", self.get_adcp_version())
        try:
            if request.get("view") == "periods":
                FeedRequest.parse(request)

            async def authorize() -> ReportingDeliveryPrincipal:
                if context is None:
                    raise ReportingFeedError("UNAUTHORIZED")
                try:
                    consumer = await _consumer(context, self._receipt_registry)
                    account = await self._receipt_account_resolver(
                        dict(request["account"]), context, consumer
                    )
                    if isinstance(context, RequestContext) and context.account.id != account:
                        raise ReportingFeedError("UNAUTHORIZED")
                    return ReportingDeliveryPrincipal(account, consumer)
                except (ReportingReceiptError, ReportingNotificationError, ValueError, TypeError):
                    raise ReportingFeedError("UNAUTHORIZED") from None

            caller = await authorize()

            async def reauthorize() -> None:
                # A still-authorized alias must not change which account or
                # canonical consumer owns the already captured boundary.
                if await authorize() != caller:
                    raise ReportingFeedError("UNAUTHORIZED")

            if request.get("view") == "periods":
                response = await self.reporting_feed_store.read_reporting_feed(
                    request,
                    caller=caller,
                    consumer_status_enabled=self._feed_consumer_status_enabled,
                    reauthorize=reauthorize,
                )
                await reauthorize()
                return response
            pagination = request.get("pagination")
            positions = (
                request.get("changes_after"),
                pagination.get("cursor") if isinstance(pagination, dict) else None,
            )
            if any(isinstance(p, str) and p.startswith("rpf1.") for p in positions):
                # A view change cannot send a frozen position to the mutable
                # legacy projector. Leave ordinary Core requests compatible.
                raise ReportingFeedError("INVALID_CHECKPOINT")
            if isinstance(self.receipt_store, ReportingLedgerStore):
                try:
                    from adcp.reporting.projection.wire import ReportingTierStatusStore

                    if isinstance(self.receipt_store, ReportingTierStatusStore):
                        response = await self.receipt_store.read_tier_status(
                            request,
                            caller=caller,
                            consumer_status_enabled=self._feed_consumer_status_enabled,
                        )
                    else:
                        response = await ReportingStatusHandler(
                            self.receipt_store,
                            consumer_status_enabled=self._feed_consumer_status_enabled,
                        ).handle(
                            request,
                            caller=ReportingStatusCaller(caller.account_id, caller.consumer_id),
                        )
                    await reauthorize()
                    return response
                except LedgerConflictError as error:
                    # Only the legacy Core projector exposes its established
                    # domain errors. ACL/provider/feed failures stay redacted.
                    code, message = error.code, str(error)
            else:
                return self._not_supported("get_reporting_status")
        except ReportingFeedError as error:
            code, message = error.code, str(error)
        except Exception:
            unavailable = ReportingFeedError("REPORTING_FEED_STORAGE_UNAVAILABLE")
            code, message = unavailable.code, str(unavailable)
        raise ADCPTaskError(
            operation="get_reporting_status", errors=[Error(code=code, message=message)]
        )

    async def sync_reporting_receipts(
        self,
        params: SyncReportingReceiptsRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any]:
        request = (
            params
            if isinstance(params, dict)
            else params.model_dump(mode="json", exclude_unset=True)
        )
        try:
            validate_receipt_request(request)
            if context is None:
                raise ReportingReceiptError("UNAUTHORIZED")
            try:
                consumer = await _consumer(context, self._receipt_registry)
            except ReportingNotificationError:
                raise ReportingReceiptError("UNAUTHORIZED") from None
            account = await self._receipt_account_resolver(
                dict(request["account"]), context, consumer
            )
            if isinstance(context, RequestContext) and context.account.id != account:
                raise ReportingReceiptError("UNAUTHORIZED")
            try:
                caller = ReportingDeliveryPrincipal(account, consumer)
            except (ValueError, TypeError):
                raise ReportingReceiptError("UNAUTHORIZED") from None
            return await self.receipt_store.ingest_receipt_batch(request, caller=caller)
        except ReportingReceiptError as error:
            code, message = error.code, str(error)
        except Exception as error:
            _storage_failure(error, boundary="handler")
            unavailable = ReportingReceiptError("RECEIPT_STORAGE_UNAVAILABLE")
            code, message = unavailable.code, str(unavailable)
        # Leave the exception scope before translating. Credential/ACL adapters
        # may raise provider errors; only safe origin coordinates were logged.
        raise ADCPTaskError(operation=TASK, errors=[Error(code=code, message=message)])

Mount receipts and the optional frozen feed; tier activation is separately gated.

resolve_account is an application ACL, called even for a completed batch. buyer_agents optionally re-resolves API/OAuth/signed commercial identity on every call. Neither transport tenancy nor body fields supply the consumer. Authentication middleware must populate trusted context.

Ancestors

Subclasses

Class variables

var advertised_tools

Methods

def advertised_tools_for_instance(self) ‑> set[str]
Expand source code
def advertised_tools_for_instance(self) -> set[str]:
    return {TASK, "get_reporting_status"} if self.reporting_feed_store is not None else {TASK}
def get_adcp_version(self) ‑> str
Expand source code
def get_adcp_version(self) -> str:
    """Select rendering and advertised MCP/A2A schemas for this mount."""
    return self._adcp_version

Select rendering and advertised MCP/A2A schemas for this mount.

async def get_reporting_status(self,
params: GetReportingStatusRequest | dict[str, Any],
context: ToolContext | None = None) ‑> dict[str, typing.Any] | NotImplementedResponse
Expand source code
async def get_reporting_status(
    self,
    params: GetReportingStatusRequest | dict[str, Any],
    context: ToolContext | None = None,
) -> dict[str, Any] | NotImplementedResponse:
    """One authenticated mount; the optional store freezes periods walks."""
    from adcp.reporting.feed.errors import ReportingFeedError
    from adcp.reporting.feed.request import FeedRequest
    from adcp.reporting.ledger.status import ReportingStatusCaller, ReportingStatusHandler
    from adcp.reporting.ledger.store import LedgerConflictError, ReportingLedgerStore

    if self.reporting_feed_store is None:
        return self._not_supported("get_reporting_status")
    request = (
        dict(params)
        if isinstance(params, dict)
        else params.model_dump(mode="json", exclude_unset=True)
    )
    if context is not None and context.resolved_adcp_version is not None:
        request["adcp_version"] = context.resolved_adcp_version
    else:
        request.setdefault("adcp_version", self.get_adcp_version())
    try:
        if request.get("view") == "periods":
            FeedRequest.parse(request)

        async def authorize() -> ReportingDeliveryPrincipal:
            if context is None:
                raise ReportingFeedError("UNAUTHORIZED")
            try:
                consumer = await _consumer(context, self._receipt_registry)
                account = await self._receipt_account_resolver(
                    dict(request["account"]), context, consumer
                )
                if isinstance(context, RequestContext) and context.account.id != account:
                    raise ReportingFeedError("UNAUTHORIZED")
                return ReportingDeliveryPrincipal(account, consumer)
            except (ReportingReceiptError, ReportingNotificationError, ValueError, TypeError):
                raise ReportingFeedError("UNAUTHORIZED") from None

        caller = await authorize()

        async def reauthorize() -> None:
            # A still-authorized alias must not change which account or
            # canonical consumer owns the already captured boundary.
            if await authorize() != caller:
                raise ReportingFeedError("UNAUTHORIZED")

        if request.get("view") == "periods":
            response = await self.reporting_feed_store.read_reporting_feed(
                request,
                caller=caller,
                consumer_status_enabled=self._feed_consumer_status_enabled,
                reauthorize=reauthorize,
            )
            await reauthorize()
            return response
        pagination = request.get("pagination")
        positions = (
            request.get("changes_after"),
            pagination.get("cursor") if isinstance(pagination, dict) else None,
        )
        if any(isinstance(p, str) and p.startswith("rpf1.") for p in positions):
            # A view change cannot send a frozen position to the mutable
            # legacy projector. Leave ordinary Core requests compatible.
            raise ReportingFeedError("INVALID_CHECKPOINT")
        if isinstance(self.receipt_store, ReportingLedgerStore):
            try:
                from adcp.reporting.projection.wire import ReportingTierStatusStore

                if isinstance(self.receipt_store, ReportingTierStatusStore):
                    response = await self.receipt_store.read_tier_status(
                        request,
                        caller=caller,
                        consumer_status_enabled=self._feed_consumer_status_enabled,
                    )
                else:
                    response = await ReportingStatusHandler(
                        self.receipt_store,
                        consumer_status_enabled=self._feed_consumer_status_enabled,
                    ).handle(
                        request,
                        caller=ReportingStatusCaller(caller.account_id, caller.consumer_id),
                    )
                await reauthorize()
                return response
            except LedgerConflictError as error:
                # Only the legacy Core projector exposes its established
                # domain errors. ACL/provider/feed failures stay redacted.
                code, message = error.code, str(error)
        else:
            return self._not_supported("get_reporting_status")
    except ReportingFeedError as error:
        code, message = error.code, str(error)
    except Exception:
        unavailable = ReportingFeedError("REPORTING_FEED_STORAGE_UNAVAILABLE")
        code, message = unavailable.code, str(unavailable)
    raise ADCPTaskError(
        operation="get_reporting_status", errors=[Error(code=code, message=message)]
    )

One authenticated mount; the optional store freezes periods walks.

Inherited members