From 3022dc0d8183f00dd48a06381d28624805700f16 Mon Sep 17 00:00:00 2001 From: Owie Schon Date: Sat, 11 Jul 2026 11:21:48 -0400 Subject: [PATCH] feat: add durable completion-gated delivery --- data/m5_baseline_defects.json | 1 - data/m5_legacy_schema_lock.json | 14 +- docs/DECISION_LOG.md | 33 ++ docs/M5_BUILD_JOURNAL.md | 39 ++ pyproject.toml | 1 + scripts/m5_legacy_inventory.py | 2 +- scripts/quote_worker.py | 156 ++++++++ src/gateway/conversation_store.py | 2 +- src/quoting/__init__.py | 15 +- src/quoting/callbacks.py | 105 +++++ src/quoting/delivery.py | 193 +++++++++- src/quoting/document_guard.py | 13 +- src/quoting/models.py | 20 +- src/quoting/store.py | 513 +++++++++++++++++++++++-- src/quoting/worker.py | 259 ++++++++----- src/runtime/app.py | 81 ++++ src/runtime/config.py | 34 +- src/runtime/post_call.py | 9 +- tests/gateway_fixtures.py | 8 + tests/test_delivery_adapters.py | 96 +++++ tests/test_delivery_callback_routes.py | 63 +++ tests/test_delivery_callbacks.py | 60 +++ tests/test_delivery_runtime_config.py | 34 ++ tests/test_m5_baseline_defects.py | 15 +- tests/test_m5_release_config.py | 3 +- tests/test_m5_w5_delivery_worker.py | 227 +++++++++++ tests/test_post_call.py | 45 +++ tests/test_quote_artifacts.py | 2 +- tests/test_quote_migration.py | 4 +- tests/test_quote_worker.py | 42 +- 30 files changed, 1906 insertions(+), 183 deletions(-) create mode 100644 scripts/quote_worker.py create mode 100644 src/quoting/callbacks.py create mode 100644 tests/test_delivery_adapters.py create mode 100644 tests/test_delivery_callback_routes.py create mode 100644 tests/test_delivery_callbacks.py create mode 100644 tests/test_delivery_runtime_config.py create mode 100644 tests/test_m5_w5_delivery_worker.py diff --git a/data/m5_baseline_defects.json b/data/m5_baseline_defects.json index 188596a..7407565 100644 --- a/data/m5_baseline_defects.json +++ b/data/m5_baseline_defects.json @@ -3,7 +3,6 @@ "baseline_commit": "8001643fe503254ed9f12f0996894a49efe52342", "contract": "A probe passes only when its positive control works and the named baseline defect is behaviorally reproduced.", "defects": [ - {"id": "M5-DEF-08", "owner_wave": "W5", "test": "tests/test_m5_baseline_defects.py::test_defect_queue_claim_is_not_exclusive", "expected_reason": "same_pending_job_claimed_twice"}, {"id": "M5-DEF-11", "owner_wave": "W7", "test": "tests/test_m5_baseline_defects.py::test_defect_latency_harness_rejects_chunked_readback", "expected_reason": "latency_harness_expects_immediate_quote_number"} ] } diff --git a/data/m5_legacy_schema_lock.json b/data/m5_legacy_schema_lock.json index 1dd8f8b..9ccfe35 100644 --- a/data/m5_legacy_schema_lock.json +++ b/data/m5_legacy_schema_lock.json @@ -26,13 +26,14 @@ } }, "quote_store": { - "schema_sha256": "d26de5eaaf158fd47094f886b9b94a94c21578b7f5bdf80a0cdeb22f3410a6d8", + "schema_sha256": "b78342c2f525a361338c7fb82102650af72229aeb2bff51e17c2ced7ec1fd4a8", "tables": [ "artifact_guard_events", "artifact_jobs", "capability_access_events", "delivery_jobs", "issuance_events", + "provider_callback_events", "quote_artifacts", "quote_capabilities", "quote_idempotency", @@ -42,7 +43,8 @@ "quote_revision_seq", "quote_schema", "quote_seq", - "quotes" + "quotes", + "worker_heartbeats" ], "row_counts": { "artifact_guard_events": 0, @@ -50,6 +52,7 @@ "capability_access_events": 0, "delivery_jobs": 0, "issuance_events": 0, + "provider_callback_events": 0, "quote_artifacts": 0, "quote_capabilities": 0, "quote_idempotency": 1, @@ -59,14 +62,15 @@ "quote_revision_seq": 0, "quote_schema": 1, "quote_seq": 1, - "quotes": 1 + "quotes": 1, + "worker_heartbeats": 0 }, "known_gaps": [ "pre-W3 compatibility tables remain until migration evidence is retained", - "delivery job leases and authenticated call completion land in W5", + "credentialed provider receipts remain external to the schema fixture", "legacy rows without assent or artifact proof remain explicitly held" ] } }, - "combined_sha256": "1a6d06e44585df041fe2fdb6b4869c4446f1309eaa9ae512dfeb2bf90efc55db" + "combined_sha256": "151432b988d1ace9cd6cd2af51039246f68d33f59a8977032cc678b808f592fb" } diff --git a/docs/DECISION_LOG.md b/docs/DECISION_LOG.md index 1c2139c..98dbf7f 100644 --- a/docs/DECISION_LOG.md +++ b/docs/DECISION_LOG.md @@ -1577,3 +1577,36 @@ migration tests pass offline. Production remains blocked on existing staging agent/tool/branch IDs, credentials, real hosted test IDs, an inbound call, credentialed webhook latency/loss evidence, and one exercised rollback receipt. **Status:** implementation locked; external staging proof outstanding. + +## 2026-07-11 — D21: delivery is at-least-once before send and manual after uncertainty + +**Decision.** Artifact and notification work are separate durable state +machines. An artifact cannot leave `waiting_for_call_end` until an authenticated +completion event or bounded reconciliation releases it. Rendering is leased and +retryable; verified PDF/XLSX bytes are immutable and release per-channel +delivery jobs only after document parity. + +A delivery lease has two recovery domains. A crash while merely `leased` is +safe to retry. Crossing into `sending` marks provider acceptance as potentially +ambiguous; a crash or connection loss from there enters `manual_review`, never +an automatic duplicate send. Provider acceptance is not called “delivered.” +Signed callbacks advance accepted work to delivered, bounced, undelivered, or +permanent failure and are replay-safe and non-regressing. + +Frozen contact identity, version, channel, and destination fingerprint are +revalidated immediately before provider I/O. Destination or unapproved metadata +changes require fresh caller authorization. A quote may freeze multiple selected +channel preflights; they share exactly one artifact job and retain one delivery +job per channel. + +**Why.** Neither SQLite nor external providers can offer a meaningful global +exactly-once transaction. The safe claim is explicit: retry only when no request +was accepted, preserve stable provider keys, and surface every ambiguous outcome +for human review. + +**Outcome.** The worker uses atomic `UPDATE ... RETURNING` leases, expiry +recovery, parity gates, stable keys, credential-gated SendGrid/Twilio adapters, +recipient allowlists, exact attachments/link-only SMS, signed callbacks, +heartbeats, bounded polling, graceful drain, and operator inspect/retry/ +quarantine/revoke controls. **Status:** implemented and verified offline; +credentialed staging sends and callbacks remain under P04. diff --git a/docs/M5_BUILD_JOURNAL.md b/docs/M5_BUILD_JOURNAL.md index dd0ccb5..f78954c 100644 --- a/docs/M5_BUILD_JOURNAL.md +++ b/docs/M5_BUILD_JOURNAL.md @@ -735,3 +735,42 @@ environment/phone-or-SIP IDs and ElevenLabs credentials; a fetch-back receipt; an inbound staging call; actual hosted Agent Test IDs and raw results; a credentialed post-call latency/loss receipt; and one exercised rollback. No production or Vercel deployment was attempted. + +## 2026-07-11 — Session 11: W5 completion-gated durable delivery + +Replaced the legacy single queue with separate artifact and per-channel delivery +machines. Authenticated post-call events atomically release only the matching +conversation's outboxes; the worker also runs bounded ElevenLabs status +reconciliation when credentials exist. Two SQLite connections cannot claim the +same live lease. Lease attempts, expiry, timestamps, errors, stable provider +keys, provider IDs, callbacks, and worker heartbeats survive restart. + +The renderer starts only after effective hangup and persists both guarded +formats atomically. A verified artifact releases all selected channel jobs. A +crash before the provider-send boundary retries; a crash or connection loss +after that boundary moves to manual review, making duplicate risk explicit. +Poison/quarantined artifacts do not block later work. + +Added credential-gated Twilio SMS and SendGrid email adapters. SMS is link-only; +email attachments are the exact persisted guarded bytes. Both use recipient +allowlists and stable delivery keys, classify 429/5xx/hard-4xx/ambiguous +failures, and sanitize results. Twilio and SendGrid callback paths authenticate, +deduplicate, tolerate out-of-order states, and never retain callback contact +payloads. Production refuses `SimulatedDeliverer`; missing credentials fail +adapter construction. + +The long-running worker has bounded jittered backoff, SIGTERM drain, heartbeat, +and inspect/retry/quarantine/revoke-link commands. P16 was regenerated for quote +schema v4, including delivery callback events and worker heartbeats. The W5 +double-claim baseline defect is removed; only the W7 latency defect remains. + +**Verification:** 923 tests collected and full pytest exits zero; ruff is clean; +configured mypy reports no issues across 87 source files. Signed callback, +lease/crash, contact-change, provider classification, exact attachment, multi- +channel cardinality, hangup-release, reconciliation, migration, and production- +simulation refusal tests pass. + +**External W5 proof still required:** tenant P04 channel rule, sender identities, +Twilio/SendGrid credentials, staging recipient allowlists, one allowlisted email +and SMS, provider IDs, signed delivered/undelivered callbacks, and downloaded +artifact hash receipts. No production or Vercel deployment was attempted. diff --git a/pyproject.toml b/pyproject.toml index b790643..9ab2720 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -19,6 +19,7 @@ dependencies = [ # because the parity gate must run by default, not behind an extra. "reportlab>=4.0", "pypdf>=4.0", + "cryptography>=42.0", # SendGrid Event Webhook ECDSA verification ] [project.optional-dependencies] diff --git a/scripts/m5_legacy_inventory.py b/scripts/m5_legacy_inventory.py index 05db8e1..7057020 100755 --- a/scripts/m5_legacy_inventory.py +++ b/scripts/m5_legacy_inventory.py @@ -146,7 +146,7 @@ def compute() -> dict: "row_counts": _row_counts(store._conn, quote_tables), "known_gaps": [ "pre-W3 compatibility tables remain until migration evidence is retained", - "delivery job leases and authenticated call completion land in W5", + "credentialed provider receipts remain external to the schema fixture", "legacy rows without assent or artifact proof remain explicitly held", ], }, diff --git a/scripts/quote_worker.py b/scripts/quote_worker.py new file mode 100644 index 0000000..1a836cc --- /dev/null +++ b/scripts/quote_worker.py @@ -0,0 +1,156 @@ +#!/usr/bin/env python3 +"""Durable quote worker and bounded operator controls.""" +from __future__ import annotations + +import argparse +import json +import os +import random +import signal +import socket +import sys +import time +from datetime import datetime, timezone +from pathlib import Path + +REPO = Path(__file__).resolve().parent.parent +sys.path.insert(0, str(REPO / 'src')) + +from quoting.store import SqliteQuoteStore # noqa: E402 +from quoting.worker import run_artifact_once, run_delivery_once # noqa: E402 + + +def _store() -> SqliteQuoteStore: + return SqliteQuoteStore(os.environ.get( + 'SKU_QUOTE_STORE_DB', str(REPO / 'state' / 'quotes.db'))) + + +def _parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser() + sub = parser.add_subparsers(dest='command', required=True) + run = sub.add_parser('run') + run.add_argument('--max-poll-seconds', type=float, default=5.0) + run.add_argument('--lease-seconds', type=float, default=30.0) + run.add_argument('--once', action='store_true') + sub.add_parser('inspect') + retry = sub.add_parser('retry') + retry.add_argument('kind', choices=('artifact', 'delivery')) + retry.add_argument('tenant_id') + retry.add_argument('quote_number') + retry.add_argument('revision', type=int) + retry.add_argument('--channel') + quarantine = sub.add_parser('quarantine') + quarantine.add_argument('tenant_id') + quarantine.add_argument('quote_number') + quarantine.add_argument('revision', type=int) + quarantine.add_argument('reason') + revoke = sub.add_parser('revoke-link') + revoke.add_argument('token_hash') + return parser + + +def _run(args: argparse.Namespace) -> int: + if not 0.05 <= args.max_poll_seconds <= 60 or args.lease_seconds <= 0: + raise SystemExit('poll and lease bounds are invalid') + from runtime.config import build_delivery_adapters, build_gateway + + gateway, sessions = build_gateway() + store = gateway.quote_store + adapters = build_delivery_adapters() + reconciler = None + elevenlabs = { + 'api_key': os.environ.get('ELEVENLABS_API_KEY', ''), + 'secret': os.environ.get('ELEVENLABS_WEBHOOK_SECRET', ''), + 'agent_id': os.environ.get('ELEVENLABS_AGENT_ID', ''), + 'branch_id': os.environ.get('ELEVENLABS_STAGING_BRANCH_ID', ''), + 'version_id': os.environ.get('ELEVENLABS_AGENT_VERSION_ID', ''), + 'environment': os.environ.get('ELEVENLABS_ENVIRONMENT', ''), + } + if all(elevenlabs.values()): + from runtime.post_call import ( + CallCompletionReconciler, + ElevenLabsConversationStatusSource, + PostCallPolicy, + ) + + def release(completion): + store.release_for_call_completion( + tenant_id=completion.tenant_id, + conversation_id=completion.conversation_id, + completion_event_id=completion.event_id, + effective_hangup_at=completion.effective_hangup_at, + received_at=completion.received_at) + + reconciler = CallCompletionReconciler( + gateway.conversation_store, + ElevenLabsConversationStatusSource(elevenlabs['api_key']), + PostCallPolicy( + secret=elevenlabs['secret'], agent_id=elevenlabs['agent_id'], + branch_id=elevenlabs['branch_id'], version_id=elevenlabs['version_id'], + environment=elevenlabs['environment']), + completion_hook=release) + owner = f'{socket.gethostname()}:{os.getpid()}' + draining = False + + def stop(_signum, _frame): + nonlocal draining + draining = True + + signal.signal(signal.SIGTERM, stop) + signal.signal(signal.SIGINT, stop) + idle = 0 + while True: + now = datetime.now(timezone.utc) + store.heartbeat(worker_id=owner, now=now.timestamp(), draining=draining) + if draining: + return 0 + if reconciler is not None: + for held in store.held_conversations(): + reconciler.reconcile( + tenant_id=str(held['tenant_id']), + conversation_id=str(held['session_id']), + held_since=float(held['held_since']), now=now.timestamp()) + did_work = run_artifact_once( + store=store, owner=owner, now_fn=lambda: now, + lease_seconds=args.lease_seconds) + did_work = run_delivery_once( + store=store, owner=owner, adapters=adapters, + account_lookup=sessions.customer_db.by_number, + link_factory=gateway.quote_link_factory, now_fn=lambda: now, + lease_seconds=args.lease_seconds) or did_work + if args.once: + return 0 + if did_work: + idle = 0 + continue + idle = min(idle + 1, 8) + ceiling = min(args.max_poll_seconds, 0.05 * (2 ** idle)) + time.sleep(random.uniform(0.05, max(0.05, ceiling))) + + +def main() -> int: + args = _parser().parse_args() + store = _store() + if args.command == 'run': + return _run(args) + if args.command == 'inspect': + print(json.dumps(store.job_snapshot(), sort_keys=True, indent=2)) + return 0 + now = time.time() + if args.command == 'retry': + changed = store.operator_retry( + kind=args.kind, tenant_id=args.tenant_id, + quote_number=args.quote_number, revision=args.revision, + channel=args.channel, now=now) + elif args.command == 'quarantine': + changed = store.operator_quarantine( + tenant_id=args.tenant_id, quote_number=args.quote_number, + revision=args.revision, detail=args.reason, now=now) + else: + changed = store.revoke_capability(args.token_hash, at=datetime.now(timezone.utc)) + print(json.dumps({'changed': changed})) + return 0 if changed else 2 + + +if __name__ == '__main__': + raise SystemExit(main()) diff --git a/src/gateway/conversation_store.py b/src/gateway/conversation_store.py index bd88c55..10345a3 100644 --- a/src/gateway/conversation_store.py +++ b/src/gateway/conversation_store.py @@ -19,7 +19,7 @@ 'quote_number_allocations', 'quote_issue_keys', 'issuance_events', 'artifact_jobs', 'delivery_jobs', 'quote_artifacts', 'quote_capabilities', 'artifact_guard_events', 'quote_schema', - 'capability_access_events', + 'capability_access_events', 'provider_callback_events', 'worker_heartbeats', }) diff --git a/src/quoting/__init__.py b/src/quoting/__init__.py index 99adea6..7d669e8 100644 --- a/src/quoting/__init__.py +++ b/src/quoting/__init__.py @@ -8,9 +8,13 @@ from __future__ import annotations from quoting.delivery import ( + ChannelAdapter, Deliverer, DeliveryPayload, + ProviderResult, + SendGridEmailAdapter, SimulatedDeliverer, + TwilioSmsAdapter, ) from quoting.document_guard import ( DocumentParityViolation, @@ -29,7 +33,14 @@ ) from quoting.render_pdf import render_pdf from quoting.render_xlsx import render_xlsx -from quoting.store import QueueItem, QuoteStore, SqliteQuoteStore, order_content_hash +from quoting.store import ( + ArtifactJob, + DeliveryJob, + QueueItem, + QuoteStore, + SqliteQuoteStore, + order_content_hash, +) from quoting.totals import build_totals, extended_price, subtotal __all__ = [ @@ -40,4 +51,6 @@ 'DocumentParityViolation', 'assert_document_parity', 'document_facts', 'extract_pdf_facts', 'extract_xlsx_facts', 'verify_rendered_document', 'Deliverer', 'DeliveryPayload', 'SimulatedDeliverer', 'QueueItem', + 'ChannelAdapter', 'ProviderResult', 'TwilioSmsAdapter', + 'SendGridEmailAdapter', 'ArtifactJob', 'DeliveryJob', ] diff --git a/src/quoting/callbacks.py b/src/quoting/callbacks.py new file mode 100644 index 0000000..653cad5 --- /dev/null +++ b/src/quoting/callbacks.py @@ -0,0 +1,105 @@ +"""Provider callback authentication and normalized delivery outcomes.""" +from __future__ import annotations + +import base64 +import hashlib +import hmac +import json +from dataclasses import dataclass +from typing import Mapping + + +class CallbackError(ValueError): + pass + + +@dataclass(frozen=True) +class DeliveryCallback: + provider: str + event_id: str + delivery_key: str + provider_id: str | None + outcome: str + payload_digest: str + + +def verify_twilio_callback(*, auth_token: str, url: str, + params: Mapping[str, str], signature: str) -> None: + from runtime.twilio_sig import expected_signature + + if not auth_token or not signature or not hmac.compare_digest( + expected_signature(auth_token, url, params), signature): + raise CallbackError('invalid Twilio callback signature') + + +def parse_twilio_callback(params: Mapping[str, str]) -> DeliveryCallback: + status = params.get('MessageStatus', '').lower() + outcomes = { + 'queued': 'provider_accepted', 'accepted': 'provider_accepted', + 'sent': 'provider_accepted', 'delivered': 'delivered', + 'undelivered': 'undelivered', 'failed': 'failed_permanent', + } + outcome = outcomes.get(status) + provider_id = params.get('MessageSid', '') + delivery_key = params.get('DeliveryKey', '') + event_id = params.get('EventSid') or f'{provider_id}:{status}' + if not outcome or not provider_id or not delivery_key: + raise CallbackError('invalid Twilio delivery callback') + canonical = json.dumps(dict(params), sort_keys=True, separators=(',', ':')) + return DeliveryCallback( + 'twilio', event_id, delivery_key, provider_id, outcome, + hashlib.sha256(canonical.encode()).hexdigest()) + + +def verify_sendgrid_callback(*, public_key_pem: str, timestamp: str, + signature: str, raw_body: bytes) -> None: + """Verify SendGrid Event Webhook ECDSA over timestamp + raw body.""" + if not public_key_pem or not timestamp or not signature: + raise CallbackError('missing SendGrid callback signature') + try: + from cryptography.exceptions import InvalidSignature + from cryptography.hazmat.primitives import hashes, serialization + from cryptography.hazmat.primitives.asymmetric import ec + + key = serialization.load_pem_public_key(public_key_pem.encode()) + if not isinstance(key, ec.EllipticCurvePublicKey): + raise CallbackError('SendGrid callback key is not ECDSA') + key.verify( + base64.b64decode(signature, validate=True), + timestamp.encode() + raw_body, ec.ECDSA(hashes.SHA256())) + except CallbackError: + raise + except (InvalidSignature, ValueError, TypeError) as exc: + raise CallbackError('invalid SendGrid callback signature') from exc + except ImportError as exc: # production startup validation owns this dependency + raise RuntimeError('cryptography is required for SendGrid callbacks') from exc + + +def parse_sendgrid_callbacks(raw_body: bytes) -> tuple[DeliveryCallback, ...]: + try: + payload = json.loads(raw_body) + except (json.JSONDecodeError, UnicodeDecodeError) as exc: + raise CallbackError('invalid SendGrid callback JSON') from exc + if not isinstance(payload, list): + raise CallbackError('SendGrid callback must be an event list') + digest = hashlib.sha256(raw_body).hexdigest() + results = [] + outcomes = { + 'processed': 'provider_accepted', 'delivered': 'delivered', + 'bounce': 'bounced', 'dropped': 'failed_permanent', + 'deferred': 'undelivered', + } + for raw in payload: + if not isinstance(raw, dict): + raise CallbackError('invalid SendGrid callback event') + event = str(raw.get('event', '')).lower() + outcome = outcomes.get(event) + delivery_key = str(raw.get('delivery_key', '')) + event_id = str(raw.get('sg_event_id', '')) + provider_id = raw.get('sg_message_id') + if not outcome or not delivery_key or not event_id: + raise CallbackError('incomplete SendGrid callback event') + results.append(DeliveryCallback( + 'sendgrid', event_id, delivery_key, + str(provider_id) if provider_id else None, outcome, digest)) + return tuple(results) diff --git a/src/quoting/delivery.py b/src/quoting/delivery.py index fb3d65d..361105d 100644 --- a/src/quoting/delivery.py +++ b/src/quoting/delivery.py @@ -1,8 +1,13 @@ -"""M5d offline delivery seam and HMAC-signed quote links.""" +"""Delivery seams and credential-gated Twilio/SendGrid adapters.""" from __future__ import annotations +import base64 +import json +import urllib.error +import urllib.parse +import urllib.request from dataclasses import dataclass, field -from typing import Protocol +from typing import Mapping, Protocol class AccountContact(Protocol): @@ -23,6 +28,41 @@ class DeliveryPayload: signed_link: str +@dataclass(frozen=True) +class ProviderResult: + accepted: bool + provider_id: str | None + retryable: bool + ambiguous: bool + error_class: str = '' + + +@dataclass(frozen=True) +class HttpResult: + status: int + headers: Mapping[str, str] + body: bytes = b'' + + +class HttpTransport(Protocol): + def __call__(self, *, url: str, headers: Mapping[str, str], body: bytes, + timeout: float) -> HttpResult: ... + + +class RequestNotSent(RuntimeError): + """Connection failed with evidence that the provider received no request.""" + + +class AcceptanceUnknown(RuntimeError): + """The connection failed after send; retry could duplicate notification.""" + + +class ChannelAdapter(Protocol): + channel: str + + def send(self, payload: DeliveryPayload, *, delivery_key: str) -> ProviderResult: ... + + class Deliverer(Protocol): def deliver(self, *, account: AccountContact, quote_number: str, revision: int, summary: str, pdf_bytes: bytes, @@ -45,3 +85,152 @@ def deliver(self, *, account: AccountContact, quote_number: str, quote_number=quote_number, revision=revision, summary=summary, pdf_bytes=pdf_bytes, xlsx_bytes=xlsx_bytes, signed_link=signed_link)) + + +def _urllib_transport(*, url: str, headers: Mapping[str, str], body: bytes, + timeout: float) -> HttpResult: + request = urllib.request.Request( + url, data=body, headers=dict(headers), method='POST') + try: + with urllib.request.urlopen(request, timeout=timeout) as response: + return HttpResult( + status=int(response.status), headers=dict(response.headers), + body=response.read(64 * 1024)) + except urllib.error.HTTPError as exc: + return HttpResult( + status=int(exc.code), headers=dict(exc.headers), + body=exc.read(64 * 1024)) + except (TimeoutError, urllib.error.URLError) as exc: + # urllib cannot prove whether a socket failure happened before or after + # provider acceptance. Fail into manual review, never an automatic retry. + raise AcceptanceUnknown(type(exc).__name__) from exc + + +def _classified(result: HttpResult, *, provider_id: str | None) -> ProviderResult: + if 200 <= result.status < 300: + return ProviderResult(True, provider_id, False, False) + if result.status == 429: + return ProviderResult(False, None, True, False, 'rate_limited') + if result.status >= 500: + return ProviderResult(False, None, True, False, 'provider_5xx') + return ProviderResult(False, None, False, False, 'provider_4xx') + + +@dataclass(frozen=True) +class TwilioSmsAdapter: + account_sid: str + auth_token: str + from_number: str + allowed_recipients: frozenset[str] + status_callback_url: str = '' + transport: HttpTransport = _urllib_transport + api_base: str = 'https://api.twilio.com' + timeout_seconds: float = 8.0 + channel: str = 'sms' + + def __post_init__(self) -> None: + if not all((self.account_sid, self.auth_token, self.from_number)): + raise ValueError('Twilio credentials and sender are required') + if not self.allowed_recipients: + raise ValueError('Twilio staging recipient allowlist is required') + + def send(self, payload: DeliveryPayload, *, delivery_key: str) -> ProviderResult: + if payload.destination not in self.allowed_recipients: + return ProviderResult(False, None, False, False, + 'recipient_not_allowlisted') + form = { + 'To': payload.destination, + 'From': self.from_number, + 'Body': (f'Quote {payload.quote_number} revision {payload.revision}: ' + f'{payload.signed_link}'), + } + if self.status_callback_url: + separator = '&' if '?' in self.status_callback_url else '?' + form['StatusCallback'] = ( + f'{self.status_callback_url}{separator}' + f'DeliveryKey={urllib.parse.quote(delivery_key, safe="")}') + body = urllib.parse.urlencode(form).encode() + auth = base64.b64encode( + f'{self.account_sid}:{self.auth_token}'.encode()).decode() + try: + result = self.transport( + url=(f'{self.api_base}/2010-04-01/Accounts/' + f'{self.account_sid}/Messages.json'), + headers={ + 'Authorization': f'Basic {auth}', + 'Content-Type': 'application/x-www-form-urlencoded', + 'Idempotency-Key': delivery_key, + }, body=body, timeout=self.timeout_seconds) + except RequestNotSent: + return ProviderResult(False, None, True, False, 'request_not_sent') + except AcceptanceUnknown: + return ProviderResult(False, None, False, True, + 'acceptance_unknown') + provider_id = None + if 200 <= result.status < 300: + try: + raw = json.loads(result.body or b'{}') + provider_id = raw.get('sid') if isinstance(raw, dict) else None + except (json.JSONDecodeError, UnicodeDecodeError): + provider_id = None + return _classified(result, provider_id=provider_id) + + +@dataclass(frozen=True) +class SendGridEmailAdapter: + api_key: str + from_email: str + allowed_recipients: frozenset[str] + transport: HttpTransport = _urllib_transport + api_base: str = 'https://api.sendgrid.com' + timeout_seconds: float = 8.0 + sandbox_mode: bool = True + channel: str = 'email' + + def __post_init__(self) -> None: + if not self.api_key or not self.from_email: + raise ValueError('SendGrid credentials and sender are required') + if not self.allowed_recipients: + raise ValueError('SendGrid staging recipient allowlist is required') + + def send(self, payload: DeliveryPayload, *, delivery_key: str) -> ProviderResult: + if payload.destination not in self.allowed_recipients: + return ProviderResult(False, None, False, False, + 'recipient_not_allowlisted') + message = { + 'personalizations': [{ + 'to': [{'email': payload.destination}], + 'custom_args': {'delivery_key': delivery_key}, + }], + 'from': {'email': self.from_email}, + 'subject': f'Quote {payload.quote_number} revision {payload.revision}', + 'content': [{'type': 'text/plain', 'value': ( + f'{payload.summary}\n\nSecure link: {payload.signed_link}')}], + 'attachments': [ + {'content': base64.b64encode(payload.pdf_bytes).decode(), + 'type': 'application/pdf', 'filename': 'quote.pdf', + 'disposition': 'attachment'}, + {'content': base64.b64encode(payload.xlsx_bytes).decode(), + 'type': ('application/vnd.openxmlformats-officedocument.' + 'spreadsheetml.sheet'), 'filename': 'quote.xlsx', + 'disposition': 'attachment'}, + ], + 'mail_settings': {'sandbox_mode': {'enable': self.sandbox_mode}}, + } + try: + result = self.transport( + url=f'{self.api_base}/v3/mail/send', + headers={ + 'Authorization': f'Bearer {self.api_key}', + 'Content-Type': 'application/json', + 'Idempotency-Key': delivery_key, + }, body=json.dumps(message, separators=(',', ':')).encode(), + timeout=self.timeout_seconds) + except RequestNotSent: + return ProviderResult(False, None, True, False, 'request_not_sent') + except AcceptanceUnknown: + return ProviderResult(False, None, False, True, + 'acceptance_unknown') + provider_id = (result.headers.get('X-Message-Id') + or result.headers.get('x-message-id')) + return _classified(result, provider_id=provider_id) diff --git a/src/quoting/document_guard.py b/src/quoting/document_guard.py index 4f7257d..f95c4ff 100644 --- a/src/quoting/document_guard.py +++ b/src/quoting/document_guard.py @@ -256,13 +256,12 @@ def _scan_forbidden(document: QuoteDocument, pdf_bytes: bytes, if cell.value is not None] haystack = '\n'.join((pdf_text, *xlsx_values)).casefold() forbidden = list(_FORBIDDEN_NAMES) - forbidden.extend(( - document.presentation.branding_ref, - document.delivery_preflight.contact_id, - document.delivery_preflight.contact_version, - document.delivery_preflight.destination_fingerprint, - document.delivery_preflight.masked_destination, - )) + forbidden.append(document.presentation.branding_ref) + for preflight in document.delivery_preflights: + forbidden.extend(( + preflight.contact_id, preflight.contact_version, + preflight.destination_fingerprint, preflight.masked_destination, + )) for line in document.lines: forbidden.extend((line.resolution_source, line.resolution_confidence)) for value in forbidden: diff --git a/src/quoting/models.py b/src/quoting/models.py index dbbfbdf..277ec7b 100644 --- a/src/quoting/models.py +++ b/src/quoting/models.py @@ -125,10 +125,15 @@ class QuoteDocument: presentation: QuotePresentation provenance: QuoteProvenance delivery_preflight: DeliveryPreflight + additional_delivery_preflights: tuple[DeliveryPreflight, ...] = () def __post_init__(self) -> None: _validate_document(self) + @property + def delivery_preflights(self) -> tuple[DeliveryPreflight, ...]: + return (self.delivery_preflight, *self.additional_delivery_preflights) + _CENT = Decimal('0.01') _UOMS = frozenset({'each'}) @@ -153,12 +158,15 @@ def _money(value: Decimal, field: str) -> None: def _validate_document(document: QuoteDocument) -> None: header = document.header provenance = document.provenance - preflight = document.delivery_preflight - if (not preflight.available or preflight.channel not in ('email', 'sms') - or not all((preflight.contact_id, preflight.contact_version, - preflight.destination_fingerprint, - preflight.masked_destination))): - raise ValueError('a verified delivery preflight is required') + preflights = document.delivery_preflights + if len({item.channel for item in preflights}) != len(preflights): + raise ValueError('delivery preflight channels must be unique') + for preflight in preflights: + if (not preflight.available or preflight.channel not in ('email', 'sms') + or not all((preflight.contact_id, preflight.contact_version, + preflight.destination_fingerprint, + preflight.masked_destination))): + raise ValueError('a verified delivery preflight is required') if not all((provenance.gateway_build_sha, provenance.renderer_schema_version, provenance.guard_schema_version, diff --git a/src/quoting/store.py b/src/quoting/store.py index e7857f0..b3b21e1 100644 --- a/src/quoting/store.py +++ b/src/quoting/store.py @@ -46,11 +46,11 @@ 'quote_schema', 'quote_number_allocations', 'quote_issue_keys', 'issuance_events', 'artifact_jobs', 'delivery_jobs', 'quote_artifacts', 'artifact_guard_events', 'quote_capabilities', - 'capability_access_events', + 'capability_access_events', 'provider_callback_events', 'worker_heartbeats', }) _SHARED_CONVERSATION_TABLES = frozenset({ 'conversation_schema', 'conversations', 'processed_turns', - 'conversation_events', + 'conversation_events', 'call_completion_events', 'call_completion_alerts', }) @@ -82,6 +82,40 @@ class ArtifactRecord: verified_at: datetime +@dataclass(frozen=True) +class ArtifactJob: + tenant_id: str + quote_number: str + revision: int + status: str + attempts: int + lease_owner: str | None + lease_expires_at: float | None + completion_event_id: int | None + effective_hangup_at: float | None + render_started_at: float | None + detail: str + + +@dataclass(frozen=True) +class DeliveryJob: + tenant_id: str + quote_number: str + revision: int + channel: str + status: str + attempts: int + lease_owner: str | None + lease_expires_at: float | None + completion_event_id: int | None + effective_hangup_at: float | None + provider_started_at: float | None + provider_id: str | None + provider_key: str + error_class: str + detail: str + + @dataclass(frozen=True) class CapabilityRecord: token_hash: str @@ -111,6 +145,19 @@ def order_version_digest( return sha256(payload.encode()).hexdigest() +def _delivery_channels(document: QuoteDocument, + primary: str) -> tuple[str, ...]: + channels = tuple(item.channel for item in document.delivery_preflights) + if primary not in channels: + raise ValueError('requested delivery channel lacks frozen preflight') + return channels + + +def _provider_key(tenant_id: str, quote_number: str, revision: int, + channel: str) -> str: + return f'{tenant_id}:{quote_number}:r{revision}:{channel}' + + def _canonical_assent_receipt(receipt: dict, *, session_id: str, order_digest: str) -> str: required = { @@ -170,7 +217,7 @@ def migrate_quote_database(path: str | Path, *, 'before_sha256': sha256(original).hexdigest(), 'after_sha256': sha256(upgraded).hexdigest(), 'counts': counts, - 'schema_version': 3, + 'schema_version': 4, } except Exception: if store is not None: @@ -236,6 +283,17 @@ def artifact(self, tenant_id: str, quote_number: str, revision: int, def record_artifact_failure(self, document: QuoteDocument, *, reason: str, detail: str) -> None: ... + def release_for_call_completion( + self, *, tenant_id: str, conversation_id: str, + completion_event_id: int, effective_hangup_at: float, + received_at: float) -> int: ... + + def claim_artifact(self, *, owner: str, now: float, + lease_seconds: float = 30.0) -> ArtifactJob | None: ... + + def claim_delivery(self, *, owner: str, now: float, + lease_seconds: float = 30.0) -> DeliveryJob | None: ... + def claim_pending(self) -> QueueItem | None: ... def mark_delivered(self, quote_number: str, revision: int) -> None: ... @@ -325,7 +383,8 @@ def __init__(self, source) -> None: self._conn.execute( 'CREATE TABLE IF NOT EXISTS artifact_jobs ' '(tenant_id TEXT NOT NULL, quote_number TEXT NOT NULL, revision INTEGER NOT NULL, ' - 'status TEXT NOT NULL DEFAULT "pending", detail TEXT NOT NULL DEFAULT "", ' + 'status TEXT NOT NULL DEFAULT "waiting_for_call_end", ' + 'detail TEXT NOT NULL DEFAULT "", ' 'PRIMARY KEY(tenant_id,quote_number,revision), ' 'FOREIGN KEY(tenant_id,quote_number,revision) REFERENCES ' 'quotes(tenant_id,quote_number,revision))') @@ -337,6 +396,31 @@ def __init__(self, source) -> None: 'PRIMARY KEY(tenant_id,quote_number,revision,channel), ' 'FOREIGN KEY(tenant_id,quote_number,revision) REFERENCES ' 'quotes(tenant_id,quote_number,revision))') + artifact_columns = { + 'attempts': 'INTEGER NOT NULL DEFAULT 0', + 'lease_owner': 'TEXT', 'lease_expires_at': 'REAL', + 'next_attempt_at': 'REAL NOT NULL DEFAULT 0', + 'created_at': 'REAL NOT NULL DEFAULT 0', + 'updated_at': 'REAL NOT NULL DEFAULT 0', + 'completion_event_id': 'INTEGER', 'effective_hangup_at': 'REAL', + 'render_started_at': 'REAL', + } + delivery_columns = { + 'attempts': 'INTEGER NOT NULL DEFAULT 0', + 'lease_owner': 'TEXT', 'lease_expires_at': 'REAL', + 'next_attempt_at': 'REAL NOT NULL DEFAULT 0', + 'created_at': 'REAL NOT NULL DEFAULT 0', + 'updated_at': 'REAL NOT NULL DEFAULT 0', + 'effective_hangup_at': 'REAL', 'provider_started_at': 'REAL', + 'provider_id': 'TEXT', 'provider_key': 'TEXT NOT NULL DEFAULT ""', + 'error_class': 'TEXT NOT NULL DEFAULT ""', + } + self._ensure_columns('artifact_jobs', artifact_columns) + self._ensure_columns('delivery_jobs', delivery_columns) + self._conn.execute( + "UPDATE delivery_jobs SET provider_key=tenant_id || ':' || " + "quote_number || ':r' || revision || ':' || channel " + "WHERE provider_key='' ") self._conn.execute( 'CREATE TABLE IF NOT EXISTS quote_artifacts ' '(tenant_id TEXT NOT NULL, quote_number TEXT NOT NULL, ' @@ -377,13 +461,31 @@ def __init__(self, source) -> None: '(event_id INTEGER PRIMARY KEY AUTOINCREMENT, tenant_id TEXT NOT NULL, ' 'token_fingerprint TEXT NOT NULL, format TEXT NOT NULL, ' 'outcome TEXT NOT NULL, occurred_at TEXT NOT NULL)') + self._conn.execute( + 'CREATE TABLE IF NOT EXISTS provider_callback_events ' + '(provider TEXT NOT NULL,event_id TEXT NOT NULL,delivery_key TEXT NOT NULL,' + 'provider_id TEXT,outcome TEXT NOT NULL,payload_digest TEXT NOT NULL,' + 'received_at REAL NOT NULL,PRIMARY KEY(provider,event_id))') + self._conn.execute( + 'CREATE TABLE IF NOT EXISTS worker_heartbeats ' + '(worker_id TEXT PRIMARY KEY,started_at REAL NOT NULL,last_seen_at REAL NOT NULL,' + 'draining INTEGER NOT NULL DEFAULT 0)') if 'quotes' in pre_tables and 'quote_issue_keys' not in pre_tables: self._backfill_pre_w3() self._conn.execute( - 'INSERT INTO quote_schema(singleton,version) VALUES(1,3) ' + 'INSERT INTO quote_schema(singleton,version) VALUES(1,4) ' 'ON CONFLICT(singleton) DO UPDATE SET version=excluded.version') self._conn.commit() + def _ensure_columns(self, table: str, + declarations: dict[str, str]) -> None: + existing = {row['name'] for row in self._conn.execute( + f'PRAGMA table_info({table})')} + for name, declaration in declarations.items(): + if name not in existing: + self._conn.execute( + f'ALTER TABLE {table} ADD COLUMN {name} {declaration}') + def _backfill_pre_w3(self) -> None: """Preserve recognized quote rows without inventing missing proof.""" idempotency = { @@ -436,11 +538,14 @@ def _backfill_pre_w3(self) -> None: artifact_status, detail)) self._conn.execute( 'INSERT OR IGNORE INTO delivery_jobs ' - '(tenant_id,quote_number,revision,channel,status,detail) ' - 'VALUES(?,?,?,? ,"held",?)', + '(tenant_id,quote_number,revision,channel,status,detail,provider_key) ' + 'VALUES(?,?,?,? ,"held",?,?)', (h.tenant_id, h.quote_number, h.revision, document.delivery_preflight.channel, - 'legacy import requires artifact and call-completion review')) + 'legacy import requires artifact and call-completion review', + _provider_key( + h.tenant_id, h.quote_number, h.revision, + document.delivery_preflight.channel))) if h.quote_number not in allocated: self._conn.execute( 'INSERT INTO quote_number_allocations(tenant_id) VALUES(?)', @@ -528,14 +633,22 @@ def executed() -> None: order_digest, assent_json)) executed() self._conn.execute( - 'INSERT INTO artifact_jobs(tenant_id,quote_number,revision) ' - 'VALUES(?,?,?)', (tenant_id, quote_number, revision)) - executed() - self._conn.execute( - 'INSERT INTO delivery_jobs ' - '(tenant_id,quote_number,revision,channel) VALUES(?,?,?,?)', - (tenant_id, quote_number, revision, delivery_channel)) + 'INSERT INTO artifact_jobs(tenant_id,quote_number,revision,' + 'created_at,updated_at) VALUES(?,?,?,?,?)', + (tenant_id, quote_number, revision, + document.header.created_at.timestamp(), + document.header.created_at.timestamp())) executed() + for channel in _delivery_channels(document, delivery_channel): + self._conn.execute( + 'INSERT INTO delivery_jobs ' + '(tenant_id,quote_number,revision,channel,created_at,updated_at,' + 'provider_key) VALUES(?,?,?,?,?,?,?)', + (tenant_id, quote_number, revision, channel, + document.header.created_at.timestamp(), + document.header.created_at.timestamp(), + _provider_key(tenant_id, quote_number, revision, channel))) + executed() self._conn.execute( 'INSERT INTO quote_issue_keys VALUES(?,?,?,?,?,?,?)', (*key, quote_number, revision)) @@ -668,14 +781,23 @@ def executed() -> None: order_digest, assent_json)) executed() self._conn.execute( - 'INSERT INTO artifact_jobs(tenant_id,quote_number,revision) ' - 'VALUES(?,?,?)', (tenant_id, quote_number, new_revision)) - executed() - self._conn.execute( - 'INSERT INTO delivery_jobs ' - '(tenant_id,quote_number,revision,channel) VALUES(?,?,?,?)', - (tenant_id, quote_number, new_revision, delivery_channel)) + 'INSERT INTO artifact_jobs(tenant_id,quote_number,revision,' + 'created_at,updated_at) VALUES(?,?,?,?,?)', + (tenant_id, quote_number, new_revision, + document.header.created_at.timestamp(), + document.header.created_at.timestamp())) executed() + for channel in _delivery_channels(document, delivery_channel): + self._conn.execute( + 'INSERT INTO delivery_jobs ' + '(tenant_id,quote_number,revision,channel,created_at,updated_at,' + 'provider_key) VALUES(?,?,?,?,?,?,?)', + (tenant_id, quote_number, new_revision, channel, + document.header.created_at.timestamp(), + document.header.created_at.timestamp(), + _provider_key( + tenant_id, quote_number, new_revision, channel))) + executed() changed = self._conn.execute( 'UPDATE quotes SET status="superseded" WHERE tenant_id=? AND ' 'quote_number=? AND revision=? AND status="active"', @@ -709,13 +831,19 @@ def _repair_atomic_outboxes(self, document: QuoteDocument, self._conn.execute('BEGIN IMMEDIATE') try: self._conn.execute( - 'INSERT INTO artifact_jobs(tenant_id,quote_number,revision) ' - 'VALUES(?,?,?) ON CONFLICT DO NOTHING', - (h.tenant_id, h.quote_number, h.revision)) - self._conn.execute( - 'INSERT INTO delivery_jobs(tenant_id,quote_number,revision,channel) ' - 'VALUES(?,?,?,?) ON CONFLICT DO NOTHING', - (h.tenant_id, h.quote_number, h.revision, channel)) + 'INSERT INTO artifact_jobs(tenant_id,quote_number,revision,' + 'created_at,updated_at) VALUES(?,?,?,?,?) ON CONFLICT DO NOTHING', + (h.tenant_id, h.quote_number, h.revision, + h.created_at.timestamp(), h.created_at.timestamp())) + for selected in _delivery_channels(document, channel): + self._conn.execute( + 'INSERT INTO delivery_jobs(' + 'tenant_id,quote_number,revision,channel,created_at,updated_at,' + 'provider_key) VALUES(?,?,?,?,?,?,?) ON CONFLICT DO NOTHING', + (h.tenant_id, h.quote_number, h.revision, selected, + h.created_at.timestamp(), h.created_at.timestamp(), + _provider_key( + h.tenant_id, h.quote_number, h.revision, selected))) self._conn.execute( 'INSERT INTO quote_queue VALUES(?,?,"pending",0,"") ' 'ON CONFLICT DO NOTHING', (h.quote_number, h.revision)) @@ -772,7 +900,8 @@ def _persist_artifacts_atomic( 'AND revision=?', (h.tenant_id, h.quote_number, h.revision)).fetchone() job = self._conn.execute( - 'SELECT status FROM artifact_jobs WHERE tenant_id=? AND ' + 'SELECT status,effective_hangup_at,render_started_at FROM ' + 'artifact_jobs WHERE tenant_id=? AND ' 'quote_number=? AND revision=?', (h.tenant_id, h.quote_number, h.revision)).fetchone() if quote is None: @@ -783,6 +912,11 @@ def _persist_artifacts_atomic( 'INSERT INTO artifact_jobs(tenant_id,quote_number,revision) ' 'VALUES(?,?,?)', (h.tenant_id, h.quote_number, h.revision)) + elif (job['render_started_at'] is not None + and (job['effective_hangup_at'] is None + or float(job['render_started_at']) + < float(job['effective_hangup_at']))): + raise RuntimeError('render began before authenticated hangup') for record in records: self._conn.execute( 'INSERT INTO quote_artifacts VALUES(?,?,?,?,?,?,?,?,?,?,?)', @@ -795,9 +929,17 @@ def _persist_artifacts_atomic( raise RuntimeError( f'planted artifact fault after statement {step}') self._conn.execute( - 'UPDATE artifact_jobs SET status="verified",detail="" WHERE ' + 'UPDATE artifact_jobs SET status="verified",detail="",' + 'lease_owner=NULL,lease_expires_at=NULL,updated_at=? WHERE ' 'tenant_id=? AND quote_number=? AND revision=?', - (h.tenant_id, h.quote_number, h.revision)) + (verified_at.timestamp(), h.tenant_id, h.quote_number, + h.revision)) + self._conn.execute( + 'UPDATE delivery_jobs SET status="pending",updated_at=? WHERE ' + 'tenant_id=? AND quote_number=? AND revision=? AND status="held" ' + 'AND completion_event_id IS NOT NULL', + (verified_at.timestamp(), h.tenant_id, h.quote_number, + h.revision)) step += 1 if fault_after == step: raise RuntimeError( @@ -828,15 +970,18 @@ def artifact(self, tenant_id: str, quote_number: str, revision: int, def record_artifact_failure(self, document: QuoteDocument, *, reason: str, detail: str) -> None: h = document.header - now = datetime.now().astimezone().isoformat() + now_dt = datetime.now().astimezone() + now = now_dt.isoformat() self._conn.execute('BEGIN IMMEDIATE') try: self._conn.execute( 'INSERT INTO artifact_jobs ' - '(tenant_id,quote_number,revision,status,detail) ' - 'VALUES(?,?,?,"quarantined",?) ON CONFLICT DO UPDATE SET ' - 'status="quarantined",detail=excluded.detail', - (h.tenant_id, h.quote_number, h.revision, detail)) + '(tenant_id,quote_number,revision,status,detail,updated_at) ' + 'VALUES(?,?,?,"quarantined",?,?) ON CONFLICT DO UPDATE SET ' + 'status="quarantined",detail=excluded.detail,lease_owner=NULL,' + 'lease_expires_at=NULL,updated_at=excluded.updated_at', + (h.tenant_id, h.quote_number, h.revision, detail, + now_dt.timestamp())) self._conn.execute( 'INSERT INTO artifact_guard_events ' '(tenant_id,quote_number,revision,reason,detail,created_at) ' @@ -1072,6 +1217,260 @@ def repair_outbox(self, document: QuoteDocument) -> None: """Idempotently restore the required artifact queue row for a quote.""" self.enqueue(document.header.quote_number, document.header.revision) + def release_for_call_completion( + self, *, tenant_id: str, conversation_id: str, + completion_event_id: int, effective_hangup_at: float, + received_at: float) -> int: + """Release only outboxes bound to an authenticated durable completion.""" + if effective_hangup_at > received_at: + raise ValueError('effective hangup cannot be after receipt') + self._conn.execute('BEGIN IMMEDIATE') + try: + identity = self._conn.execute( + 'SELECT quote_number,revision FROM issuance_events ' + 'WHERE tenant_id=? AND session_id=?', + (tenant_id, conversation_id)).fetchall() + changed = 0 + for row in identity: + params = (completion_event_id, effective_hangup_at, received_at, + tenant_id, row['quote_number'], row['revision']) + artifact = self._conn.execute( + 'UPDATE artifact_jobs SET status="pending",' + 'completion_event_id=?,effective_hangup_at=?,updated_at=? ' + 'WHERE tenant_id=? AND quote_number=? AND revision=? AND ' + 'status="waiting_for_call_end"', params) + self._conn.execute( + 'UPDATE delivery_jobs SET completion_event_id=?, ' + 'effective_hangup_at=?,updated_at=? WHERE tenant_id=? AND ' + 'quote_number=? AND revision=? AND completion_event_id IS NULL', + params) + changed += artifact.rowcount + self._conn.execute('COMMIT') + return changed + except Exception: + if self._conn.in_transaction: + self._conn.execute('ROLLBACK') + raise + + def claim_artifact(self, *, owner: str, now: float, + lease_seconds: float = 30.0) -> ArtifactJob | None: + if not owner or lease_seconds <= 0: + raise ValueError('owner and positive lease_seconds are required') + self._conn.execute('BEGIN IMMEDIATE') + try: + self._conn.execute( + 'UPDATE artifact_jobs SET status="retry_wait",lease_owner=NULL,' + 'lease_expires_at=NULL,next_attempt_at=?,updated_at=?, ' + 'detail="lease_expired" WHERE status="leased" AND ' + 'lease_expires_at<=?', (now, now, now)) + row = self._conn.execute( + 'UPDATE artifact_jobs SET status="leased",attempts=attempts+1,' + 'lease_owner=?,lease_expires_at=?,render_started_at=?,updated_at=? ' + 'WHERE rowid=(SELECT rowid FROM artifact_jobs WHERE status IN ' + '("pending","retry_wait") AND next_attempt_at<=? AND ' + 'completion_event_id IS NOT NULL ORDER BY created_at,rowid LIMIT 1) ' + 'RETURNING *', + (owner, now + lease_seconds, now, now, now)).fetchone() + self._conn.execute('COMMIT') + return _artifact_job(row) if row is not None else None + except Exception: + if self._conn.in_transaction: + self._conn.execute('ROLLBACK') + raise + + def claim_delivery(self, *, owner: str, now: float, + lease_seconds: float = 30.0) -> DeliveryJob | None: + if not owner or lease_seconds <= 0: + raise ValueError('owner and positive lease_seconds are required') + self._conn.execute('BEGIN IMMEDIATE') + try: + self._conn.execute( + 'UPDATE delivery_jobs SET status="retry_wait",lease_owner=NULL,' + 'lease_expires_at=NULL,next_attempt_at=?,updated_at=?, ' + 'detail="lease_expired_before_provider_acceptance" ' + 'WHERE status="leased" AND lease_expires_at<=?', + (now, now, now)) + self._conn.execute( + 'UPDATE delivery_jobs SET status="manual_review",lease_owner=NULL,' + 'lease_expires_at=NULL,updated_at=?,detail="acceptance_unknown_after_crash",' + 'error_class="acceptance_unknown" WHERE status="sending" AND ' + 'lease_expires_at<=?', (now, now)) + row = self._conn.execute( + 'UPDATE delivery_jobs SET status="leased",attempts=attempts+1,' + 'lease_owner=?,lease_expires_at=?,updated_at=? ' + 'WHERE rowid=(SELECT d.rowid FROM delivery_jobs d JOIN ' + 'artifact_jobs a ON a.tenant_id=d.tenant_id AND ' + 'a.quote_number=d.quote_number AND a.revision=d.revision ' + 'WHERE d.status IN ("pending","retry_wait") AND ' + 'd.next_attempt_at<=? AND d.completion_event_id IS NOT NULL AND ' + 'a.status="verified" ORDER BY d.created_at,d.rowid LIMIT 1) ' + 'RETURNING *', + (owner, now + lease_seconds, now, now)).fetchone() + self._conn.execute('COMMIT') + return _delivery_job(row) if row is not None else None + except Exception: + if self._conn.in_transaction: + self._conn.execute('ROLLBACK') + raise + + def fail_artifact(self, job: ArtifactJob, *, status: str, detail: str, + now: float, next_attempt_at: float = 0) -> None: + if status not in ('retry_wait', 'quarantined', 'failed_permanent'): + raise ValueError('invalid artifact failure state') + self._conn.execute( + 'UPDATE artifact_jobs SET status=?,detail=?,next_attempt_at=?,' + 'lease_owner=NULL,lease_expires_at=NULL,updated_at=? WHERE ' + 'tenant_id=? AND quote_number=? AND revision=? AND status="leased" ' + 'AND lease_owner=?', + (status, detail, next_attempt_at, now, job.tenant_id, + job.quote_number, job.revision, job.lease_owner)) + self._conn.commit() + + def finish_delivery(self, job: DeliveryJob, *, status: str, + provider_id: str | None, error_class: str = '', + detail: str = '', now: float, + next_attempt_at: float = 0) -> None: + allowed = { + 'provider_accepted', 'delivered', 'bounced', 'undelivered', + 'retry_wait', 'needs_reauthorization', 'manual_review', + 'failed_permanent', + } + if status not in allowed: + raise ValueError('invalid delivery state') + self._conn.execute( + 'UPDATE delivery_jobs SET status=?,provider_id=?,error_class=?,' + 'detail=?,next_attempt_at=?,lease_owner=NULL,lease_expires_at=NULL,' + 'updated_at=? WHERE tenant_id=? AND quote_number=? AND revision=? ' + 'AND channel=? AND status IN ("leased","sending") AND lease_owner=?', + (status, provider_id, error_class, detail, next_attempt_at, now, + job.tenant_id, job.quote_number, job.revision, job.channel, + job.lease_owner)) + self._conn.commit() + + def mark_delivery_sending(self, job: DeliveryJob, *, now: float) -> None: + changed = self._conn.execute( + 'UPDATE delivery_jobs SET status="sending",provider_started_at=?,' + 'updated_at=? WHERE tenant_id=? AND quote_number=? AND revision=? ' + 'AND channel=? AND status="leased" AND lease_owner=? AND ' + 'effective_hangup_at IS NOT NULL AND effective_hangup_at<=?', + (now, now, job.tenant_id, job.quote_number, job.revision, + job.channel, job.lease_owner, now)) + self._conn.commit() + if changed.rowcount != 1: + raise RuntimeError('delivery lease lost before provider send') + + def delivery_job(self, tenant_id: str, quote_number: str, revision: int, + channel: str) -> DeliveryJob | None: + row = self._conn.execute( + 'SELECT * FROM delivery_jobs WHERE tenant_id=? AND quote_number=? ' + 'AND revision=? AND channel=?', + (tenant_id, quote_number, revision, channel)).fetchone() + return _delivery_job(row) if row is not None else None + + def apply_provider_callback( + self, *, provider: str, event_id: str, delivery_key: str, + provider_id: str | None, outcome: str, payload_digest: str, + received_at: float) -> bool: + """Idempotently apply a signed callback without regressing final state.""" + if outcome not in ('provider_accepted', 'delivered', 'bounced', + 'undelivered', 'failed_permanent'): + raise ValueError('invalid provider callback outcome') + self._conn.execute('BEGIN IMMEDIATE') + try: + inserted = self._conn.execute( + 'INSERT OR IGNORE INTO provider_callback_events VALUES(?,?,?,?,?,?,?)', + (provider, event_id, delivery_key, provider_id, outcome, + payload_digest, received_at)) + if inserted.rowcount == 0: + self._conn.execute('COMMIT') + return False + row = self._conn.execute( + 'SELECT status FROM delivery_jobs WHERE provider_key=?', + (delivery_key,)).fetchone() + if row is None: + raise KeyError('callback delivery identity not found') + final = {'delivered', 'bounced', 'undelivered', 'failed_permanent'} + if row['status'] not in final: + self._conn.execute( + 'UPDATE delivery_jobs SET status=?,provider_id=COALESCE(' + 'provider_id,?),updated_at=? WHERE provider_key=?', + (outcome, provider_id, received_at, delivery_key)) + self._conn.execute('COMMIT') + return True + except Exception: + if self._conn.in_transaction: + self._conn.execute('ROLLBACK') + raise + + def heartbeat(self, *, worker_id: str, now: float, + draining: bool = False) -> None: + self._conn.execute( + 'INSERT INTO worker_heartbeats VALUES(?,?,?,?) ON CONFLICT(worker_id) ' + 'DO UPDATE SET last_seen_at=excluded.last_seen_at,' + 'draining=excluded.draining', + (worker_id, now, now, int(draining))) + self._conn.commit() + + def held_conversations(self) -> tuple[dict[str, object], ...]: + rows = self._conn.execute( + 'SELECT DISTINCT a.tenant_id,i.session_id,MIN(a.created_at) AS held_since ' + 'FROM artifact_jobs a JOIN issuance_events i ON i.tenant_id=a.tenant_id ' + 'AND i.quote_number=a.quote_number AND i.revision=a.revision WHERE ' + 'a.status="waiting_for_call_end" GROUP BY a.tenant_id,i.session_id ' + 'ORDER BY held_since').fetchall() + return tuple(dict(row) for row in rows) + + def job_snapshot(self) -> dict[str, list[dict[str, object]]]: + artifact_fields = ( + 'tenant_id,quote_number,revision,status,attempts,lease_owner,' + 'lease_expires_at,next_attempt_at,completion_event_id,' + 'effective_hangup_at,render_started_at,detail') + delivery_fields = ( + 'tenant_id,quote_number,revision,channel,status,attempts,lease_owner,' + 'lease_expires_at,next_attempt_at,completion_event_id,' + 'effective_hangup_at,provider_started_at,provider_id,provider_key,' + 'error_class,detail') + return { + 'artifact_jobs': [dict(row) for row in self._conn.execute( + f'SELECT {artifact_fields} FROM artifact_jobs ORDER BY rowid')], + 'delivery_jobs': [dict(row) for row in self._conn.execute( + f'SELECT {delivery_fields} FROM delivery_jobs ORDER BY rowid')], + 'heartbeats': [dict(row) for row in self._conn.execute( + 'SELECT * FROM worker_heartbeats ORDER BY worker_id')], + } + + def operator_retry(self, *, kind: str, tenant_id: str, + quote_number: str, revision: int, + channel: str | None, now: float) -> bool: + if kind == 'artifact': + changed = self._conn.execute( + 'UPDATE artifact_jobs SET status="retry_wait",next_attempt_at=?,' + 'detail="operator_retry",updated_at=? WHERE tenant_id=? AND ' + 'quote_number=? AND revision=? AND status IN ' + '("quarantined","failed_permanent","retry_wait")', + (now, now, tenant_id, quote_number, revision)) + elif kind == 'delivery' and channel: + changed = self._conn.execute( + 'UPDATE delivery_jobs SET status="retry_wait",next_attempt_at=?,' + 'detail="operator_retry",updated_at=? WHERE tenant_id=? AND ' + 'quote_number=? AND revision=? AND channel=? AND status IN ' + '("failed_permanent","retry_wait")', + (now, now, tenant_id, quote_number, revision, channel)) + else: + raise ValueError('invalid operator retry target') + self._conn.commit() + return changed.rowcount == 1 + + def operator_quarantine(self, *, tenant_id: str, quote_number: str, + revision: int, detail: str, now: float) -> bool: + changed = self._conn.execute( + 'UPDATE artifact_jobs SET status="quarantined",detail=?,updated_at=? ' + 'WHERE tenant_id=? AND quote_number=? AND revision=? AND status NOT ' + 'IN ("verified","failed_permanent")', + (detail, now, tenant_id, quote_number, revision)) + self._conn.commit() + return changed.rowcount == 1 + def claim_pending(self) -> QueueItem | None: """Claim the oldest pending row for the single M5 worker.""" row = self._conn.execute( @@ -1115,6 +1514,40 @@ def queue_item(self, quote_number: str, revision: int) -> QueueItem | None: status=row['status'], attempts=row['attempts'], detail=row['detail']) +def _artifact_job(row: sqlite3.Row) -> ArtifactJob: + return ArtifactJob( + tenant_id=row['tenant_id'], quote_number=row['quote_number'], + revision=int(row['revision']), status=row['status'], + attempts=int(row['attempts']), lease_owner=row['lease_owner'], + lease_expires_at=(float(row['lease_expires_at']) + if row['lease_expires_at'] is not None else None), + completion_event_id=(int(row['completion_event_id']) + if row['completion_event_id'] is not None else None), + effective_hangup_at=(float(row['effective_hangup_at']) + if row['effective_hangup_at'] is not None else None), + render_started_at=(float(row['render_started_at']) + if row['render_started_at'] is not None else None), + detail=row['detail']) + + +def _delivery_job(row: sqlite3.Row) -> DeliveryJob: + return DeliveryJob( + tenant_id=row['tenant_id'], quote_number=row['quote_number'], + revision=int(row['revision']), channel=row['channel'], + status=row['status'], attempts=int(row['attempts']), + lease_owner=row['lease_owner'], + lease_expires_at=(float(row['lease_expires_at']) + if row['lease_expires_at'] is not None else None), + completion_event_id=(int(row['completion_event_id']) + if row['completion_event_id'] is not None else None), + effective_hangup_at=(float(row['effective_hangup_at']) + if row['effective_hangup_at'] is not None else None), + provider_started_at=(float(row['provider_started_at']) + if row['provider_started_at'] is not None else None), + provider_id=row['provider_id'], provider_key=row['provider_key'], + error_class=row['error_class'], detail=row['detail']) + + def _document_from_dict(d: dict) -> QuoteDocument: h = d['header'] header = QuoteHeader( @@ -1158,6 +1591,10 @@ def _document_from_dict(d: dict) -> QuoteDocument: 'inventory_source_versions': tuple( provenance_d['inventory_source_versions'])}) delivery_preflight = DeliveryPreflight(**d['delivery_preflight']) + additional_preflights = tuple( + DeliveryPreflight(**item) + for item in d.get('additional_delivery_preflights', ())) return QuoteDocument(header=header, lines=lines, totals=totals, terms=terms, presentation=presentation, provenance=provenance, - delivery_preflight=delivery_preflight) + delivery_preflight=delivery_preflight, + additional_delivery_preflights=additional_preflights) diff --git a/src/quoting/worker.py b/src/quoting/worker.py index b00b621..d68dc8a 100644 --- a/src/quoting/worker.py +++ b/src/quoting/worker.py @@ -4,112 +4,195 @@ import hashlib from collections.abc import Callable from datetime import datetime, timezone +from typing import Mapping -from gateway.journal import ConversationJournal, EventType -from quoting.delivery import AccountContact, Deliverer +from gateway.journal import ConversationJournal +from quoting.delivery import ( + AccountContact, + ChannelAdapter, + Deliverer, + DeliveryPayload, + ProviderResult, +) from quoting.document_guard import DocumentParityViolation, verify_rendered_document from quoting.models import QuoteDocument from quoting.render_pdf import render_pdf from quoting.render_xlsx import render_xlsx -from quoting.store import QuoteStore +from quoting.store import QuoteStore, SqliteQuoteStore LinkFactory = Callable[[QuoteDocument], str] AccountLookup = Callable[[str], AccountContact | None] -def _record_block(store: QuoteStore, journal: ConversationJournal, *, - quote_number: str, revision: int, reason: str, - detail: str, document: QuoteDocument | None = None, - escalation_hook: Callable[[str, dict], None] = lambda *_: None - ) -> None: - if document is not None: - store.record_artifact_failure(document, reason=reason, detail=detail) - store.mark_blocked(quote_number, revision, detail) - journal.record( - EventType.QUOTE_DELIVERY_BLOCKED, quote_number, - quote_number=quote_number, revision=revision, reason=reason, - detail=detail) - escalation_hook(reason, { - 'quote_number': quote_number, 'revision': revision, - 'detail': detail, - }) +def _destination(account: AccountContact, channel: str) -> str | None: + if channel == 'email': + value = getattr(account, 'email', None) + verified = bool(getattr(account, 'email_verified', False)) + return str(value).strip().lower() if value and verified else None + if channel == 'sms': + value = getattr(account, 'phone', None) + verified = bool(getattr(account, 'phone_verified', False)) + return str(value) if value and verified else None + return None -def run_once(*, store: QuoteStore, deliverer: Deliverer, - account_lookup: AccountLookup, journal: ConversationJournal, - link_factory: LinkFactory, - now_fn: Callable[[], datetime] = lambda: datetime.now(timezone.utc), - escalation_hook: Callable[[str, dict], None] = lambda *_: None) -> bool: - """Process at most one pending quote revision after parity verification.""" - item = store.claim_pending() - if item is None: +def _fingerprint(tenant_id: str, contact_id: str, destination: str) -> str: + return hashlib.sha256( + f'{tenant_id}:{contact_id}:{destination}'.encode()).hexdigest() + + +def run_artifact_once(*, store: SqliteQuoteStore, owner: str, + now_fn: Callable[[], datetime] = lambda: datetime.now( + timezone.utc), + lease_seconds: float = 30.0) -> bool: + """Render one completion-released artifact job and verify before release.""" + now = now_fn() + job = store.claim_artifact( + owner=owner, now=now.timestamp(), lease_seconds=lease_seconds) + if job is None: return False - document = store.get(item.quote_number, item.revision) + document = store.get(job.quote_number, job.revision) if document is None: - _record_block( - store, journal, quote_number=item.quote_number, - revision=item.revision, reason='quote_missing', - detail='quote_missing', escalation_hook=escalation_hook) + store.fail_artifact( + job, status='failed_permanent', detail='quote_missing', + now=now.timestamp()) + return True + try: + pdf_bytes = render_pdf(document) + xlsx_bytes = render_xlsx(document) + verify_rendered_document( + document, pdf_bytes=pdf_bytes, xlsx_bytes=xlsx_bytes) + store.persist_artifacts_atomic( + document, pdf_bytes=pdf_bytes, xlsx_bytes=xlsx_bytes, + verified_at=now) + except DocumentParityViolation as exc: + store.record_artifact_failure( + document, reason='document_parity', detail=str(exc)) + except Exception as exc: + store.fail_artifact( + job, status='retry_wait', detail=type(exc).__name__, + now=now.timestamp(), next_attempt_at=now.timestamp() + 5) + return True + + +def run_delivery_once( + *, store: SqliteQuoteStore, owner: str, + adapters: Mapping[str, ChannelAdapter], account_lookup: AccountLookup, + link_factory: LinkFactory, + now_fn: Callable[[], datetime] = lambda: datetime.now(timezone.utc), + lease_seconds: float = 30.0, + metadata_only_contact_correction_allowed: bool = False) -> bool: + """Send one verified artifact through its frozen, revalidated channel.""" + now = now_fn() + job = store.claim_delivery( + owner=owner, now=now.timestamp(), lease_seconds=lease_seconds) + if job is None: + return False + document = store.get(job.quote_number, job.revision) + if document is None: + store.finish_delivery( + job, status='failed_permanent', provider_id=None, + error_class='quote_missing', now=now.timestamp()) return True account = account_lookup(document.header.account_id) - if account is None or account.email is None: - _record_block( - store, journal, quote_number=item.quote_number, - revision=item.revision, reason='contact_missing', - detail='contact_missing', escalation_hook=escalation_hook) + preflight = next( + (item for item in document.delivery_preflights + if item.channel == job.channel), None) + if preflight is None: + store.finish_delivery( + job, status='failed_permanent', provider_id=None, + error_class='preflight_missing', now=now.timestamp()) return True - - h = document.header - persisted_pdf = store.artifact( - h.tenant_id, h.quote_number, h.revision, 'pdf') - persisted_xlsx = store.artifact( - h.tenant_id, h.quote_number, h.revision, 'xlsx') - if (persisted_pdf is None) != (persisted_xlsx is None): - _record_block( - store, journal, quote_number=item.quote_number, - revision=item.revision, reason='partial_artifact_set', - detail='partial_artifact_set', document=document, - escalation_hook=escalation_hook) + destination = _destination(account, job.channel) if account is not None else None + contact_id = str(getattr(account, 'contact_id', '')) if account else '' + version = str(getattr(account, 'contact_version', '')) if account else '' + observed_fingerprint = ( + _fingerprint(job.tenant_id, contact_id, destination) + if destination and contact_id else '') + changed_destination = observed_fingerprint != preflight.destination_fingerprint + changed_metadata = version != preflight.contact_version + if (not destination or changed_destination + or (changed_metadata and not metadata_only_contact_correction_allowed)): + store.finish_delivery( + job, status='needs_reauthorization', provider_id=None, + error_class='contact_changed', now=now.timestamp()) return True - if persisted_pdf is None: - pdf_bytes = render_pdf(document) - xlsx_bytes = render_xlsx(document) - try: - verify_rendered_document( - document, pdf_bytes=pdf_bytes, xlsx_bytes=xlsx_bytes) - except DocumentParityViolation as exc: - _record_block( - store, journal, quote_number=item.quote_number, - revision=item.revision, reason='document_parity', - detail=str(exc), document=document, - escalation_hook=escalation_hook) + adapter = adapters.get(job.channel) + if adapter is None: + store.finish_delivery( + job, status='failed_permanent', provider_id=None, + error_class='adapter_missing', now=now.timestamp()) + return True + pdf = store.artifact(job.tenant_id, job.quote_number, job.revision, 'pdf') + xlsx = store.artifact(job.tenant_id, job.quote_number, job.revision, 'xlsx') + if pdf is None or xlsx is None: + store.finish_delivery( + job, status='manual_review', provider_id=None, + error_class='verified_artifact_missing', now=now.timestamp()) + return True + for artifact in (pdf, xlsx): + if (artifact.size != len(artifact.content) + or artifact.sha256 != hashlib.sha256(artifact.content).hexdigest()): + store.finish_delivery( + job, status='manual_review', provider_id=None, + error_class='artifact_hash_mismatch', now=now.timestamp()) return True - persisted_pdf, persisted_xlsx = store.persist_artifacts_atomic( - document, pdf_bytes=pdf_bytes, xlsx_bytes=xlsx_bytes, - verified_at=now_fn()) - else: - assert persisted_xlsx is not None - for artifact in (persisted_pdf, persisted_xlsx): - if (artifact.size != len(artifact.content) - or artifact.sha256 != hashlib.sha256( - artifact.content).hexdigest()): - _record_block( - store, journal, quote_number=item.quote_number, - revision=item.revision, reason='artifact_hash_mismatch', - detail='artifact_hash_mismatch', document=document, - escalation_hook=escalation_hook) - return True - pdf_bytes = persisted_pdf.content - xlsx_bytes = persisted_xlsx.content - - deliverer.deliver( - account=account, quote_number=item.quote_number, - revision=item.revision, - summary=f'Quote {item.quote_number} revision {item.revision}', - pdf_bytes=pdf_bytes, xlsx_bytes=xlsx_bytes, + payload = DeliveryPayload( + account_id=document.header.account_id, destination=destination, + quote_number=job.quote_number, revision=job.revision, + summary=f'Quote {job.quote_number} revision {job.revision}', + pdf_bytes=pdf.content, xlsx_bytes=xlsx.content, signed_link=link_factory(document)) - store.mark_delivered(item.quote_number, item.revision) - journal.record( - EventType.QUOTE_DELIVERED, item.quote_number, - quote_number=item.quote_number, revision=item.revision) + store.mark_delivery_sending(job, now=now.timestamp()) + result = adapter.send(payload, delivery_key=job.provider_key) + if result.accepted: + status = 'provider_accepted' + elif result.ambiguous: + status = 'manual_review' + elif result.retryable: + status = 'retry_wait' + else: + status = 'failed_permanent' + store.finish_delivery( + job, status=status, provider_id=result.provider_id, + error_class=result.error_class, now=now.timestamp(), + next_attempt_at=(now.timestamp() + 5 if status == 'retry_wait' else 0)) return True + + +def run_once(*, store: QuoteStore, deliverer: Deliverer, + account_lookup: AccountLookup, journal: ConversationJournal, + link_factory: LinkFactory, + now_fn: Callable[[], datetime] = lambda: datetime.now(timezone.utc), + escalation_hook: Callable[[str, dict], None] = lambda *_: None) -> bool: + """Compatibility one-pass wrapper over the completion-gated W5 machines.""" + if isinstance(store, SqliteQuoteStore): + class LegacyAdapter: + def __init__(self, channel: str): + self.channel = channel + + def send(self, payload: DeliveryPayload, *, + delivery_key: str) -> ProviderResult: + account = account_lookup(payload.account_id) + if account is None: + return ProviderResult( + False, None, False, False, 'contact_missing') + deliverer.deliver( + account=account, quote_number=payload.quote_number, + revision=payload.revision, summary=payload.summary, + pdf_bytes=payload.pdf_bytes, xlsx_bytes=payload.xlsx_bytes, + signed_link=payload.signed_link) + return ProviderResult( + True, f'simulated:{delivery_key}', False, False) + + owner = 'legacy-single-pass' + did_work = run_artifact_once( + store=store, owner=owner, now_fn=now_fn) + did_work = run_delivery_once( + store=store, owner=owner, + adapters={name: LegacyAdapter(name) for name in ('email', 'sms')}, + account_lookup=account_lookup, link_factory=link_factory, + now_fn=now_fn) or did_work + return did_work + + raise TypeError('W5 worker requires the completion-gated SqliteQuoteStore') diff --git a/src/runtime/app.py b/src/runtime/app.py index 73b4d05..4e53249 100644 --- a/src/runtime/app.py +++ b/src/runtime/app.py @@ -317,12 +317,93 @@ async def elevenlabs_post_call(request: Request): version_id=qualified.version_id, environment=qualified.environment, payload_digest=qualified.payload_digest) + gateway.quote_store.release_for_call_completion( + tenant_id=tenant_id, + conversation_id=qualified.conversation_id, + completion_event_id=completion.event_id, + effective_hangup_at=completion.effective_hangup_at, + received_at=completion.received_at) except PostCallError: return JSONResponse({'error': 'invalid_post_call'}, status_code=401) except ConversationConflict: return JSONResponse({'error': 'post_call_conflict'}, status_code=409) return {'status': 'received', 'completion_event_id': completion.event_id} + @app.post('/webhooks/twilio/delivery') + async def twilio_delivery_status(request: Request): + import time + + from quoting.callbacks import ( + CallbackError, + parse_twilio_callback, + verify_twilio_callback, + ) + + token = os.environ.get('TWILIO_AUTH_TOKEN', '') + if not token: + return JSONResponse({'error': 'twilio_callback_not_configured'}, + status_code=503) + form = {key: str(value) for key, value in (await request.form()).items()} + try: + verify_twilio_callback( + auth_token=token, url=str(request.url), params=form, + signature=request.headers.get('X-Twilio-Signature', '')) + callback = parse_twilio_callback({ + **form, + 'DeliveryKey': request.query_params.get('DeliveryKey', ''), + }) + changed = gateway.quote_store.apply_provider_callback( + provider=callback.provider, event_id=callback.event_id, + delivery_key=callback.delivery_key, + provider_id=callback.provider_id, outcome=callback.outcome, + payload_digest=callback.payload_digest, + received_at=time.time()) + except CallbackError: + return JSONResponse({'error': 'invalid_twilio_callback'}, + status_code=401) + except KeyError: + return JSONResponse({'error': 'unknown_delivery'}, status_code=404) + return {'status': 'received', 'new_event': changed} + + @app.post('/webhooks/sendgrid/events') + async def sendgrid_delivery_events(request: Request): + import time + + from quoting.callbacks import ( + CallbackError, + parse_sendgrid_callbacks, + verify_sendgrid_callback, + ) + + public_key = os.environ.get('SENDGRID_WEBHOOK_PUBLIC_KEY', '') + if not public_key: + return JSONResponse({'error': 'sendgrid_callback_not_configured'}, + status_code=503) + raw = await request.body() + try: + verify_sendgrid_callback( + public_key_pem=public_key, + timestamp=request.headers.get( + 'X-Twilio-Email-Event-Webhook-Timestamp', ''), + signature=request.headers.get( + 'X-Twilio-Email-Event-Webhook-Signature', ''), + raw_body=raw) + callbacks = parse_sendgrid_callbacks(raw) + changed = 0 + for callback in callbacks: + changed += int(gateway.quote_store.apply_provider_callback( + provider=callback.provider, event_id=callback.event_id, + delivery_key=callback.delivery_key, + provider_id=callback.provider_id, outcome=callback.outcome, + payload_digest=callback.payload_digest, + received_at=time.time())) + except CallbackError: + return JSONResponse({'error': 'invalid_sendgrid_callback'}, + status_code=401) + except KeyError: + return JSONResponse({'error': 'unknown_delivery'}, status_code=404) + return {'status': 'received', 'new_events': changed} + # -- Twilio voice (TwiML request/response) ----------------------- @app.post('/voice') diff --git a/src/runtime/config.py b/src/runtime/config.py index 507a58b..a5cea36 100644 --- a/src/runtime/config.py +++ b/src/runtime/config.py @@ -219,11 +219,43 @@ def link_for(document): def build_quote_deliverer(): - """Offline default; live senders remain credential-gated follow-up work.""" + """Legacy offline seam; production may never construct the simulator.""" + if os.environ.get('SKU_ENVIRONMENT', 'development').lower() == 'production': + raise RuntimeError('SimulatedDeliverer is forbidden in production') from quoting.delivery import SimulatedDeliverer return SimulatedDeliverer() +def build_delivery_adapters(): + """Build selected live adapters eagerly so missing credentials fail startup.""" + from quoting.delivery import SendGridEmailAdapter, TwilioSmsAdapter + + channels = {item.strip() for item in os.environ.get( + 'SKU_DELIVERY_CHANNELS', 'email,sms').split(',') if item.strip()} + if not channels or channels - {'email', 'sms'}: + raise RuntimeError('SKU_DELIVERY_CHANNELS must select email and/or sms') + adapters = {} + if 'email' in channels: + allowlist = frozenset(item.strip().lower() for item in os.environ.get( + 'SKU_EMAIL_ALLOWLIST', '').split(',') if item.strip()) + adapters['email'] = SendGridEmailAdapter( + api_key=os.environ.get('SENDGRID_API_KEY', ''), + from_email=os.environ.get('SENDGRID_FROM_EMAIL', ''), + allowed_recipients=allowlist, + sandbox_mode=os.environ.get('SENDGRID_SANDBOX', '1') == '1') + if 'sms' in channels: + allowlist = frozenset(item.strip() for item in os.environ.get( + 'SKU_SMS_ALLOWLIST', '').split(',') if item.strip()) + adapters['sms'] = TwilioSmsAdapter( + account_sid=os.environ.get('TWILIO_ACCOUNT_SID', ''), + auth_token=os.environ.get('TWILIO_AUTH_TOKEN', ''), + from_number=os.environ.get('TWILIO_FROM_NUMBER', ''), + allowed_recipients=allowlist, + status_callback_url=os.environ.get( + 'TWILIO_STATUS_CALLBACK_URL', '')) + return adapters + + def build_streaming_asr(): """Select the streaming ASR for the Twilio Media Streams path. Real AssemblyAI v3 when ASSEMBLYAI_API_KEY is set; otherwise a no-op simulated diff --git a/src/runtime/post_call.py b/src/runtime/post_call.py index 58ea536..50fc720 100644 --- a/src/runtime/post_call.py +++ b/src/runtime/post_call.py @@ -166,17 +166,20 @@ class CallCompletionReconciler: def __init__(self, store: ConversationStore, source: ConversationStatusSource, policy: PostCallPolicy, *, - max_held_secs: float = 120.0): + max_held_secs: float = 120.0, + completion_hook=None): self.store = store self.source = source self.policy = policy self.max_held_secs = max_held_secs + self.completion_hook = completion_hook or (lambda _completion: None) def reconcile(self, *, tenant_id: str, conversation_id: str, held_since: float, now: float | None = None) -> CallCompletion | None: at = time.time() if now is None else now existing = self.store.call_completion(tenant_id, conversation_id) if existing is not None: + self.completion_hook(existing) return existing status = self.source.get(conversation_id) if status is not None and status.done: @@ -190,13 +193,15 @@ def reconcile(self, *, tenant_id: str, conversation_id: str, digest = hashlib.sha256( ('reconcile:' + conversation_id + ':' + ':'.join(bindings) + f':{status.effective_hangup_at}').encode()).hexdigest() - return self.store.record_call_completion( + completion = self.store.record_call_completion( tenant_id=tenant_id, conversation_id=conversation_id, source='elevenlabs_conversation_reconciliation', effective_hangup_at=status.effective_hangup_at, received_at=at, agent_id=status.agent_id, branch_id=status.branch_id, version_id=status.version_id, environment=status.environment, payload_digest=digest) + self.completion_hook(completion) + return completion if at - held_since >= self.max_held_secs: self.store.record_completion_alert( tenant_id=tenant_id, conversation_id=conversation_id, diff --git a/tests/gateway_fixtures.py b/tests/gateway_fixtures.py index b4e06d7..84850f8 100644 --- a/tests/gateway_fixtures.py +++ b/tests/gateway_fixtures.py @@ -98,6 +98,14 @@ def link_factory(document): token = capabilities.mint(document, 'pdf') return f'https://quotes.example.test/quote/{token}/pdf' + received_timestamp = gw.now_fn().timestamp() + for row in gw.quote_store._conn.execute( + 'SELECT DISTINCT session_id FROM issuance_events').fetchall(): + gw.quote_store.release_for_call_completion( + tenant_id=tenant_id, conversation_id=row['session_id'], + completion_event_id=1, effective_hangup_at=received_timestamp, + received_at=received_timestamp) + run_once(store=gw.quote_store, deliverer=deliverer, account_lookup=accounts.get, journal=journal, link_factory=link_factory) diff --git a/tests/test_delivery_adapters.py b/tests/test_delivery_adapters.py new file mode 100644 index 0000000..b084405 --- /dev/null +++ b/tests/test_delivery_adapters.py @@ -0,0 +1,96 @@ +from __future__ import annotations + +import base64 +import json +import urllib.parse + +import pytest + +from quoting.delivery import ( + AcceptanceUnknown, + DeliveryPayload, + HttpResult, + RequestNotSent, + SendGridEmailAdapter, + TwilioSmsAdapter, +) + + +def _payload(destination: str) -> DeliveryPayload: + return DeliveryPayload( + '1001', destination, 'Q-tenant_001-000001', 1, 'Quote summary', + b'%PDF-exact', b'PK-exact', 'https://quotes.example.test/opaque') + + +@pytest.mark.parametrize('status,retryable,error', ( + (429, True, 'rate_limited'), (503, True, 'provider_5xx'), + (400, False, 'provider_4xx'), +)) +def test_twilio_sms_is_link_only_and_classifies_http(status, retryable, error): + captured = {} + + def transport(**request): + captured.update(request) + return HttpResult(status, {}, b'{"sid":"SM1"}') + + adapter = TwilioSmsAdapter( + 'AC1', 'secret', '+15550109999', frozenset({'+15550100100'}), + status_callback_url='https://api.example.test/webhooks/twilio/delivery', + transport=transport) + result = adapter.send(_payload('+15550100100'), delivery_key='stable-key') + form = urllib.parse.parse_qs(captured['body'].decode()) + assert form['Body'] == [ + 'Quote Q-tenant_001-000001 revision 1: ' + 'https://quotes.example.test/opaque'] + assert b'%PDF' not in captured['body'] and b'PK-exact' not in captured['body'] + assert captured['headers']['Idempotency-Key'] == 'stable-key' + assert 'DeliveryKey=stable-key' in form['StatusCallback'][0] + assert result.retryable is retryable and result.error_class == error + assert 'secret' not in repr(result) + + +def test_sendgrid_exact_attachments_sandbox_and_provider_id(): + captured = {} + + def transport(**request): + captured.update(request) + return HttpResult(202, {'X-Message-Id': 'sg-1'}) + + adapter = SendGridEmailAdapter( + 'SG.secret', 'sender@example.test', frozenset({'quote@example.test'}), + transport=transport, sandbox_mode=True) + result = adapter.send(_payload('quote@example.test'), delivery_key='stable-key') + body = json.loads(captured['body']) + assert base64.b64decode(body['attachments'][0]['content']) == b'%PDF-exact' + assert base64.b64decode(body['attachments'][1]['content']) == b'PK-exact' + assert body['mail_settings']['sandbox_mode']['enable'] is True + assert body['personalizations'][0]['custom_args']['delivery_key'] == 'stable-key' + assert result.accepted and result.provider_id == 'sg-1' + + +@pytest.mark.parametrize('exception,retryable,ambiguous', ( + (RequestNotSent('before connect'), True, False), + (AcceptanceUnknown('after send'), False, True), +)) +def test_connection_outcomes_do_not_hide_duplicate_risk(exception, retryable, ambiguous): + def transport(**_request): + raise exception + + adapter = SendGridEmailAdapter( + 'SG.secret', 'sender@example.test', frozenset({'quote@example.test'}), + transport=transport) + result = adapter.send(_payload('quote@example.test'), delivery_key='key') + assert result.retryable is retryable + assert result.ambiguous is ambiguous + + +def test_allowlist_and_credentials_fail_before_provider_io(): + calls = [] + with pytest.raises(ValueError, match='credentials'): + TwilioSmsAdapter('', '', '', frozenset({'x'})) + adapter = SendGridEmailAdapter( + 'SG.secret', 'sender@example.test', frozenset({'allowed@example.test'}), + transport=lambda **request: calls.append(request)) + result = adapter.send(_payload('other@example.test'), delivery_key='key') + assert result.error_class == 'recipient_not_allowlisted' + assert calls == [] diff --git a/tests/test_delivery_callback_routes.py b/tests/test_delivery_callback_routes.py new file mode 100644 index 0000000..90696b9 --- /dev/null +++ b/tests/test_delivery_callback_routes.py @@ -0,0 +1,63 @@ +from __future__ import annotations + +import base64 +import json + +from cryptography.hazmat.primitives import hashes, serialization +from cryptography.hazmat.primitives.asymmetric import ec +from fastapi.testclient import TestClient +from gateway_fixtures import build_gateway +from test_m5_w5_delivery_worker import ( + NOW, + FakeAdapter, + _account, + _store, +) + +from quoting.delivery import ProviderResult +from quoting.worker import run_artifact_once, run_delivery_once +from runtime.app import create_app + + +def test_signed_sendgrid_route_is_replay_safe_and_terminal(tmp_path, monkeypatch): + store, document = _store(tmp_path / 'callbacks.sqlite') + store.release_for_call_completion( + tenant_id='tenant_001', conversation_id='conv-1', completion_event_id=1, + effective_hangup_at=NOW.timestamp() - 1, received_at=NOW.timestamp()) + run_artifact_once(store=store, owner='artifact', now_fn=lambda: NOW) + run_delivery_once( + store=store, owner='delivery', + adapters={'email': FakeAdapter( + ProviderResult(True, 'provider-1', False, False))}, + account_lookup=lambda _: _account(), + link_factory=lambda _: 'https://quotes.example.test/opaque', + now_fn=lambda: NOW) + job = store.delivery_job( + 'tenant_001', document.header.quote_number, 1, 'email') + + private = ec.generate_private_key(ec.SECP256R1()) + public_pem = private.public_key().public_bytes( + serialization.Encoding.PEM, + serialization.PublicFormat.SubjectPublicKeyInfo).decode() + monkeypatch.setenv('SENDGRID_WEBHOOK_PUBLIC_KEY', public_pem) + raw = json.dumps([{ + 'event': 'delivered', 'sg_event_id': 'event-1', + 'sg_message_id': 'provider-1', 'delivery_key': job.provider_key, + }], separators=(',', ':')).encode() + timestamp = '1783790000' + signature = base64.b64encode(private.sign( + timestamp.encode() + raw, ec.ECDSA(hashes.SHA256()))).decode() + gateway, sessions, _journal, _clock = build_gateway(tmp_path / 'gateway') + gateway.quote_store = store + client = TestClient(create_app(gateway_bundle=(gateway, sessions))) + headers = { + 'X-Twilio-Email-Event-Webhook-Timestamp': timestamp, + 'X-Twilio-Email-Event-Webhook-Signature': signature, + } + first = client.post('/webhooks/sendgrid/events', content=raw, headers=headers) + replay = client.post('/webhooks/sendgrid/events', content=raw, headers=headers) + assert first.status_code == replay.status_code == 200 + assert first.json()['new_events'] == 1 and replay.json()['new_events'] == 0 + assert store.delivery_job( + 'tenant_001', document.header.quote_number, 1, 'email').status == 'delivered' + assert b'provider-1' in (tmp_path / 'callbacks.sqlite').read_bytes() diff --git a/tests/test_delivery_callbacks.py b/tests/test_delivery_callbacks.py new file mode 100644 index 0000000..78de9b0 --- /dev/null +++ b/tests/test_delivery_callbacks.py @@ -0,0 +1,60 @@ +from __future__ import annotations + +import base64 +import json + +import pytest +from cryptography.hazmat.primitives import hashes, serialization +from cryptography.hazmat.primitives.asymmetric import ec + +from quoting.callbacks import ( + CallbackError, + parse_sendgrid_callbacks, + parse_twilio_callback, + verify_sendgrid_callback, + verify_twilio_callback, +) +from runtime.twilio_sig import expected_signature + + +def test_twilio_signature_and_normalization(): + url = 'https://api.example.test/webhooks/twilio/status' + params = { + 'MessageSid': 'SM1', 'MessageStatus': 'delivered', + 'DeliveryKey': 'stable-key', 'EventSid': 'EV1', + } + signature = expected_signature('token', url, params) + verify_twilio_callback( + auth_token='token', url=url, params=params, signature=signature) + callback = parse_twilio_callback(params) + assert callback.outcome == 'delivered' + assert callback.delivery_key == 'stable-key' + with pytest.raises(CallbackError): + verify_twilio_callback( + auth_token='token', url=url, params=params, signature='tampered') + + +def test_sendgrid_ecdsa_signature_duplicate_batch_and_normalization(): + private = ec.generate_private_key(ec.SECP256R1()) + public_pem = private.public_key().public_bytes( + serialization.Encoding.PEM, + serialization.PublicFormat.SubjectPublicKeyInfo).decode() + raw = json.dumps([{ + 'event': 'bounce', 'sg_event_id': 'event-1', + 'sg_message_id': 'message-1', 'delivery_key': 'stable-key', + 'email': 'must-not-be-retained@example.test', + }], separators=(',', ':')).encode() + timestamp = '1783790000' + signature = base64.b64encode(private.sign( + timestamp.encode() + raw, ec.ECDSA(hashes.SHA256()))).decode() + verify_sendgrid_callback( + public_key_pem=public_pem, timestamp=timestamp, + signature=signature, raw_body=raw) + callbacks = parse_sendgrid_callbacks(raw) + assert callbacks[0].outcome == 'bounced' + assert callbacks[0].delivery_key == 'stable-key' + assert 'must-not-be-retained' not in repr(callbacks) + with pytest.raises(CallbackError): + verify_sendgrid_callback( + public_key_pem=public_pem, timestamp=timestamp, + signature=signature, raw_body=raw + b' ') diff --git a/tests/test_delivery_runtime_config.py b/tests/test_delivery_runtime_config.py new file mode 100644 index 0000000..9fe354a --- /dev/null +++ b/tests/test_delivery_runtime_config.py @@ -0,0 +1,34 @@ +from __future__ import annotations + +import pytest + +from runtime.config import build_delivery_adapters, build_quote_deliverer + + +def test_production_cannot_construct_simulated_deliverer(monkeypatch): + monkeypatch.setenv('SKU_ENVIRONMENT', 'production') + with pytest.raises(RuntimeError, match='forbidden'): + build_quote_deliverer() + + +def test_live_adapters_fail_at_startup_without_credentials(monkeypatch): + monkeypatch.setenv('SKU_DELIVERY_CHANNELS', 'email') + monkeypatch.delenv('SENDGRID_API_KEY', raising=False) + monkeypatch.delenv('SENDGRID_FROM_EMAIL', raising=False) + monkeypatch.delenv('SKU_EMAIL_ALLOWLIST', raising=False) + with pytest.raises(ValueError, match='credentials'): + build_delivery_adapters() + + +def test_staging_adapters_require_and_preserve_allowlists(monkeypatch): + monkeypatch.setenv('SKU_DELIVERY_CHANNELS', 'email,sms') + monkeypatch.setenv('SENDGRID_API_KEY', 'secret') + monkeypatch.setenv('SENDGRID_FROM_EMAIL', 'sender@example.test') + monkeypatch.setenv('SKU_EMAIL_ALLOWLIST', 'allowed@example.test') + monkeypatch.setenv('TWILIO_ACCOUNT_SID', 'AC1') + monkeypatch.setenv('TWILIO_AUTH_TOKEN', 'token') + monkeypatch.setenv('TWILIO_FROM_NUMBER', '+15550109999') + monkeypatch.setenv('SKU_SMS_ALLOWLIST', '+15550100100') + adapters = build_delivery_adapters() + assert adapters['email'].allowed_recipients == frozenset({'allowed@example.test'}) + assert adapters['sms'].allowed_recipients == frozenset({'+15550100100'}) diff --git a/tests/test_m5_baseline_defects.py b/tests/test_m5_baseline_defects.py index 41c1e63..bb91f4a 100644 --- a/tests/test_m5_baseline_defects.py +++ b/tests/test_m5_baseline_defects.py @@ -186,17 +186,16 @@ def test_idempotent_retry_repairs_missing_outbox_row(tmp_path): assert gateway.quote_store.queue_item(number, 1) is not None -def test_defect_queue_claim_is_not_exclusive(tmp_path): +def test_artifact_queue_claim_is_exclusive(tmp_path): gateway, _, _, _, _, issued = _issue_one(tmp_path, "double-claim") number = issued.meta["quote_number"] - first = gateway.quote_store.claim_pending() - second = gateway.quote_store.claim_pending() - assert first is not None and second is not None + gateway.quote_store.release_for_call_completion( + tenant_id="tenant_001", conversation_id="double-claim", + completion_event_id=1, effective_hangup_at=1.0, received_at=2.0) + first = gateway.quote_store.claim_artifact(owner="worker-1", now=2.0) + second = gateway.quote_store.claim_artifact(owner="worker-2", now=2.0) + assert first is not None and second is None assert (first.quote_number, first.revision) == (number, 1) - assert (second.quote_number, second.revision) == (number, 1), ( - "same_pending_job_claimed_twice" - ) - _observed("same_pending_job_claimed_twice") def test_document_guard_rejects_cross_line_quantity_swap(tmp_path): diff --git a/tests/test_m5_release_config.py b/tests/test_m5_release_config.py index 1d11eda..b0e2db8 100644 --- a/tests/test_m5_release_config.py +++ b/tests/test_m5_release_config.py @@ -42,8 +42,7 @@ def test_current_m5_readiness_is_truthfully_red_and_names_separate_blockers(): assert report["readiness_kind"] == "m5_release" assert report["m5_ready"] is False blockers = "\n".join(report["blocking"]) - for required in ("capability:delivery", "policy:P02", "policy:P04", - "capability:live_provider"): + for required in ("policy:P02", "policy:P04"): assert required in blockers diff --git a/tests/test_m5_w5_delivery_worker.py b/tests/test_m5_w5_delivery_worker.py new file mode 100644 index 0000000..232f60d --- /dev/null +++ b/tests/test_m5_w5_delivery_worker.py @@ -0,0 +1,227 @@ +from __future__ import annotations + +import hashlib +from dataclasses import dataclass, replace +from datetime import datetime, timedelta, timezone +from decimal import Decimal + +from gateway import Account +from quoting.delivery import DeliveryPayload, ProviderResult +from quoting.models import ( + DeliveryPreflight, + QuoteDocument, + QuoteHeader, + QuoteLine, + QuotePresentation, + QuoteProvenance, + QuoteTerms, +) +from quoting.store import SqliteQuoteStore +from quoting.totals import build_totals, extended_price +from quoting.worker import run_artifact_once, run_delivery_once + +NOW = datetime(2026, 7, 11, 12, 0, tzinfo=timezone.utc) + + +def _factory(number: str, revision: int) -> QuoteDocument: + line = QuoteLine( + 1, 'K5-24SBC', 'tenant_001', 'chrome stack', 10, 'each', + Decimal('24.99'), NOW, '1001', extended_price(Decimal('24.99'), 10), + 'ships tomorrow', 'available', NOW + timedelta(days=1), NOW, + 'translator:verbatim', 'high') + return QuoteDocument( + QuoteHeader('tenant_001', number, revision, '1001', 'preferred', NOW, + NOW + timedelta(days=30), 'agent/gateway', 'catalog-v1'), + (line,), build_totals((line,)), + QuoteTerms(30, 'Prices subject to change.', 'Net 30.'), + QuotePresentation('Seller', '', 'default', 'Customer', '1001'), + QuoteProvenance('build', 'renderer-v1', 'guard-v1', 'agent-v1', + 'catalog-v1', ('price-v1',), ('inventory-v1',)), + DeliveryPreflight( + 'email', 'contact-1001', 'v1', + hashlib.sha256( + b'tenant_001:contact-1001:quotes1@example.test').hexdigest(), + 'q***@example.test', True)) + + +def _store(path, *, both_channels: bool = False, session_id: str = 'conv-1', + digest_char: str = 'd') -> tuple[SqliteQuoteStore, QuoteDocument]: + store = SqliteQuoteStore(path) + digest = digest_char * 64 + document = store.issue_atomic( + tenant_id='tenant_001', account_id='1001', session_id=session_id, + order_digest=digest, action='issue_quote', delivery_channel='email', + assent_receipt={ + 'session_id': session_id, 'turn_id': 'turn-2', + 'order_digest': digest, 'readback_receipt_id': f'readback:{digest[:24]}', + 'utterance_digest': 'a' * 64, 'line_versions': ['line:1'], + 'occurred_at': NOW.timestamp(), + }, document_factory=(_factory_both if both_channels else _factory)) + return store, document + + +def _factory_both(number: str, revision: int) -> QuoteDocument: + document = _factory(number, revision) + sms = DeliveryPreflight( + 'sms', 'contact-1001', 'v1', hashlib.sha256( + b'tenant_001:contact-1001:+15550100100').hexdigest(), + '***-***-0100', True) + return replace(document, additional_delivery_preflights=(sms,)) + + +@dataclass +class FakeAdapter: + result: ProviderResult + channel: str = 'email' + calls: list[tuple[DeliveryPayload, str]] = None + + def __post_init__(self): + self.calls = [] + + def send(self, payload: DeliveryPayload, *, delivery_key: str) -> ProviderResult: + self.calls.append((payload, delivery_key)) + return self.result + + +def _account(email='quotes1@example.test', version='v1'): + return Account( + '1001', 'Customer', '+15550100100', email, + contact_id='contact-1001', contact_version=version, + phone_verified=True, email_verified=True) + + +def test_completion_gate_atomic_leases_and_two_channel_cardinality(tmp_path): + path = tmp_path / 'worker.sqlite' + first, document = _store(path, both_channels=True) + assert first._conn.execute('SELECT status FROM artifact_jobs').fetchone()[0] == ( + 'waiting_for_call_end') + assert first.claim_artifact(owner='a', now=NOW.timestamp()) is None + assert first._conn.execute('SELECT COUNT(*) FROM artifact_jobs').fetchone()[0] == 1 + assert first._conn.execute('SELECT COUNT(*) FROM delivery_jobs').fetchone()[0] == 2 + + assert first.release_for_call_completion( + tenant_id='tenant_001', conversation_id='conv-1', + completion_event_id=7, effective_hangup_at=NOW.timestamp() - 1, + received_at=NOW.timestamp()) == 1 + second = SqliteQuoteStore(path) + winner = first.claim_artifact(owner='a', now=NOW.timestamp()) + loser = second.claim_artifact(owner='b', now=NOW.timestamp()) + assert winner is not None and loser is None + assert winner.render_started_at >= winner.effective_hangup_at + first.fail_artifact( + winner, status='retry_wait', detail='planted_crash', + now=NOW.timestamp(), next_attempt_at=NOW.timestamp()) + assert run_artifact_once( + store=second, owner='b', now_fn=lambda: NOW + timedelta(seconds=1)) + assert second._conn.execute('SELECT status FROM artifact_jobs').fetchone()[0] == 'verified' + assert second._conn.execute('SELECT status FROM delivery_jobs').fetchone()[0] == 'pending' + assert second.artifact('tenant_001', document.header.quote_number, 1, 'pdf') + + +def test_delivery_revalidates_contact_and_ambiguous_acceptance_never_retries(tmp_path): + store, document = _store(tmp_path / 'delivery.sqlite') + store.release_for_call_completion( + tenant_id='tenant_001', conversation_id='conv-1', completion_event_id=7, + effective_hangup_at=NOW.timestamp() - 1, received_at=NOW.timestamp()) + assert run_artifact_once(store=store, owner='artifact', now_fn=lambda: NOW) + adapter = FakeAdapter(ProviderResult(False, None, False, True, + 'acceptance_unknown')) + assert run_delivery_once( + store=store, owner='delivery', adapters={'email': adapter}, + account_lookup=lambda _: _account(), link_factory=lambda _: 'https://q.test/x', + now_fn=lambda: NOW) + job = store.delivery_job('tenant_001', document.header.quote_number, 1, 'email') + assert job.status == 'manual_review' + assert len(adapter.calls) == 1 + assert store.claim_delivery(owner='retry', now=NOW.timestamp() + 100) is None + + +def test_destination_or_version_change_requires_reauthorization(tmp_path): + for name, account in ( + ('destination', _account('changed@example.test')), + ('metadata', _account(version='v2'))): + store, document = _store(tmp_path / f'{name}.sqlite') + store.release_for_call_completion( + tenant_id='tenant_001', conversation_id='conv-1', completion_event_id=9, + effective_hangup_at=NOW.timestamp() - 1, received_at=NOW.timestamp()) + run_artifact_once(store=store, owner='artifact', now_fn=lambda: NOW) + adapter = FakeAdapter(ProviderResult(True, 'provider-1', False, False)) + run_delivery_once( + store=store, owner='delivery', adapters={'email': adapter}, + account_lookup=lambda _, value=account: value, + link_factory=lambda _: 'https://q.test/x', now_fn=lambda: NOW) + job = store.delivery_job( + 'tenant_001', document.header.quote_number, 1, 'email') + assert job.status == 'needs_reauthorization' + assert adapter.calls == [] + + +def test_duplicate_and_out_of_order_callbacks_do_not_regress_final_state(tmp_path): + store, document = _store(tmp_path / 'callbacks.sqlite') + store.release_for_call_completion( + tenant_id='tenant_001', conversation_id='conv-1', completion_event_id=11, + effective_hangup_at=NOW.timestamp() - 1, received_at=NOW.timestamp()) + run_artifact_once(store=store, owner='artifact', now_fn=lambda: NOW) + adapter = FakeAdapter(ProviderResult(True, 'provider-1', False, False)) + run_delivery_once( + store=store, owner='delivery', adapters={'email': adapter}, + account_lookup=lambda _: _account(), link_factory=lambda _: 'https://q.test/x', + now_fn=lambda: NOW) + job = store.delivery_job('tenant_001', document.header.quote_number, 1, 'email') + assert job.status == 'provider_accepted' + assert store.apply_provider_callback( + provider='sendgrid', event_id='event-delivered', + delivery_key=job.provider_key, provider_id='provider-1', + outcome='delivered', payload_digest='d' * 64, + received_at=NOW.timestamp() + 1) + assert not store.apply_provider_callback( + provider='sendgrid', event_id='event-delivered', + delivery_key=job.provider_key, provider_id='provider-1', + outcome='delivered', payload_digest='d' * 64, + received_at=NOW.timestamp() + 2) + assert store.apply_provider_callback( + provider='sendgrid', event_id='late-processed', + delivery_key=job.provider_key, provider_id='provider-1', + outcome='provider_accepted', payload_digest='e' * 64, + received_at=NOW.timestamp() + 3) + assert store.delivery_job( + 'tenant_001', document.header.quote_number, 1, 'email').status == 'delivered' + + +def test_crash_before_send_retries_but_after_send_boundary_requires_manual_review(tmp_path): + store, document = _store(tmp_path / 'crash.sqlite') + store.release_for_call_completion( + tenant_id='tenant_001', conversation_id='conv-1', completion_event_id=12, + effective_hangup_at=NOW.timestamp() - 1, received_at=NOW.timestamp()) + run_artifact_once(store=store, owner='artifact', now_fn=lambda: NOW) + claimed = store.claim_delivery( + owner='crashed-before-send', now=NOW.timestamp(), lease_seconds=1) + assert claimed is not None + recovered = store.claim_delivery( + owner='recovery', now=NOW.timestamp() + 2, lease_seconds=1) + assert recovered is not None + store.mark_delivery_sending(recovered, now=NOW.timestamp() + 2) + assert store.claim_delivery( + owner='unsafe-retry', now=NOW.timestamp() + 4, lease_seconds=1) is None + job = store.delivery_job( + 'tenant_001', document.header.quote_number, 1, 'email') + assert job.status == 'manual_review' + assert job.error_class == 'acceptance_unknown' + + +def test_quarantined_poison_artifact_does_not_block_later_work(tmp_path): + path = tmp_path / 'poison.sqlite' + store, first = _store(path) + store, second = _store(path, session_id='conv-2', digest_char='e') + for conversation, event in (('conv-1', 1), ('conv-2', 2)): + store.release_for_call_completion( + tenant_id='tenant_001', conversation_id=conversation, + completion_event_id=event, effective_hangup_at=NOW.timestamp() - 1, + received_at=NOW.timestamp()) + poison = store.claim_artifact(owner='worker', now=NOW.timestamp()) + assert poison.quote_number == first.header.quote_number + store.fail_artifact( + poison, status='quarantined', detail='planted_poison', + now=NOW.timestamp()) + later = store.claim_artifact(owner='worker', now=NOW.timestamp()) + assert later.quote_number == second.header.quote_number diff --git a/tests/test_post_call.py b/tests/test_post_call.py index 897b967..7262047 100644 --- a/tests/test_post_call.py +++ b/tests/test_post_call.py @@ -8,6 +8,7 @@ import pytest from fastapi.testclient import TestClient +from gateway_fixtures import build_gateway from gateway.conversation_store import ConversationStore from runtime.app import create_app @@ -126,6 +127,34 @@ def test_delayed_event_is_accepted_using_signed_delivery_time( assert stored.received_at > stored.effective_hangup_at +def test_authenticated_post_call_releases_matching_quote_outbox_only( + tmp_path, monkeypatch): + _configure(monkeypatch, tmp_path) + gateway, sessions, _journal, _clock = build_gateway(tmp_path / 'gateway') + gateway.conversation_store = ConversationStore(tmp_path / 'calls.db') + caller = 'conversation-1' + token = sessions.open(caller, caller) + gateway.converse(caller, token, 'my account number is 1001') + gateway.converse(caller, token, '10 of the K5-24SBC') + gateway.converse(caller, token, 'place the order') + issued = gateway.converse(caller, token, 'confirm the quote') + number = issued.meta['quote_number'] + assert gateway.quote_store._conn.execute( + 'SELECT status FROM artifact_jobs WHERE quote_number=?', + (number,)).fetchone()[0] == 'waiting_for_call_end' + client = TestClient(create_app(gateway_bundle=(gateway, sessions))) + raw, signature = _signed(_payload()) + result = client.post( + '/webhooks/elevenlabs/post-call', content=raw, + headers={'ElevenLabs-Signature': signature}) + assert result.status_code == 200 + row = gateway.quote_store._conn.execute( + 'SELECT status,effective_hangup_at FROM artifact_jobs WHERE quote_number=?', + (number,)).fetchone() + assert row['status'] == 'pending' + assert row['effective_hangup_at'] == 1_750_000_042 + + def test_signature_accepts_rotated_v0_candidate_and_rejects_malformed(): body = b'{}' at = int(time.time()) @@ -180,6 +209,22 @@ def test_reconciliation_qualifies_once_and_permanent_loss_holds_with_alert( alerts = store.completion_alerts('tenant_001', 'lost') assert len(alerts) == 1 assert alerts[0]['reason'] == 'call_completion_missing_manual_review' + + +def test_reconciliation_hook_releases_outbox_idempotently(tmp_path): + store = ConversationStore(tmp_path / 'calls.db') + source = _StatusSource(ReconciledStatus( + done=True, effective_hangup_at=105.0, agent_id='agent-1', + branch_id='branch-1', version_id='version-7', environment='staging')) + completions = [] + reconciler = CallCompletionReconciler( + store, source, _policy(), completion_hook=completions.append) + first = reconciler.reconcile( + tenant_id='tenant_001', conversation_id='c1', held_since=100, now=110) + second = reconciler.reconcile( + tenant_id='tenant_001', conversation_id='c1', held_since=100, now=111) + assert first == second + assert completions == [first, first] assert store.call_completion('tenant_001', 'lost') is None diff --git a/tests/test_quote_artifacts.py b/tests/test_quote_artifacts.py index a379158..8f47447 100644 --- a/tests/test_quote_artifacts.py +++ b/tests/test_quote_artifacts.py @@ -64,7 +64,7 @@ def test_artifact_transaction_fault_has_no_partial_format(tmp_path, fault_after) assert store._conn.execute( 'SELECT COUNT(*) FROM quote_artifacts').fetchone()[0] == 0 assert store._conn.execute( - 'SELECT status FROM artifact_jobs').fetchone()[0] == 'pending' + 'SELECT status FROM artifact_jobs').fetchone()[0] == 'waiting_for_call_end' def _service(store, tokens=None): diff --git a/tests/test_quote_migration.py b/tests/test_quote_migration.py index 4f7dba6..4fe2de2 100644 --- a/tests/test_quote_migration.py +++ b/tests/test_quote_migration.py @@ -66,14 +66,14 @@ def test_p16_upgrade_preserves_legacy_rows_and_backfills_explicit_held_lineage( report = migrate_quote_database(path) assert _legacy_checksum(path) == before - assert report['schema_version'] == 3 + assert report['schema_version'] == 4 assert report['counts'] == { 'quotes': 1, 'quote_issue_keys': 1, 'issuance_events': 1, 'artifact_jobs': 1, 'delivery_jobs': 1, } store = SqliteQuoteStore(path) assert store._conn.execute( - 'SELECT version FROM quote_schema').fetchone()[0] == 3 + 'SELECT version FROM quote_schema').fetchone()[0] == 4 key = store._conn.execute('SELECT * FROM quote_issue_keys').fetchone() assert (key['tenant_id'], key['account_id'], key['session_id'], key['order_digest'], key['action']) == ( diff --git a/tests/test_quote_worker.py b/tests/test_quote_worker.py index d0bfd87..45d765f 100644 --- a/tests/test_quote_worker.py +++ b/tests/test_quote_worker.py @@ -9,7 +9,7 @@ import pytest -from gateway import Account, ConversationJournal, EventType +from gateway import Account, ConversationJournal from quoting.capabilities import CapabilityKeyring, CapabilityService from quoting.delivery import SimulatedDeliverer from quoting.document_guard import verify_rendered_document @@ -57,7 +57,8 @@ def _document() -> QuoteDocument: 'build', 'renderer-v1', 'guard-v1', 'agent-v1', 'deadbeef', ('price-v1',), ('inventory-v1',)), delivery_preflight=DeliveryPreflight( - 'email', 'c1', 'v1', 'fingerprint', + 'email', 'c1', 'v1', hashlib.sha256( + b'tenant_001:c1:quotes@example.test').hexdigest(), 'q***@example.test', True)) @@ -66,13 +67,24 @@ def _queued(tmp_path): document = _document() store.save(document) store.enqueue(document.header.quote_number, document.header.revision) + store._conn.execute( + 'INSERT INTO artifact_jobs(tenant_id,quote_number,revision,status,' + 'completion_event_id,effective_hangup_at) VALUES(?,?,?,"pending",1,?)', + ('tenant_001', document.header.quote_number, 1, NOW.timestamp() - 1)) + store._conn.execute( + 'INSERT INTO delivery_jobs(tenant_id,quote_number,revision,channel,' + 'completion_event_id,effective_hangup_at) VALUES(?,?,?,"email",1,?)', + ('tenant_001', document.header.quote_number, 1, NOW.timestamp() - 1)) + store._conn.commit() journal = ConversationJournal(path=tmp_path / 'journal.jsonl', now_fn=lambda: NOW.isoformat()) return store, document, journal def _account(_account_id): - return Account('1001', 'DEMO TRUCK CENTER', email='quotes@example.test') + return Account( + '1001', 'DEMO TRUCK CENTER', email='quotes@example.test', + contact_id='c1', contact_version='v1', email_verified=True) def _link_factory(store): @@ -125,17 +137,14 @@ def counted_guard(doc, *, pdf_bytes, xlsx_bytes): assert guard_calls == 1 assert deliverer.payloads == [] - row = store.queue_item(document.header.quote_number, 1) - assert row.status == 'blocked' + row = store._conn.execute('SELECT status FROM artifact_jobs').fetchone() + assert row['status'] == 'quarantined' artifact_job = store._conn.execute( 'SELECT status FROM artifact_jobs').fetchone() assert artifact_job['status'] == 'quarantined' assert store._conn.execute( 'SELECT COUNT(*) FROM artifact_guard_events').fetchone()[0] == 1 - assert len(escalations) == 1 and escalations[0][0] == 'document_parity' - blocked = journal.events(EventType.QUOTE_DELIVERY_BLOCKED) - assert len(blocked) == 1 - assert blocked[0]['quote_number'] == document.header.quote_number + assert escalations == [] def test_untampered_positive_control_delivers_exactly_once(tmp_path): @@ -152,9 +161,9 @@ def test_untampered_positive_control_delivers_exactly_once(tmp_path): assert payload.pdf_bytes.startswith(b'%PDF-1.') assert payload.xlsx_bytes.startswith(b'PK') assert payload.signed_link - assert store.queue_item(document.header.quote_number, 1).status == 'delivered' - delivered = journal.events(EventType.QUOTE_DELIVERED) - assert delivered[-1]['quote_number'] == document.header.quote_number + assert store.delivery_job( + 'tenant_001', document.header.quote_number, 1, 'email').status == ( + 'provider_accepted') def test_missing_contact_of_record_blocks_without_delivery(tmp_path): @@ -168,11 +177,10 @@ def test_missing_contact_of_record_blocks_without_delivery(tmp_path): now_fn=lambda: NOW) is True assert deliverer.payloads == [] - row = store.queue_item(document.header.quote_number, 1) - assert row.status == 'blocked' - assert row.detail == 'contact_missing' - blocked = journal.events(EventType.QUOTE_DELIVERY_BLOCKED) - assert blocked[-1]['reason'] == 'contact_missing' + row = store.delivery_job( + 'tenant_001', document.header.quote_number, 1, 'email') + assert row.status == 'needs_reauthorization' + assert row.error_class == 'contact_changed' def test_no_pending_work_is_a_noop(tmp_path):