"""PostgreSQL ledger and Finance commit for direct MCP result ingestion.""" from __future__ import annotations import hashlib import json import re import secrets from datetime import datetime from typing import Any, Callable, Mapping, Optional from arr_ingestion.contracts import ( OPAQUE_ID_RE, ArtifactRef, IngestionError, ) from arr_ingestion.direct_contracts import ( DIRECT_CONTRACT_VERSION, GRANT_TOKEN_RE, DirectSubmissionReceipt, DirectSubmissionRequest, ReceivedDirectSubmission, SubmissionGrant, VerifiedDirectSubmission, canonical_json_bytes, ) from arr_ingestion.postgres import ( DatabaseConfig, PostgresIngestionRepository, _outcome_counts, ) _RUN_TERMINAL = frozenset({"accepted", "rejected", "failed", "cancelled"}) _ATTEMPT_TERMINAL = frozenset({"succeeded", "failed", "cancelled"}) _SUBMISSION_TERMINAL = frozenset( {"committed", "already_committed", "rejected", "expired"} ) _FAILURE_CODE_RE = re.compile(r"^[A-Z][A-Z0-9_]{0,63}$") def _json_mapping(value: Any, label: str) -> Mapping[str, Any]: if isinstance(value, Mapping): return value if isinstance(value, str): try: decoded = json.loads(value) except (TypeError, ValueError, json.JSONDecodeError): decoded = None if isinstance(decoded, Mapping): return decoded raise IngestionError( "DATABASE_STATE_INVALID", f"stored {label} is invalid", ) class PostgresDirectIngestionRepository(PostgresIngestionRepository): """Attempt-bound direct submission ledger with atomic Finance activation.""" def __init__( self, config: DatabaseConfig, connect: Optional[Callable[[str], Any]] = None, token_factory: Optional[Callable[[], str]] = None, ) -> None: super().__init__(config, connect=connect) self._token_factory = token_factory or ( lambda: secrets.token_urlsafe(32) ) def assert_ready(self) -> None: self._run_transaction( self._assert_ready, "direct result database is unavailable", ) @staticmethod def _assert_ready(cursor: Any) -> None: cursor.execute( """ SELECT to_regclass('ingestion.result_submission_grants'), to_regclass('ingestion.result_submissions'), to_regclass('finance.daily_versions'), to_regclass('finance.daily_records'), to_regclass('finance.current_daily_versions') """ ) row = cursor.fetchone() if row is None or any(value is None for value in row): raise IngestionError( "DATABASE_MIGRATION_MISSING", "direct result database migration is unavailable", ) def issue_grant( self, job_id: str, attempt_no: int, *, ttl_seconds: int, ) -> SubmissionGrant: if ( not isinstance(job_id, str) or not OPAQUE_ID_RE.fullmatch(job_id) or not isinstance(attempt_no, int) or isinstance(attempt_no, bool) or not 1 <= attempt_no <= 9999 or not isinstance(ttl_seconds, int) or isinstance(ttl_seconds, bool) or not 30 <= ttl_seconds <= 1800 ): raise IngestionError( "DIRECT_SUBMISSION_INVALID", "submission grant request is invalid", ) token = self._token_factory() if not isinstance(token, str) or not GRANT_TOKEN_RE.fullmatch(token): raise IngestionError( "SUBMISSION_GRANT_UNAVAILABLE", "submission grant could not be issued", retryable=True, ) return self._run_transaction( lambda cursor: self._issue_grant( cursor, job_id, attempt_no, ttl_seconds, token, ), "submission grant could not be stored", ) @staticmethod def _issue_grant( cursor: Any, job_id: str, attempt_no: int, ttl_seconds: int, token: str, ) -> SubmissionGrant: cursor.execute( """ SELECT run.id, run.pipeline_type, run.run_status, run.result_delivery_mode, attempt.id, attempt.attempt_status FROM ingestion.processing_runs AS run JOIN ingestion.processing_attempts AS attempt ON attempt.processing_run_id = run.id AND attempt.attempt_no = %s WHERE run.run_key = %s FOR UPDATE OF run, attempt """, (attempt_no, job_id), ) row = cursor.fetchone() if row is None: raise IngestionError( "JOB_NOT_FOUND", "processing attempt was not registered", ) if str(row[1]) != "opera_daily": raise IngestionError( "JOB_DELIVERY_MISMATCH", "processing attempt does not support direct results", ) if str(row[2]) in _RUN_TERMINAL or str(row[5]) in _ATTEMPT_TERMINAL: raise IngestionError( "JOB_TERMINAL", "processing attempt is already terminal", ) cursor.execute( """ SELECT 1 FROM ingestion.processing_deliveries WHERE attempt_id = %s LIMIT 1 """, (int(row[4]),), ) if cursor.fetchone() is not None: raise IngestionError( "SUBMISSION_GRANT_CONFLICT", "processing attempt already uses artifact delivery", ) cursor.execute( """ SELECT 1 FROM ingestion.result_submission_grants WHERE attempt_id = %s """, (int(row[4]),), ) if cursor.fetchone() is not None: raise IngestionError( "SUBMISSION_GRANT_CONFLICT", "processing attempt already has a submission grant", ) cursor.execute( """ UPDATE ingestion.processing_runs SET result_delivery_mode = 'direct_mcp', updated_at = now() WHERE id = %s """, (int(row[0]),), ) grant_sha256 = hashlib.sha256(token.encode("utf-8")).hexdigest() cursor.execute( """ INSERT INTO ingestion.result_submission_grants ( grant_sha256, processing_run_id, attempt_id, expires_at ) VALUES ( %s, %s, %s, now() + (%s * interval '1 second') ) RETURNING expires_at """, (grant_sha256, int(row[0]), int(row[4]), ttl_seconds), ) expires_at = cursor.fetchone()[0] if not isinstance(expires_at, datetime): raise IngestionError( "DATABASE_STATE_INVALID", "stored submission grant expiry is invalid", ) return SubmissionGrant(token, job_id, attempt_no, expires_at) def receive( self, request: DirectSubmissionRequest, ) -> ReceivedDirectSubmission: return self._run_transaction( lambda cursor: self._receive(cursor, request), "direct result could not be received", ) @staticmethod def _source_from_grant_row(row: tuple[Any, ...]) -> ArtifactRef: if str(row[15]) != "opera_xml": raise IngestionError( "DATABASE_STATE_INVALID", "stored source artifact kind is invalid", ) return ArtifactRef.from_dict( "source_xml", { "object_key": str(row[16]), "original_filename": str(row[17]), "sha256": str(row[18]), "byte_size": int(row[19]), "mime_type": str(row[20]), }, ) @classmethod def _receive( cls, cursor: Any, request: DirectSubmissionRequest, ) -> ReceivedDirectSubmission: cursor.execute( """ SELECT submission_grant.id, submission_grant.expires_at, submission_grant.consumed_at, submission_grant.revoked_at, run.id, run.run_key, run.pipeline_type, run.run_status, run.result_delivery_mode, run.source_artifact_id, run.requested_processor_version, run.requested_rule_set_sha256, attempt.id, attempt.attempt_no, attempt.attempt_status, artifact.artifact_kind, artifact.object_key, artifact.original_filename, artifact.sha256, artifact.byte_size, artifact.mime_type, now() FROM ingestion.result_submission_grants AS submission_grant JOIN ingestion.processing_runs AS run ON run.id = submission_grant.processing_run_id JOIN ingestion.processing_attempts AS attempt ON attempt.id = submission_grant.attempt_id AND attempt.processing_run_id = run.id JOIN ingestion.artifacts AS artifact ON artifact.id = run.source_artifact_id WHERE submission_grant.grant_sha256 = %s FOR UPDATE OF submission_grant, run, attempt """, (request.grant_sha256,), ) row = cursor.fetchone() if row is None: raise IngestionError( "SUBMISSION_GRANT_INVALID", "submission grant is invalid", ) if ( str(row[5]) != request.job_id or int(row[13]) != request.attempt_no ): raise IngestionError( "SUBMISSION_GRANT_INVALID", "submission grant is invalid", ) source = cls._source_from_grant_row(row) cursor.execute( """ SELECT id, submission_key, submission_status, contract_version, business_date, processor_version, rule_set_sha256, result_schema_version, payload_json, payload_sha256, payload_byte_size, record_count, receipt_json, failure_code FROM ingestion.result_submissions WHERE grant_id = %s FOR UPDATE """, (int(row[0]),), ) existing = cursor.fetchone() if existing is not None: cls._require_same_submission(existing, request) status = str(existing[2]) if status in {"committed", "already_committed"}: receipt = DirectSubmissionReceipt.from_dict( _json_mapping(existing[12], "direct receipt") ) if ( receipt.job_id != request.job_id or receipt.attempt_no != request.attempt_no or receipt.record_count != request.record_count ): raise IngestionError( "DATABASE_STATE_INVALID", "stored direct receipt identity is invalid", ) return ReceivedDirectSubmission( int(existing[0]), status, request, source, receipt, ) if status == "rejected": raise IngestionError( "SUBMISSION_REJECTED", "direct result was previously rejected", ) if status == "expired": raise IngestionError( "SUBMISSION_EXPIRED", "direct result submission expired", ) if status == "received": cursor.execute( """ UPDATE ingestion.result_submissions SET submission_status = 'validating', validation_started_at = now() WHERE id = %s """, (int(existing[0]),), ) status = "validating" if status != "validating" or existing[8] is None: raise IngestionError( "DATABASE_STATE_INVALID", "stored direct submission state is invalid", ) return ReceivedDirectSubmission( int(existing[0]), status, request, source, ) expires_at = row[1] now_value = row[21] if row[3] is not None: raise IngestionError( "SUBMISSION_GRANT_INVALID", "submission grant is invalid", ) if not isinstance(expires_at, datetime) or not isinstance( now_value, datetime, ): raise IngestionError( "DATABASE_STATE_INVALID", "stored submission grant time is invalid", ) if expires_at <= now_value: raise IngestionError( "SUBMISSION_GRANT_EXPIRED", "submission grant expired", ) if row[2] is not None: raise IngestionError( "DATABASE_STATE_INVALID", "consumed submission grant has no ledger", ) if ( str(row[6]) != "opera_daily" or str(row[8]) != "direct_mcp" or str(row[7]) in _RUN_TERMINAL or str(row[14]) in _ATTEMPT_TERMINAL ): raise IngestionError( "JOB_TERMINAL", "processing attempt cannot accept a direct result", ) if ( str(row[10]) != request.processor_version or str(row[11]) != request.rule_set_sha256 ): raise IngestionError( "JOB_DELIVERY_MISMATCH", "direct result does not match its processing job", ) cursor.execute( """ INSERT INTO ingestion.result_submissions ( submission_key, grant_id, processing_run_id, attempt_id, submission_status, contract_version, business_date, processor_version, rule_set_sha256, result_schema_version, payload_json, payload_sha256, payload_byte_size, record_count, expires_at, validation_started_at ) VALUES ( %s, %s, %s, %s, 'validating', %s, %s, %s, %s, %s, %s::jsonb, %s, %s, %s, %s, now() ) RETURNING id """, ( request.submission_key, int(row[0]), int(row[4]), int(row[12]), DIRECT_CONTRACT_VERSION, request.business_date, request.processor_version, request.rule_set_sha256, request.result_schema_version, request.payload_json(), request.payload_sha256, len(request.payload_bytes), request.record_count, expires_at, ), ) submission_id = int(cursor.fetchone()[0]) cursor.execute( """ UPDATE ingestion.result_submission_grants SET consumed_at = now() WHERE id = %s """, (int(row[0]),), ) delivery_json = cls._delivery_json(request) cursor.execute( """ UPDATE ingestion.processing_runs SET run_status = 'validating', result_delivery_mode = 'direct_mcp', result_artifact_id = NULL, delivered_processor_version = %s, delivered_rule_set_sha256 = %s, result_schema_version = %s, delivery_sha256 = %s, delivery_json = %s::jsonb, business_date = %s, failure_code = NULL, failure_message = NULL, updated_at = now(), validated_at = now(), finished_at = NULL WHERE id = %s """, ( request.processor_version, request.rule_set_sha256, request.result_schema_version, request.payload_sha256, delivery_json, request.business_date, int(row[4]), ), ) cursor.execute( """ UPDATE ingestion.processing_attempts SET attempt_status = 'delivered', failure_code = NULL, failure_message = NULL, finished_at = NULL WHERE id = %s """, (int(row[12]),), ) return ReceivedDirectSubmission( submission_id, "validating", request, source, ) @staticmethod def _require_same_submission( row: tuple[Any, ...], request: DirectSubmissionRequest, ) -> None: same = ( str(row[1]) == request.submission_key and str(row[3]) == DIRECT_CONTRACT_VERSION and row[4] == request.business_date and str(row[5]) == request.processor_version and str(row[6]) == request.rule_set_sha256 and str(row[7]) == request.result_schema_version and str(row[9]) == request.payload_sha256 and int(row[10]) == len(request.payload_bytes) and int(row[11]) == request.record_count ) if same and row[8] is not None: stored_payload = _json_mapping(row[8], "direct payload") same = canonical_json_bytes(stored_payload) == request.payload_bytes if not same: raise IngestionError( "SUBMISSION_CONFLICT", "processing attempt was resubmitted with different parameters", ) @staticmethod def _delivery_json(request: DirectSubmissionRequest) -> str: return json.dumps( { "contract_version": DIRECT_CONTRACT_VERSION, "submission_key": request.submission_key, "payload_sha256": request.payload_sha256, "record_count": request.record_count, }, ensure_ascii=False, sort_keys=True, separators=(",", ":"), ) def commit( self, verified: VerifiedDirectSubmission, ) -> DirectSubmissionReceipt: return self._run_transaction( lambda cursor: self._commit_direct(cursor, verified), "validated direct result was not committed", ) def _commit_direct( self, cursor: Any, verified: VerifiedDirectSubmission, ) -> DirectSubmissionReceipt: state = verified.submission request = state.request cursor.execute( """ SELECT submission.submission_status, submission.payload_sha256, submission.payload_json, submission.business_date, submission.processor_version, submission.rule_set_sha256, submission.result_schema_version, submission.record_count, submission.receipt_json, submission.submission_key, run.id, run.run_key, run.source_artifact_id, run.run_status, run.result_delivery_mode, attempt.id, attempt.attempt_no, attempt.attempt_status FROM ingestion.result_submissions AS submission JOIN ingestion.processing_runs AS run ON run.id = submission.processing_run_id JOIN ingestion.processing_attempts AS attempt ON attempt.id = submission.attempt_id AND attempt.processing_run_id = run.id WHERE submission.id = %s FOR UPDATE OF submission, run, attempt """, (state.submission_id,), ) row = cursor.fetchone() if row is None: raise IngestionError( "SUBMISSION_NOT_FOUND", "direct result submission was not found", ) status = str(row[0]) if status in {"committed", "already_committed"}: return DirectSubmissionReceipt.from_dict( _json_mapping(row[8], "direct receipt") ) if status == "rejected": raise IngestionError( "SUBMISSION_REJECTED", "direct result was previously rejected", ) if status == "expired": raise IngestionError( "SUBMISSION_EXPIRED", "direct result submission expired", ) if ( status != "validating" or str(row[1]) != request.payload_sha256 or row[2] is None or row[3] != request.business_date or str(row[4]) != request.processor_version or str(row[5]) != request.rule_set_sha256 or str(row[6]) != request.result_schema_version or int(row[7]) != request.record_count or str(row[9]) != request.submission_key or str(row[11]) != request.job_id or int(row[16]) != request.attempt_no or str(row[14]) != "direct_mcp" ): raise IngestionError( "DATABASE_STATE_INVALID", "stored direct submission identity is invalid", ) if canonical_json_bytes( _json_mapping(row[2], "direct payload") ) != request.payload_bytes: raise IngestionError( "DATABASE_STATE_INVALID", "stored direct payload identity is invalid", ) if str(row[13]) in _RUN_TERMINAL or str(row[17]) in _ATTEMPT_TERMINAL: raise IngestionError( "JOB_TERMINAL", "processing attempt is already terminal", ) business_date = request.business_date source_artifact_id = int(row[12]) cursor.execute( "SELECT pg_advisory_xact_lock(hashtextextended(%s, 0))", (f"finance-daily:{business_date.isoformat()}",), ) cursor.execute( """ SELECT id, version_no, version_status FROM finance.daily_versions WHERE source_artifact_id = %s AND business_date = %s AND processor_version = %s AND rule_set_sha256 = %s FOR UPDATE """, ( source_artifact_id, business_date, request.processor_version, request.rule_set_sha256, ), ) existing = cursor.fetchone() if existing is not None: if str(existing[2]) not in {"active", "superseded"}: raise IngestionError( "DATABASE_STATE_INVALID", "existing daily version is invalid", ) daily_version_id = int(existing[0]) version_no = int(existing[1]) disposition = "already_committed" else: payload = request.payload counts = _outcome_counts(payload) records = payload.get("records") channels = payload.get("channels") if not isinstance(records, list) or not isinstance(channels, list): raise IngestionError( "RESULT_CONTRACT_INVALID", "direct result arrays are invalid", ) cursor.execute( """ SELECT COALESCE(MAX(version_no), 0) + 1 FROM finance.daily_versions WHERE business_date = %s """, (business_date,), ) version_no = int(cursor.fetchone()[0]) cursor.execute( """ INSERT INTO finance.daily_versions ( business_date, version_no, processing_run_id, source_artifact_id, version_status, processor_version, rule_set_sha256, result_schema_version, result_sha256, source_rows, retained_rows, excluded_rate_code_rows, duplicate_rows, validation_failed_rows, price_unmatched_rows, validated_at, result_delivery_mode ) VALUES ( %s, %s, %s, %s, 'validated', %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now(), 'direct_mcp' ) RETURNING id """, ( business_date, version_no, int(row[10]), source_artifact_id, request.processor_version, request.rule_set_sha256, request.result_schema_version, request.payload_sha256, int(payload["source_rows"]), *counts, ), ) daily_version_id = int(cursor.fetchone()[0]) booking_links = self._load_booking_links(cursor, records) self._insert_records( cursor, daily_version_id, records, booking_links, ) self._insert_channel_metrics( cursor, daily_version_id, channels, ) self._activate_daily_version( cursor, business_date, daily_version_id, ) disposition = "committed" receipt = DirectSubmissionReceipt( disposition, request.job_id, request.attempt_no, business_date, daily_version_id, version_no, request.record_count, ) cursor.execute( """ UPDATE ingestion.result_submissions SET submission_status = %s, daily_version_id = %s, receipt_json = %s::jsonb, payload_json = NULL, finished_at = now(), payload_purged_at = now(), failure_code = NULL, failure_message = NULL WHERE id = %s """, ( disposition, daily_version_id, json.dumps( receipt.to_dict(), ensure_ascii=False, sort_keys=True, separators=(",", ":"), ), state.submission_id, ), ) cursor.execute( """ UPDATE ingestion.processing_attempts SET attempt_status = 'succeeded', failure_code = NULL, failure_message = NULL, finished_at = now() WHERE id = %s """, (int(row[15]),), ) cursor.execute( """ UPDATE ingestion.processing_runs SET run_status = 'accepted', result_delivery_mode = 'direct_mcp', result_artifact_id = NULL, delivered_processor_version = %s, delivered_rule_set_sha256 = %s, result_schema_version = %s, delivery_sha256 = %s, delivery_json = %s::jsonb, business_date = %s, failure_code = NULL, failure_message = NULL, updated_at = now(), validated_at = now(), finished_at = now() WHERE id = %s """, ( request.processor_version, request.rule_set_sha256, request.result_schema_version, request.payload_sha256, self._delivery_json(request), business_date, int(row[10]), ), ) self._insert_outbox( cursor, f"processing-run:{int(row[10])}:accepted", "processing_run", int(row[10]), "arr.daily_version_committed", { "job_id": request.job_id, "business_date": business_date.isoformat(), "daily_version_id": daily_version_id, "version_no": version_no, "disposition": disposition, "delivery_mode": "direct_mcp", }, ) return receipt def reject( self, submission: ReceivedDirectSubmission, failure_code: str, ) -> None: code = ( failure_code if isinstance(failure_code, str) and _FAILURE_CODE_RE.fullmatch(failure_code) else "DIRECT_RESULT_REJECTED" ) self._run_transaction( lambda cursor: self._reject(cursor, submission, code), "direct result rejection could not be stored", ) def _reject( self, cursor: Any, submission: ReceivedDirectSubmission, failure_code: str, ) -> None: cursor.execute( """ SELECT submission.submission_status, submission.payload_sha256, run.id, run.run_key, run.run_status, attempt.id, attempt.attempt_status FROM ingestion.result_submissions AS submission JOIN ingestion.processing_runs AS run ON run.id = submission.processing_run_id JOIN ingestion.processing_attempts AS attempt ON attempt.id = submission.attempt_id AND attempt.processing_run_id = run.id WHERE submission.id = %s FOR UPDATE OF submission, run, attempt """, (submission.submission_id,), ) row = cursor.fetchone() if row is None: raise IngestionError( "SUBMISSION_NOT_FOUND", "direct result submission was not found", ) status = str(row[0]) if status in _SUBMISSION_TERMINAL: return if str(row[1]) != submission.request.payload_sha256: raise IngestionError( "SUBMISSION_CONFLICT", "direct result submission identity conflicts", ) cursor.execute( """ UPDATE ingestion.result_submissions SET submission_status = 'rejected', failure_code = %s, failure_message = 'direct result validation failed', payload_json = NULL, finished_at = now(), payload_purged_at = now() WHERE id = %s """, (failure_code, submission.submission_id), ) if str(row[6]) not in _ATTEMPT_TERMINAL: cursor.execute( """ UPDATE ingestion.processing_attempts SET attempt_status = 'failed', failure_code = %s, failure_message = 'direct result validation failed', finished_at = now() WHERE id = %s """, (failure_code, int(row[5])), ) if str(row[4]) not in _RUN_TERMINAL: cursor.execute( """ UPDATE ingestion.processing_runs SET run_status = 'failed', failure_code = %s, failure_message = 'direct result validation failed', updated_at = now(), finished_at = now() WHERE id = %s """, (failure_code, int(row[2])), ) self._insert_outbox( cursor, f"processing-run:{int(row[2])}:failed", "processing_run", int(row[2]), "arr.processing_failed", { "job_id": str(row[3]), "failure_code": failure_code, "delivery_mode": "direct_mcp", }, ) def expire_stale(self, *, limit: int) -> int: if ( not isinstance(limit, int) or isinstance(limit, bool) or not 1 <= limit <= 1000 ): raise IngestionError( "DIRECT_SUBMISSION_INVALID", "direct submission expiry limit is invalid", ) return self._run_transaction( lambda cursor: self._expire_stale(cursor, limit), "expired direct submissions could not be cleaned", ) def _expire_stale(self, cursor: Any, limit: int) -> int: cursor.execute( """ SELECT submission.id, run.id, run.run_key, run.run_status, attempt.id, attempt.attempt_status FROM ingestion.result_submissions AS submission JOIN ingestion.processing_runs AS run ON run.id = submission.processing_run_id JOIN ingestion.processing_attempts AS attempt ON attempt.id = submission.attempt_id AND attempt.processing_run_id = run.id WHERE submission.submission_status IN ('received', 'validating') AND submission.expires_at <= now() ORDER BY submission.id LIMIT %s FOR UPDATE OF submission, run, attempt SKIP LOCKED """, (limit,), ) rows = list(cursor.fetchall()) for row in rows: cursor.execute( """ UPDATE ingestion.result_submissions SET submission_status = 'expired', failure_code = 'SUBMISSION_EXPIRED', failure_message = 'direct result submission expired', payload_json = NULL, finished_at = now(), payload_purged_at = now() WHERE id = %s """, (int(row[0]),), ) if str(row[5]) not in _ATTEMPT_TERMINAL: cursor.execute( """ UPDATE ingestion.processing_attempts SET attempt_status = 'failed', failure_code = 'SUBMISSION_EXPIRED', failure_message = 'direct result submission expired', finished_at = now() WHERE id = %s """, (int(row[4]),), ) if str(row[3]) not in _RUN_TERMINAL: cursor.execute( """ UPDATE ingestion.processing_runs SET run_status = 'failed', failure_code = 'SUBMISSION_EXPIRED', failure_message = 'direct result submission expired', updated_at = now(), finished_at = now() WHERE id = %s """, (int(row[1]),), ) self._insert_outbox( cursor, f"processing-run:{int(row[1])}:failed", "processing_run", int(row[1]), "arr.processing_failed", { "job_id": str(row[2]), "failure_code": "SUBMISSION_EXPIRED", "delivery_mode": "direct_mcp", }, ) return len(rows)