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 payloadImmutable 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)