Module adcp.reporting.inline_storage

Borrowed-pool PostgreSQL storage for inline objects and committed replay seals.

Create either store with an open psycopg_pool.AsyncConnectionPool and call create_schema() before use; both stores share the same additive migration. check_ready() audits required catalog objects without changing the schema. Operations also audit it within their short transaction. Pools, credentials, timeouts, database durability, backups and retention remain adopter-owned. Install adcp[pg] to construct a store. Importing this module needs no driver. Catalog fingerprints are qualified on PostgreSQL 16. Other server majors are unqualified; catalog formatting differences may require separate qualification.

An authenticated caller supplies account_id. A reference is not a credential; source_scope does not add authorization. Staging deduplicates exact bytes within an account, independently of execution key or ordinal. Objects are capped at 16 MiB (a constructor may lower that cap), manifests at the source contract's 1 MiB cap. No garbage collector or destructive retention operation is supplied. Payload hashing yields between chunks of at most 1 MiB; it starts no threads.

Staging and seal insertion are separate commits. A committed seal enables exact replay after process restart, but a crash before the seal commits can refetch, even if an object already committed. These stores do not reserve requests before dispatch, recover a lost provider answer, bind unseen request/routing facts, own source leases, or compose a production service. The producer's conformance check still validates a replay against its complete frozen request. Seal validation does not establish that referenced objects exist or share this storage backend. A seal write has no lease token and does not fence a stale source owner.

Private connection participants accept detached backend preparations and return tentative references or a neutral SealedSlice. Their caller owns the same-task READ COMMITTED transaction, table locks and schema audit, account-lock ordering, cancellation/error boundary, commit/rollback and connection return. They acquire no connection, start no task/transaction, and provide no grant or commit proof.

Cancellation settles the borrowed connection before propagating a fresh, redacted CancelledError. Cancellation or a resource error may follow a committed write; resume with the same account/key or content identity to discover its outcome. An active AnyIO cancellation scope retains only a fixed framework marker so its timeout/containment semantics survive; caller text and exception chains do not. The read cancellation event is checked before and after the transaction; cancel the task to interrupt a pending database operation. An object above a store's lower configured read cap fails closed with INTEGRITY_FAILED.

Classes

class InlineStorageError (code: _Code)
Expand source code
class InlineStorageError(RuntimeError):
    """Closed diagnostic; resource failures can have an unknown commit outcome."""

    def __init__(self, code: _Code) -> None:
        self.code = code
        super().__init__(_MESSAGES[code])

Closed diagnostic; resource failures can have an unknown commit outcome.

Ancestors

  • builtins.RuntimeError
  • builtins.Exception
  • builtins.BaseException
class PgReportingSealStore (*, pool: AsyncConnectionPool)
Expand source code
class PgReportingSealStore(_PgStorage):
    """First committed seal wins for each trusted account and execution key."""

    async def _get_on(self, connection: Any, account_id: str, key: str) -> SealedSlice | None:
        row = await (
            await connection.execute(
                "SELECT staged_commit_ref, manifest_sha256, byte_count,"
                " CASE WHEN octet_length(manifest) <= %s THEN manifest END"
                " FROM reporting_inline_seals WHERE account_id = %s AND source_execution_key = %s",
                (SOURCE_BATCH_MANIFEST_MAX_BYTES_V1, account_id, key),
            )
        ).fetchone()
        if row is None:
            return None
        reference: SourceBatchManifestReferenceV1 | None
        try:
            reference = SourceBatchManifestReferenceV1(
                staged_commit_ref=row[0], manifest_sha256=row[1], byte_count=row[2]
            )
        except Exception:
            reference = None
        if reference is None:
            raise InlineStorageError("INTEGRITY_FAILED")
        return _seal(
            account_id,
            key,
            SealedSlice(reference=reference, manifest_bytes=row[3]),
            "INTEGRITY_FAILED",
        )

    @_redact
    async def get(self, *, account_id: str, source_execution_key: str) -> SealedSlice | None:
        _identity(account_id, source_execution_key)
        async with self._transaction() as connection:
            return await self._get_on(connection, account_id, source_execution_key)

    def _prepare_seal(
        self, *, account_id: str, source_execution_key: str, sealed: SealedSlice
    ) -> PreparedBackendSealV1:
        _prepared_identity(account_id, source_execution_key)
        candidate = _seal(account_id, source_execution_key, sealed, "INVALID_INPUT")
        reference = candidate.reference
        return PreparedBackendSealV1(
            account_id,
            source_execution_key,
            reference.staged_commit_ref,
            reference.manifest_sha256,
            reference.byte_count,
            candidate.manifest_bytes,
        )

    async def _put_on(self, connection: Any, prepared: PreparedBackendSealV1) -> SealedSlice:
        """Return the verified stored winner on the caller's audited transaction.

        The neutral value proves no admission, referenced-object existence or
        commit. Caller owns account locks, cancellation, settlement and any fence.
        """
        if type(prepared) is not PreparedBackendSealV1:
            raise InlineStorageError("INVALID_INPUT")
        prepared._validate()
        # Build a fresh reference, then revalidate it inside the existing closed
        # candidate validator. No mutable caller model is retained or trusted.
        reference = SourceBatchManifestReferenceV1.model_construct(
            staged_commit_ref=prepared.staged_commit_ref,
            manifest_sha256=prepared.manifest_sha256,
            byte_count=prepared.byte_count,
        )
        candidate = _seal(
            prepared.account_id,
            prepared.source_execution_key,
            SealedSlice(reference=reference, manifest_bytes=prepared.manifest_bytes),
            "INVALID_INPUT",
        )
        await connection.execute(
            "INSERT INTO reporting_inline_seals (account_id, source_execution_key,"
            " staged_commit_ref, manifest_sha256, byte_count, manifest)"
            " VALUES (%s, %s, %s, %s, %s, %s)"
            " ON CONFLICT (account_id, source_execution_key) DO NOTHING",
            (
                prepared.account_id,
                prepared.source_execution_key,
                candidate.reference.staged_commit_ref,
                candidate.reference.manifest_sha256,
                candidate.reference.byte_count,
                candidate.manifest_bytes,
            ),
        )
        winner = await self._get_on(connection, prepared.account_id, prepared.source_execution_key)
        if winner is None:
            raise InlineStorageError("INTEGRITY_FAILED")
        return winner

    @_redact
    async def put(
        self, *, account_id: str, source_execution_key: str, sealed: SealedSlice
    ) -> SealedSlice:
        prepared = self._prepare_seal(
            account_id=account_id, source_execution_key=source_execution_key, sealed=sealed
        )
        async with self._transaction() as connection:
            return await self._put_on(connection, prepared)

First committed seal wins for each trusted account and execution key.

Ancestors

  • adcp.reporting.inline_storage._PgStorage

Methods

async def get(self, *, account_id: str, source_execution_key: str) ‑> SealedSlice | None
Expand source code
@_redact
async def get(self, *, account_id: str, source_execution_key: str) -> SealedSlice | None:
    _identity(account_id, source_execution_key)
    async with self._transaction() as connection:
        return await self._get_on(connection, account_id, source_execution_key)
async def put(self, *, account_id: str, source_execution_key: str, sealed: SealedSlice) ‑> SealedSlice
Expand source code
@_redact
async def put(
    self, *, account_id: str, source_execution_key: str, sealed: SealedSlice
) -> SealedSlice:
    prepared = self._prepare_seal(
        account_id=account_id, source_execution_key=source_execution_key, sealed=sealed
    )
    async with self._transaction() as connection:
        return await self._put_on(connection, prepared)
class PgReportingStagingStore (*, pool: AsyncConnectionPool, max_payload_bytes: int = 16777216)
Expand source code
class PgReportingStagingStore(_PgStorage):
    """Immutable account-qualified content storage using a borrowed async pool."""

    def __init__(
        self, *, pool: AsyncConnectionPool, max_payload_bytes: int = INLINE_STAGING_MAX_BYTES
    ) -> None:
        if (
            type(max_payload_bytes) is not int
            or not 1 <= max_payload_bytes <= INLINE_STAGING_MAX_BYTES
        ):
            raise InlineStorageError("INVALID_INPUT")
        super().__init__(pool=pool)
        self._max_payload_bytes = max_payload_bytes

    async def _read_on(self, connection: Any, account_id: str, digest: str) -> bytes:
        row = await (
            await connection.execute(
                "SELECT CASE WHEN octet_length(payload) <= %s THEN payload END"
                " FROM reporting_inline_objects WHERE account_id = %s AND payload_sha256 = %s",
                (self._max_payload_bytes, account_id, digest),
            )
        ).fetchone()
        if row is None:
            raise InlineStorageError("NOT_FOUND")
        payload = row[0]
        if type(payload) is not bytes or await _payload_digest(payload) != digest:
            raise InlineStorageError("INTEGRITY_FAILED")
        return payload

    async def _prepare_object(
        self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes
    ) -> PreparedBackendObjectV1:
        _object_inputs(account_id, source_execution_key, ordinal, payload, self._max_payload_bytes)
        digest = await _payload_digest(payload)
        return PreparedBackendObjectV1(account_id, source_execution_key, ordinal, payload, digest)

    async def _stage_on(
        self, connection: Any, prepared: PreparedBackendObjectV1
    ) -> tuple[str, str]:
        """Stage on the caller's audited transaction; the returned identity is tentative.

        Caller owns account-lock ordering and the full transaction/error boundary.
        No checkout, task, transaction, commit, drain or fence is owned here.
        """
        if type(prepared) is not PreparedBackendObjectV1:
            raise InlineStorageError("INVALID_INPUT")
        prepared._validate(self._max_payload_bytes)
        await connection.execute(
            "INSERT INTO reporting_inline_objects (account_id, payload_sha256, payload)"
            " VALUES (%s, %s, %s) ON CONFLICT (account_id, payload_sha256) DO NOTHING",
            (prepared.account_id, prepared.payload_sha256, prepared.payload),
        )
        if (
            await self._read_on(connection, prepared.account_id, prepared.payload_sha256)
            != prepared.payload
        ):
            raise InlineStorageError("INTEGRITY_FAILED")
        return _reference(prepared.account_id, prepared.payload_sha256), prepared.payload_sha256

    @_redact
    async def stage(
        self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes
    ) -> tuple[str, str]:
        prepared = await self._prepare_object(
            account_id=account_id,
            source_execution_key=source_execution_key,
            ordinal=ordinal,
            payload=payload,
        )
        async with self._transaction() as connection:
            return await self._stage_on(connection, prepared)

    @_redact
    async def read(
        self,
        *,
        object_ref: str,
        object_generation: str,
        account_id: str,
        source_scope: Mapping[str, Any],
        cancel: asyncio.Event,
    ) -> bytes:
        _account(account_id)
        if type(object_generation) is not str or _DIGEST.fullmatch(object_generation) is None:
            raise InlineStorageError("INVALID_INPUT")
        if type(object_ref) is not str or object_ref != _reference(account_id, object_generation):
            raise InlineStorageError("NOT_FOUND")
        if cancel.is_set():
            raise asyncio.CancelledError
        async with self._transaction() as connection:
            payload = await self._read_on(connection, account_id, object_generation)
        if cancel.is_set():
            raise asyncio.CancelledError
        return payload

Immutable account-qualified content storage using a borrowed async pool.

Ancestors

  • adcp.reporting.inline_storage._PgStorage

Methods

async def read(self,
*,
object_ref: str,
object_generation: str,
account_id: str,
source_scope: Mapping[str, Any],
cancel: asyncio.Event) ‑> bytes
Expand source code
@_redact
async def read(
    self,
    *,
    object_ref: str,
    object_generation: str,
    account_id: str,
    source_scope: Mapping[str, Any],
    cancel: asyncio.Event,
) -> bytes:
    _account(account_id)
    if type(object_generation) is not str or _DIGEST.fullmatch(object_generation) is None:
        raise InlineStorageError("INVALID_INPUT")
    if type(object_ref) is not str or object_ref != _reference(account_id, object_generation):
        raise InlineStorageError("NOT_FOUND")
    if cancel.is_set():
        raise asyncio.CancelledError
    async with self._transaction() as connection:
        payload = await self._read_on(connection, account_id, object_generation)
    if cancel.is_set():
        raise asyncio.CancelledError
    return payload
async def stage(self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes) ‑> tuple[str, str]
Expand source code
@_redact
async def stage(
    self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes
) -> tuple[str, str]:
    prepared = await self._prepare_object(
        account_id=account_id,
        source_execution_key=source_execution_key,
        ordinal=ordinal,
        payload=payload,
    )
    async with self._transaction() as connection:
        return await self._stage_on(connection, prepared)