diff --git a/CHANGELOG.md b/CHANGELOG.md index cdb4d96..579fdff 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- Add opt-in live response authentication through `require_authenticated_response`. + The session-MAC profile binds readable provenance or denial to the appraised + peer key, complete request and fresh recipient/session. Verification is expiring + and single-use; failed responses retain unknown execution outcomes without retry. + This is not portable signed evidence, response encryption or a new hardware result. + +- Record acceptance-harness operation timing and unfinished work so transport + timeouts remain distinguishable from security rejections. The historical SNP + burst timeout remains unexplained; this change adds diagnostics without retries. + - **Delegation revocation.** Until now a delegated grant could not be withdrawn inside its validity window (`docs/spec/threat-model.md`). The issuer of a credential, or any issuer above it in the chain, can now sign a diff --git a/LIMITATIONS.md b/LIMITATIONS.md index c2850f2..f8325bb 100644 --- a/LIMITATIONS.md +++ b/LIMITATIONS.md @@ -1,5 +1,16 @@ # Limitations +## Opt-in response authentication + +`send_task(require_authenticated_response=True)` authenticates readable +provenance or denial to the live caller using a request-bound session MAC. The +default legacy path remains unauthenticated. The profile supplies one-use, +expiring client verification, not request-execution deduplication, portable +signed receipts or confidential output encryption. Both session parties know the +MAC key. Hardware assurance still requires explicit hardware appraisal; the new +profile has software/HTTP tests and no new live-hardware validation. See the +[profile and deployment limits](docs/spec/response-authentication.md). + cA2A 0.2 is a Developer Preview with a runnable, tested profile and runtime. This document states plainly what is built, what remains before 1.0, and what is out of scope, so no claim in the documentation runs ahead of the code. This is a deliberate discipline: proof, not promises. ## What is built diff --git a/docs/mutual-hardware-acceptance.md b/docs/mutual-hardware-acceptance.md index e96157f..2e25912 100644 --- a/docs/mutual-hardware-acceptance.md +++ b/docs/mutual-hardware-acceptance.md @@ -12,7 +12,7 @@ Policies forbidding SMT or requiring ciphertext hiding correctly rejected the hosts. Measurement pins came from owner-SSH bootstrap reports, not a precomputed application image. Initial burst runs timed out; paced runs exited 0. The sanitized results are in -[the diagnostic record](../experiments/hardware-validation/mutual-snp-2026-09-17/README.md). +[the diagnostic record](https://github.com/agentrust-io/ca2a/blob/main/experiments/hardware-validation/mutual-snp-2026-09-17/README.md). Raw reports and device-identifying certificates remain private. Both VMs and their boot disks were deleted after collection. This is not end-to-end inference confidentiality validation. @@ -125,3 +125,26 @@ secure version. These fields are checked after signature, chain and binding verification. They do not establish firmware TCB currency, revocation status or safe migration policy. Other VMPL deployments require a separately justified profile and are deliberately refused here. + +## Diagnosing interrupted calls + +The harness writes `operation_started`, `operation_completed` and +`operation_failed` observations with an operation ID and monotonic elapsed time. +Stages distinguish receiver offer generation, local quote collection, peer quote +verification, receiver task processing and the overall sender call. Nonces are +hashed for correlation; payloads, certificates and raw reports are not logged. + +A start without a matching completion identifies unfinished observed work; it +does not identify why it stalled. A sender transport failure records the receiver +outcome as unknown and preserves the original exception. There are no automatic +retries: the receiver may have processed a task before its response was lost. +Collect both hosts' receipts and service/kernel logs before diagnosing a cause. +Operation IDs correlate events within a receipt, not an authenticated distributed +trace. These observations remain unsigned and cannot prove receiver inactivity. + +The September 17 burst timeout has not been reproduced locally: 100 consecutive +software-provider calls over the reference HTTP server completed without pacing. +A controlled local test blocks software quote generation to verify that a real +HTTP timeout leaves useful unfinished-stage evidence. That test does not establish +that quote generation caused the historical SNP timeout. Hardware results have +not been rerun or changed by this diagnostic addition. diff --git a/docs/spec/error-codes.md b/docs/spec/error-codes.md index f6627d7..e4b77c7 100644 --- a/docs/spec/error-codes.md +++ b/docs/spec/error-codes.md @@ -58,9 +58,16 @@ except CA2AError as exc: Verification fails closed. `verify_chain`, `verify_dag`, and `cross_check_chain` raise the first error they find rather than returning a partial result, so a caught `CA2AError` means the chain or DAG was rejected. -## See also - -- [Delegation Chain](delegation-chain.md) for the checks behind `ScopeEscalation`, `BrokenDelegationLink`, `DelegationDepthExceeded`, `CredentialReplay`, `CredentialNotYetValid`, `CredentialExpired`, `CredentialRevoked`, `RevocationStatusUnknown`, and `InvalidRevocation`. +## See also + +The opt-in [response authentication profile](response-authentication.md) adds +`ResponseAuthenticationFailed` (`RESPONSE_AUTHENTICATION_FAILED`, HTTP 400). +At the caller it means no authenticated response was established and execution +may already have occurred. `response.AuthenticatedPeerError` instead carries a +MAC-authenticated denial's status/code and locally verified response metadata; +it does not by itself prove absence of side effects. + +- [Delegation Chain](delegation-chain.md) for the checks behind `ScopeEscalation`, `BrokenDelegationLink`, `DelegationDepthExceeded`, `CredentialReplay`, `CredentialNotYetValid`, `CredentialExpired`, `CredentialRevoked`, `RevocationStatusUnknown`, and `InvalidRevocation`. - [Provenance DAG](provenance-dag.md) for the checks behind `ProvenanceLinkBroken`. - [Verification Library](verification-library.md) for `verify_chain`, `verify_chain_file`, `verify_dag`, and `cross_check_chain`. - [Failure Modes](failure-modes.md) for how these errors map to observable runtime behavior. diff --git a/docs/spec/mutual-attestation.md b/docs/spec/mutual-attestation.md index 522f581..981a59e 100644 --- a/docs/spec/mutual-attestation.md +++ b/docs/spec/mutual-attestation.md @@ -118,6 +118,12 @@ that was wrong in this protocol. If a genuinely confidential response is ever added, sealing *that* to the caller's key is the right move. Encrypting the provenance record is not. + Authentication of the returned provenance is a separate extension. See the + [proposed response-binding requirements](response-binding-requirements.md) + for its acceptance matrix and the opt-in + [live response authentication profile](response-authentication.md) for the + implemented session-MAC path. This does not encrypt provenance or sign lineage. + ## The state problem A challenge is worth nothing unless it is single-use and expiring, and the diff --git a/docs/spec/response-authentication.md b/docs/spec/response-authentication.md new file mode 100644 index 0000000..8ef4be7 --- /dev/null +++ b/docs/spec/response-authentication.md @@ -0,0 +1,132 @@ +# Live response authentication + +`ca2a-response-mac-v1` implements the live-caller branch of +[the response requirements](response-binding-requirements.md), tracked in +[#188](https://github.com/agentrust-io/ca2a/issues/188). This is an opt-in +reference-transport extension. It does not sign the execution lineage or create +third-party-verifiable receipts. + +## Using the profile + +Call `transport.client.send_task` with `require_authenticated_response=True`. +Configure `require_hardware=True` and an independently pinned verifier when a +hardware-appraised peer is required. These are separate controls: authenticated +responses to a software-only handshake still have `peer_assurance="none"`. + +The caller uses the existing handshake and holder-bound request, then submits an +envelope to `POST /ca2a/task/authenticated`. The server validates its response +context before admitting the nested task. An old server or unsigned reply cannot +downgrade a strict caller to the legacy endpoint. The legacy endpoint and default +client behavior remain available, without authenticated-response assurance. + +Successful strict calls return the ordinary provenance body plus locally added +`response_authentication` metadata containing `profile`, `peer_assurance`, +`request_sha256` and `audience="live-caller-only"`. Legacy calls strip this +reserved field so a peer cannot inject local verification metadata. Applications +must set the strict option themselves; a dict received from another source is +not a verified object merely because it carries the same field names. + +Authenticated denials raise `response.AuthenticatedPeerError`, with the bound +HTTP status, peer error code and `response_authentication`. Unauthenticated, +malformed, expired or tampered replies raise `ResponseAuthenticationFailed`. +Transport failure after submission has an unknown execution outcome. Neither a +timeout nor a MAC failure proves non-execution. The client does not retry or +follow POST redirects and ignores ambient proxy settings for this submission. + +## Wire and key derivation + +The request envelope has exactly these fields: + +| Field | Value | +|---|---| +| `profile` | `ca2a-response-mac-v1` | +| `peer` | Appraised responder X25519 public key, 32 bytes in lowercase hex. | +| `recipient` | Fresh caller X25519 public key for this occurrence, same encoding. | +| `session` | Fresh random 32-byte occurrence identifier, lowercase hex. | +| `request` | Complete cA2A A2A message, including delegation, holder proof, challenge, requested action and sealed payload where present. | + +Both endpoints compute `request_sha256 = SHA256(JCS(envelope))`. The key and +session fields are therefore inside the commitment, alongside the request's +challenge and security-relevant contents. The server requires `peer` to match +its current channel key. Restart/rotation invalidates old contexts; a new call +must appraise the new key. + +Use [X25519](https://www.rfc-editor.org/rfc/rfc7748) between the fresh caller key +and appraised responder key. Reject an invalid/all-zero shared secret. Derive +32 bytes with [HKDF-SHA256](https://www.rfc-editor.org/rfc/rfc5869): salt is the +raw request digest; info is ASCII `ca2a/response-mac/v1/key`. The request digest +commits both public keys. No Ed25519 key is inferred from the X25519 key. + +The response statement has exactly `profile`, `request_sha256`, `status`, +`kind` and `body`. `kind` is `provenance` for status 200 and `denial` for +400 through 599. Other kinds, including confidential output, are unsupported. +`body` is an object: provenance has `accepted: true`; denial has an `error` +object. The reference server never echoes the opened task payload. + +Append `mac`, the lowercase hex HMAC-SHA256 under the derived key over: + +```text +ASCII("ca2a/response-mac/v1/statement") || 0x00 || JCS(statement) +``` + +The verifier checks the exact field set, profile, request digest, integer status, +actual HTTP status, MAC, kind and body shape before returning content. Merely +changing the outer HTTP status cannot turn an authenticated result into a denial. + +This uses the repository's integer-only [RFC 8785 JCS](https://www.rfc-editor.org/rfc/rfc8785) +subset: strings, objects, arrays, booleans, null and safe integers. Floating +point, nonfinite values, duplicate object fields, non-string object keys and +invalid Unicode are refused. The authenticated request/response path is bounded +to 1 MiB and nesting depth 64; unknown envelope fields are refused. HTTP input +is byte-bounded before decoding. These limits do not provide admission rate limits. + +## Replay, time and compatibility + +`PendingResponse` is process-local state for one occurrence. It consumes its +verification opportunity atomically even on failure, and checks its monotonic +deadline before and after authentication. The default is 30 seconds, configurable +on the primitive within `(0, 300]`. Nonfinite or backward clock observations +fail closed. Simultaneous deliveries have at most one accepted response per +pending object. Copying or serializing pending state is rejected. + +A fresh call generates a new recipient key and session, even if record IDs or +request content repeat. A response captured for an earlier call cannot satisfy +the new commitment. There is no pending-state persistence, multi-process recovery +or automatic resubmission after restart. Losing the pending object loses the +ability to authenticate its reply. Cloning a whole process/VM snapshot is outside +this software guarantee and requires rollback-resistant deployment controls. + +This is response replay protection. The server's holder challenge remains +stateless and replayable within its documented window. Duplicating an incoming +request can still execute it again. Applications needing idempotent/exactly-once +effects must supply an authoritative request-execution store; neither record ID +reuse nor a response MAC establishes that property. + +## Evidence and limits + +The committed `tests/fixtures/response-mac-v1.json` contains public synthetic test +keys, canonical request bytes, a derived key, a valid response and negative +mutations. Its HKDF/MAC values were calculated independently using direct HMAC +extract/expand and ASCII JSON canonicalization. Unit tests also cross-check that +derivation independently of the production HKDF helper. These fixtures contain +no hardware evidence or production secrets. + +Real-HTTP tests exercise successful provenance, authenticated denial, legacy +behavior, content/status substitution, cross-request replay, unsigned replies, +redirect refusal and reply loss after successful execution. Unit tests cover +concurrent delivery, expiry, malformed input, key rotation and unsupported output +kinds. Removing MAC verification, single-use enforcement or deadline checks makes +their regressions fail. This software result has not been rerun on live hardware. + +Both session parties can compute the MAC. It authenticates a live reply against +network substitution for the caller that retains its secret; it cannot convince +an independent third party which party authored the statement. It does not fix +unsigned execution history in #168. The MAC also does not encrypt provenance, +whose metadata can be sensitive. HTTPS and access controls remain deployment +requirements where that metadata needs confidentiality. + +Hardware assurance inherits the configured verifier and its approved measurement +and policy. No measured application identity, protected clock, rollback-resistant +memory, secure erasure, model execution or business completion follows from the +MAC itself. A future confidential-output profile needs explicit recipient-bound +encryption and disclosure rules; this profile does not add one. diff --git a/docs/spec/response-binding-requirements.md b/docs/spec/response-binding-requirements.md new file mode 100644 index 0000000..63c7f34 --- /dev/null +++ b/docs/spec/response-binding-requirements.md @@ -0,0 +1,113 @@ +# Authenticated response requirements + +Status: proposed requirements and acceptance matrix for [#188](https://github.com/agentrust-io/ca2a/issues/188). +The [live response authentication profile](response-authentication.md) now +implements a session-MAC verifier and reference adapter against this contract. +It selects live-caller authentication, not portable signatures. The requirements +below remain the acceptance basis; hardware validation is separate. + +## Current boundary + +At `0d2327600e66be4988706ef52b68a893ebebf603`, the reference client +`transport.client.send_task` appraises the callee offer, binds the outgoing +request to the holder proof, submits it, then returns a parsed HTTP 200 body. +It does not authenticate that response against the appraised callee or bind it +to the submitted request. Non-200 bodies become peer errors without that binding. +HTTPS can protect a connection to its configured endpoint; it does not by itself +bind a response to the workload appraised during the cA2A handshake. + +The server returns provenance through `wire.serialize_peer_result`, not the +opened confidential payload. The earlier decision to withdraw response sealing +in [mutual attestation](mutual-attestation.md) therefore remains applicable. +Readable evidence can still require authentication. Any future confidential +task output requires a separate recipient-bound encryption decision. + +## Required contract + +A response-required relying-party policy must be chosen locally before sending +the request. A peer cannot lower it by omitting a field or choosing an older +profile. Legacy mode must remain visibly lower assurance. + +The future authenticated statement must commit to: + +- A versioned response profile and a domain distinct from requests and other signatures. +- The responder key and its independently verified relationship to the appraised + workload and effective policy. The existing X25519 channel key cannot simply + be treated as a signing key. +- The intended recipient and one request occurrence, with an unambiguous commitment + to the complete security-relevant request. Record ID equality alone is insufficient. +- The session/handshake context, including freshness and identity inputs needed + to prevent substitution across sessions or peers. +- The response kind, authenticated outcome and exact response content commitment. + Any status that affects interpretation must be inside the authenticated statement. + +The design must define canonical bytes, field types, duplicate-field rejection, +encoding and size limits, key establishment/rotation and verification ordering. +If request bytes are committed directly, both sides must agree on those exact +bytes; if a canonical representation is used, specify it before producing vectors. +Do not introduce an ad hoc encoding through the test generator. + +Portable signatures and session authentication have different verification +audiences. The design must choose whether offline third parties can authenticate +the response, or whether only the live caller can do so. It must not claim the +former from a shared session MAC. Neither choice by itself fixes the separate +execution-lineage issue [#168](https://github.com/agentrust-io/ca2a/issues/168). + +## Outcomes and state + +Keep four outcomes distinguishable: + +| Outcome | Permitted interpretation | +|---|---| +| Authenticated acknowledgement/result | The bound peer made this statement about this request; execution or external completion requires the corresponding evidence. | +| Authenticated denial | The bound peer reports a denial at a defined stage; absence of side effects requires a stated and supported pre-execution boundary. | +| Unauthenticated response or transport failure | No authenticated peer verdict is established. | +| Timeout or lost response after submission | The receiver's execution outcome remains unknown; automatic retry is not justified by the timeout. | + +Define request occurrence identity separately from an idempotency key. A replay +policy must cover concurrent responses, duplicate delivery, expiry, restart, +multiple client instances and persistence/rollback assumptions. If replay state +is unavailable, a strict client cannot silently accept a previously consumed +response. A receiver's replay defense and a caller's response replay defense are +separate controls. + +## Acceptance matrix + +These are contract requirements. The implemented live-session profile exercises +them in `tests/unit/test_response_authentication.py`, using software evidence; +it does not establish hardware behavior. Each negative case needs a valid +control under the same policy. Refusal means no authenticated +result is exposed to the application; it does not assert that the receiver did +not execute the submitted request. + +| ID | Case | Required observation | +|---|---|---| +| RESP-001 | Valid response from the appraised responder for the current occurrence | Accepted with explicit response assurance and bound context. | +| RESP-002 | Response contents changed after authentication | Refused before application consumption. | +| RESP-003 | Valid proof from another responder | Refused, even if the response body is otherwise identical. | +| RESP-004 | Valid response moved to another request with a reused record ID | Refused using the complete request-occurrence binding. | +| RESP-005 | Response moved across sessions or intended recipients | Refused despite a valid proof under the original context. | +| RESP-006 | Missing proof, unknown profile or legacy response under strict policy | Refused without fallback. | +| RESP-007 | Success/denial kind or interpreted status substituted | Refused; the outer HTTP status alone cannot supply an authenticated verdict. | +| RESP-008 | Previously accepted response delivered again, including concurrent delivery | Replay policy enforced atomically within its declared state scope. | +| RESP-009 | Expired response, restarted verifier or unavailable replay state | Declared freshness/recovery policy enforced; no silent assurance promotion. | +| RESP-010 | Authenticated denial versus unsigned error body | The former is a peer statement; the latter remains unverified. | +| RESP-011 | Submission followed by timeout or truncated response | Unknown execution outcome retained; no blind resubmission. | +| RESP-012 | Duplicate fields, ambiguous encoding, oversized proof or malformed envelope | Bounded parsing refusal without treating input as authenticated. | +| RESP-013 | Responder key rotated without an approved binding to the appraised workload | Refused; approved rotation has a separate positive control. | +| RESP-014 | Confidential output added to otherwise readable provenance | Not released until a separately approved encrypted-output profile binds the recipient. | + +Run the contract through an actual HTTP exchange as well as verifier unit tests. +Observe whether unverified content reaches the application. Removing body, +request, peer or replay verification must fail the corresponding regression. +Synthetic cryptography, captured hardware evidence and live hardware execution +must be labeled separately. + +## Implementation state + +The [profile](response-authentication.md) selects X25519/HKDF/HMAC with the +existing JCS subset, live-caller verification and process-local one-use state. +Public synthetic vectors, a strict verifier, the HTTP adapter, compatibility +behavior and causal checks are included in the same handoff follow-up PR. +Portable signing, confidential outputs and hardware deployment validation remain +outside this profile. Issue #188 stays open during implementation review. diff --git a/mkdocs.yml b/mkdocs.yml index 5311f5a..5563778 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -165,6 +165,8 @@ nav: - Profile and transport: - A2A profile: docs/spec/profile.md - Transport binding: docs/spec/transport.md + - Response binding (proposal): docs/spec/response-binding-requirements.md + - Live response authentication: docs/spec/response-authentication.md - Sealed peer channel: docs/spec/sealed-channel.md - TRACE A2A profile: docs/spec/trace-a2a-profile.md - Delegation and policy: diff --git a/src/ca2a_runtime/errors.py b/src/ca2a_runtime/errors.py index d7523e8..2636f57 100644 --- a/src/ca2a_runtime/errors.py +++ b/src/ca2a_runtime/errors.py @@ -231,3 +231,9 @@ class TraceRecordInvalid(CA2AError): code = "TRACE_RECORD_INVALID" http_status = 422 + + +class ResponseAuthenticationFailed(TransportError): + """No authenticated verdict; a submitted request may already have executed.""" + + code = "RESPONSE_AUTHENTICATION_FAILED" diff --git a/src/ca2a_runtime/hardware_acceptance.py b/src/ca2a_runtime/hardware_acceptance.py index 4d740c6..c286920 100644 --- a/src/ca2a_runtime/hardware_acceptance.py +++ b/src/ca2a_runtime/hardware_acceptance.py @@ -9,7 +9,11 @@ import argparse import hashlib import json +import secrets import threading +import time +from collections.abc import Iterator +from contextlib import contextmanager from datetime import UTC, datetime from pathlib import Path from typing import Any @@ -18,13 +22,13 @@ from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey from cryptography.hazmat.primitives.hashes import SHA256 -from ca2a_runtime.attestation import Verifier +from ca2a_runtime.attestation import ChannelOffer, Verifier from ca2a_runtime.delegation.credential import DelegationCredential from ca2a_runtime.errors import AttestationFailed, AttestationUnsupported, CA2AError from ca2a_runtime.node import PeerNode from ca2a_runtime.peer import PeerResult from ca2a_runtime.policy import LocalPolicy -from ca2a_runtime.tee.base import AttestationReport +from ca2a_runtime.tee.base import AttestationReport, BaseProvider from ca2a_runtime.tee.sev_snp import SevSnpProvider from ca2a_runtime.transport import client, server from ca2a_verify.sev_snp import sev_snp_verifier @@ -50,6 +54,47 @@ def write(self, event: str, **values: Any) -> None: with self.lock, self.path.open("a", encoding="utf-8") as stream: stream.write(json.dumps(entry, sort_keys=True) + "\n") + @contextmanager + def span(self, stage: str, **values: Any) -> Iterator[None]: + """Retain a start even when work stalls; use a monotonic duration.""" + started = time.monotonic() + values = {**values, "operation_id": secrets.token_hex(8)} + self.write("operation_started", stage=stage, **values) + try: + yield + except Exception as exc: + self.write( + "operation_failed", + stage=stage, + elapsed_seconds=time.monotonic() - started, + error_type=type(exc).__name__, + **values, + ) + raise + else: + self.write( + "operation_completed", + stage=stage, + elapsed_seconds=time.monotonic() - started, + **values, + ) + + +class ObservedSnpProvider(SevSnpProvider): + """Time collection on the already selected SNP provider; do not retry it.""" + + def __init__(self, inner: BaseProvider, log: ReceiptLog) -> None: + self.inner = inner + self.log = log + + def attest(self, public_key: str, nonce: str) -> AttestationReport: + with self.log.span("local_quote_collection", nonce_sha256=_nonce_hash(nonce)): + return self.inner.attest(public_key, nonce) + + +def _nonce_hash(nonce: str) -> str: + return hashlib.sha256(nonce.encode()).hexdigest() + def real_provider() -> SevSnpProvider: if not SevSnpProvider.detect(): @@ -94,7 +139,8 @@ def build_verifier(config: dict[str, Any], base: Path, log: ReceiptLog) -> Verif def recorded(report: AttestationReport, nonce: str) -> str: try: - measurement = verify(report, nonce) + with log.span("peer_quote_verification", nonce_sha256=_nonce_hash(nonce)): + measurement = verify(report, nonce) except AttestationFailed: log.write("peer_appraisal", accepted=False) raise @@ -116,9 +162,14 @@ def __init__(self, *args: Any, log: ReceiptLog, **kwargs: Any) -> None: super().__init__(*args, **kwargs) self.log = log + def offer(self, nonce: str) -> ChannelOffer: + with self.log.span("receiver_offer", nonce_sha256=_nonce_hash(nonce)): + return super().offer(nonce) + def handle(self, message: dict[str, Any]) -> PeerResult: try: - result = super().handle(message) + with self.log.span("receiver_task_processing"): + result = super().handle(message) except CA2AError as exc: self.log.write("receiver_task", accepted=False, error_code=exc.code) raise @@ -133,7 +184,7 @@ def handle(self, message: dict[str, Any]) -> PeerResult: def run(config: dict[str, Any], base: Path, mode: str, log: ReceiptLog) -> None: - provider = real_provider() + provider = ObservedSnpProvider(real_provider(), log) verifier = build_verifier(config["peer"], base, log) if mode == "serve": issuers = config["trusted_root_issuers"] @@ -159,17 +210,24 @@ def run(config: dict[str, Any], base: Path, mode: str, log: ReceiptLog) -> None: bytes.fromhex((base / config["holder_key"]).read_text(encoding="utf-8").strip()) ) payload = (base / config["payload"]).read_bytes() - result = client.send_task( - config["peer_url"], - chain, - config["capability"], - config["record_id"], - holder_key=key, - payload=payload, - verifier=verifier, - require_hardware=True, - caller_provider=provider, - ) + try: + with log.span("sender_call"): + result = client.send_task( + config["peer_url"], + chain, + config["capability"], + config["record_id"], + holder_key=key, + payload=payload, + verifier=verifier, + require_hardware=True, + caller_provider=provider, + ) + except OSError: + # A timeout is neither an attestation rejection nor evidence that the + # receiver did no work. Preserve the original exception and never retry. + log.write("sender_transport_failure", receiver_outcome="unknown", retried=False) + raise if result.get("accepted") is not True or result.get("caller_attestation") != "hardware": raise AttestationFailed( "peer response does not report successful mutual hardware appraisal" diff --git a/src/ca2a_runtime/node.py b/src/ca2a_runtime/node.py index 765e8ae..fe496f7 100644 --- a/src/ca2a_runtime/node.py +++ b/src/ca2a_runtime/node.py @@ -20,7 +20,7 @@ from ca2a_runtime.challenge import DEFAULT_TTL_SECONDS, generate_secret, issue_challenge from ca2a_runtime.channel import generate_channel_keypair from ca2a_runtime.delegation.revocation import RevocationSnapshot -from ca2a_runtime.errors import ConfigError, TransportError +from ca2a_runtime.errors import CA2AError, ConfigError, TransportError from ca2a_runtime.peer import ( REQUIRE_HARDWARE, REQUIRE_NONE, @@ -120,3 +120,20 @@ def handle(self, message: dict[str, Any]) -> PeerResult: revocations=None if self.revocation_source is None else self.revocation_source(), max_revocation_staleness=self.max_revocation_staleness, ) + + def handle_authenticated(self, envelope: dict[str, Any]) -> tuple[int, dict[str, Any]]: + """Authenticate a provenance response using the appraised channel key. + + Invalid response contexts fail before admission. Runtime failures after + admission remain unknown outcomes unless an authenticated denial is made. + """ + from ca2a_runtime.response import ResponseProducer + from ca2a_runtime.transport.wire import serialize_error, serialize_peer_result + + producer = ResponseProducer(self._private_key, envelope) + try: + body = serialize_peer_result(self.handle(producer.envelope["request"])) + status = 200 + except CA2AError as exc: + body, status = serialize_error(exc), exc.http_status + return status, producer.finish(status, body) diff --git a/src/ca2a_runtime/response.py b/src/ca2a_runtime/response.py new file mode 100644 index 0000000..0d36ee8 --- /dev/null +++ b/src/ca2a_runtime/response.py @@ -0,0 +1,269 @@ +"""One-use live response authentication; not a portable signature or encryption.""" + +from __future__ import annotations + +import hashlib +import hmac +import json +import math +import secrets +import threading +import time +from collections.abc import Callable +from dataclasses import dataclass +from typing import Any + +from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey, X25519PublicKey +from cryptography.hazmat.primitives.hashes import SHA256 +from cryptography.hazmat.primitives.kdf.hkdf import HKDF + +from ca2a_runtime.attestation import VerifiedPeer +from ca2a_runtime.canonical import canonicalize +from ca2a_runtime.errors import CA2AError +from ca2a_runtime.errors import ResponseAuthenticationFailed as ResponseAuthenticationFailed + +PROFILE = "ca2a-response-mac-v1" +MAX_BYTES = 1 << 20 +_REQUEST_FIELDS = {"profile", "peer", "recipient", "session", "request"} +_RESPONSE_FIELDS = {"profile", "request_sha256", "status", "kind", "body", "mac"} + + +class AuthenticatedPeerError(CA2AError): + """An authenticated peer denial, not proof that no side effect occurred.""" + + def __init__(self, response: VerifiedResponse) -> None: + error = response.body["error"] + super().__init__(str(error.get("message", "peer denial"))) + self.code = str(error.get("code", "PEER_DENIAL")) + self.http_status = response.status + self.response_authentication = response.authentication() + + +def _fail() -> ResponseAuthenticationFailed: + return ResponseAuthenticationFailed( + "response authentication failed; execution outcome is unknown" + ) + + +def _object(pairs: list[tuple[str, Any]]) -> dict[str, Any]: + result: dict[str, Any] = {} + for key, value in pairs: + if key in result: + raise _fail() + result[key] = value + return result + + +def encode(value: dict[str, Any]) -> bytes: + """Bound the existing integer-only JCS profile, with string object keys.""" + + def check(item: Any, depth: int = 0) -> None: + if depth > 64: + raise _fail() + if isinstance(item, dict): + if any(not isinstance(key, str) for key in item): + raise _fail() + for child in item.values(): + check(child, depth + 1) + elif isinstance(item, list): + for child in item: + check(child, depth + 1) + + try: + check(value) + raw = canonicalize(value) + if len(raw) > MAX_BYTES: + raise _fail() + return raw + except (TypeError, ValueError, RecursionError) as exc: + raise _fail() from exc + + +def decode(raw: bytes) -> dict[str, Any]: + """Reject ambiguous/oversized JSON before it can be authenticated.""" + if len(raw) > MAX_BYTES: + raise _fail() + try: + value = json.loads(raw.decode("utf-8"), object_pairs_hook=_object) + if not isinstance(value, dict): + raise _fail() + encode(value) + return value + except (ValueError, UnicodeError, RecursionError) as exc: + raise _fail() from exc + + +def _hex(value: Any, length: int = 32) -> bytes: + if not isinstance(value, str) or len(value) != length * 2: + raise _fail() + try: + raw = bytes.fromhex(value) + except ValueError as exc: + raise _fail() from exc + if len(raw) != length or raw.hex() != value: + raise _fail() + return raw + + +def _key(private: X25519PrivateKey, public: str, request_hash: str) -> bytes: + try: + shared = private.exchange(X25519PublicKey.from_public_bytes(_hex(public))) + except ValueError as exc: + raise _fail() from exc + return HKDF( + algorithm=SHA256(), + length=32, + salt=bytes.fromhex(request_hash), + info=b"ca2a/response-mac/v1/key", + ).derive(shared) + + +def _mac(key: bytes, statement: dict[str, Any]) -> str: + return hmac.new( + key, b"ca2a/response-mac/v1/statement\x00" + encode(statement), "sha256" + ).hexdigest() + + +@dataclass(frozen=True) +class VerifiedResponse: + """A live caller's verified statement; both session parties can create its MAC.""" + + status: int + body: dict[str, Any] + request_sha256: str + peer_assurance: str + + def authentication(self) -> dict[str, str]: + return { + "profile": PROFILE, + "peer_assurance": self.peer_assurance, + "request_sha256": self.request_sha256, + "audience": "live-caller-only", + } + + +class ResponseProducer: + """Validate the request and establish the response key before executing it.""" + + def __init__(self, private: X25519PrivateKey, envelope: dict[str, Any]) -> None: + # Freeze caller-owned input before both hashing and execution. + self.envelope = decode(encode(envelope)) + if set(self.envelope) != _REQUEST_FIELDS or self.envelope["profile"] != PROFILE: + raise _fail() + peer = private.public_key().public_bytes_raw().hex() + if self.envelope["peer"] != peer or not isinstance(self.envelope["request"], dict): + raise _fail() + _hex(self.envelope["session"]) + self.request_sha256 = hashlib.sha256(encode(self.envelope)).hexdigest() + self._key = _key(private, self.envelope["recipient"], self.request_sha256) + + def finish(self, status: int, body: dict[str, Any]) -> dict[str, Any]: + statement = { + "profile": PROFILE, + "request_sha256": self.request_sha256, + "status": status, + "kind": "provenance" if status == 200 else "denial", + "body": body, + } + result = {**statement, "mac": _mac(self._key, statement)} + encode(result) + return result + + +class PendingResponse: + """Process-local, expiring, single-use verification state for one request. + + No persistence, restart recovery or resubmission. Every new call gets a fresh + key and occurrence ID. State loss means an unknown outcome, not an empty + replay cache with permission to accept an old response. + """ + + def __init__( + self, + peer: VerifiedPeer, + message: dict[str, Any], + *, + timeout: float = 30.0, + clock: Callable[[], float] = time.monotonic, + ) -> None: + if not math.isfinite(timeout) or not 0 < timeout <= 300: + raise ValueError("response timeout must be finite and in (0, 300]") + private = X25519PrivateKey.generate() + self._request = encode( + { + "profile": PROFILE, + "peer": peer.public_key, + "recipient": private.public_key().public_bytes_raw().hex(), + "session": secrets.token_hex(32), + "request": message, + } + ) + self.request_sha256 = hashlib.sha256(self._request).hexdigest() + self._key = _key(private, peer.public_key, self.request_sha256) + self._assurance = peer.assurance + self._clock = clock + self._started = clock() + if not math.isfinite(self._started): + raise ValueError("response clock must be finite") + self._deadline = self._started + timeout + self._used = False + self._lock = threading.Lock() + + @property + def request_bytes(self) -> bytes: + return self._request + + def __copy__(self) -> PendingResponse: + raise TypeError("pending response state cannot be copied or restored") + + def __deepcopy__(self, memo: dict[int, Any]) -> PendingResponse: + raise TypeError("pending response state cannot be copied or restored") + + def __reduce__(self) -> Any: + raise TypeError("pending response state cannot be copied or restored") + + def _check_deadline(self) -> None: + try: + now = self._clock() + if not math.isfinite(now) or not self._started <= now < self._deadline: + raise _fail() + except Exception as exc: + raise _fail() from exc + + def verify(self, status: int, raw: bytes) -> VerifiedResponse: + """Consume this occurrence even on failure; never expose unverified content.""" + with self._lock: + if self._used: + raise _fail() + self._used = True + key, self._key = self._key, b"" + try: + self._check_deadline() + envelope = decode(raw) + if set(envelope) != _RESPONSE_FIELDS: + raise _fail() + mac = envelope.pop("mac") + _hex(mac) + if ( + envelope["profile"] != PROFILE + or envelope["request_sha256"] != self.request_sha256 + or type(envelope["status"]) is not int + or type(status) is not int + or envelope["status"] != status + or not (status == 200 or 400 <= status <= 599) + ): + raise _fail() + if not hmac.compare_digest(mac, _mac(key, envelope)): + raise _fail() + body = envelope["body"] + kind = "provenance" if status == 200 else "denial" + if envelope["kind"] != kind or not isinstance(body, dict): + raise _fail() + if status == 200 and body.get("accepted") is not True: + raise _fail() + if status != 200 and not isinstance(body.get("error"), dict): + raise _fail() + self._check_deadline() + return VerifiedResponse(status, body, self.request_sha256, self._assurance) + except (ValueError, TypeError, OSError) as exc: + raise _fail() from exc diff --git a/src/ca2a_runtime/transport/client.py b/src/ca2a_runtime/transport/client.py index dd0dcfa..c1a5c2c 100644 --- a/src/ca2a_runtime/transport/client.py +++ b/src/ca2a_runtime/transport/client.py @@ -8,6 +8,7 @@ from __future__ import annotations +import http.client import json import secrets import urllib.error @@ -30,9 +31,14 @@ from ca2a_runtime.delegation.holder import build_holder_proof from ca2a_runtime.errors import AttestationFailed, CA2AError, TransportError from ca2a_runtime.peer import PeerRequest +from ca2a_runtime.response import ( + AuthenticatedPeerError, + PendingResponse, + ResponseAuthenticationFailed, +) from ca2a_runtime.tee.base import BaseProvider from ca2a_runtime.transport import a2a_adapter, wire -from ca2a_runtime.transport.server import CHANNEL_PATH, TASK_PATH +from ca2a_runtime.transport.server import AUTHENTICATED_TASK_PATH, CHANNEL_PATH, TASK_PATH _TIMEOUT = 10.0 _ALLOWED_SCHEMES = ("http", "https") @@ -91,6 +97,28 @@ def _post_json(url: str, body: dict[str, Any]) -> tuple[int, dict[str, Any]]: return exc.code, json.loads(_read_bounded(exc)) +class _NoRedirect(urllib.request.HTTPRedirectHandler): + def redirect_request( + self, req: Any, fp: Any, code: int, msg: str, headers: Any, newurl: str + ) -> None: + return None + + +def _post_authenticated(url: str, raw: bytes) -> tuple[int, bytes]: + _require_http_url(url) + # No proxy inheritance, redirect, or automatic replay of a submitted task. + opener = urllib.request.build_opener(urllib.request.ProxyHandler({}), _NoRedirect()) + request = urllib.request.Request( + url, data=raw, headers={"Content-Type": "application/json"}, method="POST" + ) + try: + with opener.open(request, timeout=_TIMEOUT) as response: + return response.status, _read_bounded(response) + except urllib.error.HTTPError as exc: + with exc: + return exc.code, _read_bounded(exc) + + @dataclass(frozen=True) class Handshake: """What one round trip to the handshake endpoint yields. @@ -143,6 +171,7 @@ def send_task( require_hardware: bool = False, parent_record_hash: str | None = None, caller_provider: BaseProvider | None = None, + require_authenticated_response: bool = False, ) -> dict[str, Any]: """Run the caller side end to end: verify the peer, prove holdership, seal, send. @@ -163,6 +192,10 @@ def send_task( ``require_hardware=True`` rejects software-only offers before sealing or sending a task. Configure the verifier's measurement and platform policy separately; hardware assurance alone does not specify either requirement. + + ``require_authenticated_response=True`` uses the one-use session-MAC profile + and rejects legacy/unsigned responses. Success carries locally established + ``response_authentication`` metadata. This is not a portable signature. """ sealed: bytes | None = None caller_offer: ChannelOffer | None = None @@ -206,10 +239,27 @@ def send_task( holder_proof=holder_proof, ) message = a2a_adapter.attach_ca2a_metadata({}, request) + if require_authenticated_response: + pending = PendingResponse(hello.peer, message) + try: + status, raw = _post_authenticated( + f"{base_url}{AUTHENTICATED_TASK_PATH}", pending.request_bytes + ) + response = pending.verify(status, raw) + except (OSError, http.client.HTTPException, TransportError) as exc: + raise ResponseAuthenticationFailed( + "no authenticated response; execution outcome is unknown" + ) from exc + if response.status != 200: + raise AuthenticatedPeerError(response) + return {**response.body, "response_authentication": response.authentication()} status, body = _post_json(f"{base_url}{TASK_PATH}", message) if status != 200: err = body.get("error", {}) raise _rehydrate_error(err) + # Only local verification may attach this metadata. A legacy peer cannot + # manufacture assurance by returning the reserved field itself. + body.pop("response_authentication", None) return body diff --git a/src/ca2a_runtime/transport/server.py b/src/ca2a_runtime/transport/server.py index 9cc8414..bbaf3c0 100644 --- a/src/ca2a_runtime/transport/server.py +++ b/src/ca2a_runtime/transport/server.py @@ -23,10 +23,12 @@ from ca2a_runtime.errors import CA2AError from ca2a_runtime.node import PeerNode +from ca2a_runtime.response import decode from ca2a_runtime.transport import wire CHANNEL_PATH = "/.well-known/ca2a/channel" TASK_PATH = "/ca2a/task" +AUTHENTICATED_TASK_PATH = "/ca2a/task/authenticated" _MAX_BODY = 1 << 20 # 1 MiB; fail closed on larger bodies _MAX_NONCE = 256 _READ_TIMEOUT_SECONDS = 10.0 @@ -102,7 +104,8 @@ def do_GET(self) -> None: ) def do_POST(self) -> None: - if urlparse(self.path).path != TASK_PATH: + path = urlparse(self.path).path + if path not in (TASK_PATH, AUTHENTICATED_TASK_PATH): self._send_json(404, {"error": {"code": "NOT_FOUND", "message": "unknown path"}}) return raw_length = self.headers.get("Content-Length", "") @@ -118,14 +121,19 @@ def do_POST(self) -> None: ) return try: - message = json.loads(self.rfile.read(length)) - except (json.JSONDecodeError, UnicodeDecodeError, TimeoutError): + raw = self.rfile.read(length) + message = decode(raw) if path == AUTHENTICATED_TASK_PATH else json.loads(raw) + except (json.JSONDecodeError, UnicodeDecodeError, TimeoutError, CA2AError): self._send_json(400, {"error": {"code": "BAD_REQUEST", "message": "invalid JSON"}}) return if not isinstance(message, dict): self._send_json(400, {"error": {"code": "BAD_REQUEST", "message": "object required"}}) return try: + if path == AUTHENTICATED_TASK_PATH: + status, body = self._node().handle_authenticated(message) + self._send_json(status, body) + return result = self._node().handle(message) except CA2AError as exc: self._send_json(exc.http_status, wire.serialize_error(exc)) diff --git a/tests/fixtures/response-mac-v1.json b/tests/fixtures/response-mac-v1.json new file mode 100644 index 0000000..95ed61e --- /dev/null +++ b/tests/fixtures/response-mac-v1.json @@ -0,0 +1,82 @@ +{ + "description": "Public synthetic test keys. No hardware evidence. ASCII JCS fixture.", + "callee_private_key": "0101010101010101010101010101010101010101010101010101010101010101", + "caller_private_key": "0202020202020202020202020202020202020202020202020202020202020202", + "request": { + "profile": "ca2a-response-mac-v1", + "peer": "a4e09292b651c278b9772c569f5fa9bb13d906b46ab68c9df9dc2b4409f8a209", + "recipient": "ce8d3ad1ccb633ec7b70c17814a5c76ecd029685050d344745ba05870e587d59", + "session": "0303030303030303030303030303030303030303030303030303030303030303", + "request": { + "record_id": "vector-1" + } + }, + "request_canonical_hex": "7b2270656572223a2261346530393239326236353163323738623937373263353639663566613962623133643930366234366162363863396466396463326234343039663861323039222c2270726f66696c65223a22636132612d726573706f6e73652d6d61632d7631222c22726563697069656e74223a2263653864336164316363623633336563376237306331373831346135633736656364303239363835303530643334343734356261303538373065353837643539222c2272657175657374223a7b227265636f72645f6964223a22766563746f722d31227d2c2273657373696f6e223a2230333033303330333033303330333033303330333033303330333033303330333033303330333033303330333033303330333033303330333033303330333033227d", + "shared_secret_hex": "2ed76ab549b1e73c031eb49c9448f0798aea81b698279a0c3dc3e49fbfc4b953", + "hkdf_key_hex": "7674b703aae99f7c0f56d7f1c3c419b20d563c2cb25a1b559d4b89667c0f0467", + "response": { + "profile": "ca2a-response-mac-v1", + "request_sha256": "c3cae3cbd25953b92b45a2cf1b1b794c34756823d0de7f495934bad5cb380ea7", + "kind": "provenance", + "status": 200, + "body": { + "accepted": true, + "record": { + "record_id": "vector-1" + } + }, + "mac": "a6a319d55db2b2b925bb2ee471b65bfcf5f6e2f07169df6f4c00eae3ca017e68" + }, + "negative_cases": [ + { + "id": "body-substitution", + "path": [ + "body", + "record", + "record_id" + ], + "value": "other", + "expected": "reject" + }, + { + "id": "request-substitution", + "path": [ + "request_sha256" + ], + "value": "0000000000000000000000000000000000000000000000000000000000000000", + "expected": "reject" + }, + { + "id": "downgrade", + "path": [ + "profile" + ], + "value": "legacy", + "expected": "reject" + }, + { + "id": "status-substitution", + "path": [ + "status" + ], + "value": 403, + "expected": "reject" + }, + { + "id": "confidential-output-kind", + "path": [ + "kind" + ], + "value": "confidential-output", + "expected": "reject" + }, + { + "id": "invalid-mac", + "path": [ + "mac" + ], + "value": "0000000000000000000000000000000000000000000000000000000000000000", + "expected": "reject" + } + ] +} diff --git a/tests/unit/test_hardware_acceptance.py b/tests/unit/test_hardware_acceptance.py index 6a21b74..34dfaff 100644 --- a/tests/unit/test_hardware_acceptance.py +++ b/tests/unit/test_hardware_acceptance.py @@ -1,12 +1,16 @@ """Harness controls with injected test doubles; these are not hardware results.""" import json +import threading +from concurrent.futures import ThreadPoolExecutor import pytest from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey from ca2a_runtime import hardware_acceptance as harness from ca2a_runtime.errors import AttestationUnsupported +from ca2a_runtime.policy import LocalPolicy +from ca2a_runtime.tee.software import SoftwareProvider def test_no_software_fallback(monkeypatch): @@ -27,7 +31,7 @@ def test_call_requires_hardware_and_labels_response_untrusted(tmp_path, monkeypa def send(*args, **kwargs): assert kwargs["require_hardware"] is True - assert kwargs["caller_provider"] is provider + assert kwargs["caller_provider"].inner is provider assert kwargs["verifier"] is verifier return {"accepted": True, "caller_attestation": "hardware"} @@ -46,7 +50,7 @@ def send(*args, **kwargs): "call", log, ) - entry = json.loads(log.path.read_text()) + entry = json.loads(log.path.read_text().splitlines()[-1]) assert entry["response_authenticated"] is False assert entry["peer_reported_caller_attestation"] == "hardware" assert "private-input" not in log.path.read_text() @@ -113,3 +117,83 @@ def test_cli_records_failed_preflight_without_hardware_claim(tmp_path, monkeypat assert entry["event"] == "run_failed" assert entry["error_type"] == "AttestationUnsupported" assert "hardware" not in entry + + +def test_real_http_quote_timeout_is_incomplete_work_not_security_rejection(tmp_path, monkeypatch): + """A stalled test provider demonstrates diagnostics, not the live SNP cause.""" + entered, release, finished = threading.Event(), threading.Event(), threading.Event() + + class BlockingSoftwareProvider(SoftwareProvider): + def attest(self, public_key, nonce): + entered.set() + assert release.wait(10), "test must release the blocked provider" + return super().attest(public_key, nonce) + + class FinishingNode(harness.ObservedNode): + def offer(self, nonce): + try: + return super().offer(nonce) + finally: + finished.set() + + receiver_log = harness.ReceiptLog(tmp_path / "receiver.jsonl", "software-timeout-test") + node = FinishingNode( + LocalPolicy.of({"read"}), + log=receiver_log, + provider=harness.ObservedSnpProvider(BlockingSoftwareProvider(), receiver_log), + ) + service = harness.server.serve(node, port=0) + # Sending a response after the test client disconnects can raise BrokenPipe. + service.handle_error = lambda *args: None + thread = threading.Thread(target=service.serve_forever, daemon=True) + thread.start() + monkeypatch.setattr(harness.client, "_TIMEOUT", 1.0) + monkeypatch.setattr(harness, "real_provider", SoftwareProvider) + monkeypatch.setattr(harness, "build_verifier", lambda *args: None) + (tmp_path / "chain.json").write_text("[]") + (tmp_path / "key").write_text(Ed25519PrivateKey.generate().private_bytes_raw().hex()) + (tmp_path / "payload").write_bytes(b"test-secret-not-for-logs") + sender_log = harness.ReceiptLog(tmp_path / "sender.jsonl", "software-timeout-test") + executor = ThreadPoolExecutor(max_workers=1) + try: + pending = executor.submit( + harness.run, + { + "peer": {}, + "chain": "chain.json", + "holder_key": "key", + "payload": "payload", + "peer_url": f"http://127.0.0.1:{service.server_port}", + "capability": "read", + "record_id": "timeout", + }, + tmp_path, + "call", + sender_log, + ) + assert entered.wait(5), "receiver must reach quote collection before inspecting the timeout" + with pytest.raises(TimeoutError): + pending.result(timeout=5) + assert pending.done(), "the HTTP call must time out, not the executor wait" + sender = [json.loads(line) for line in sender_log.path.read_text().splitlines()] + failure = next(item for item in sender if item["event"] == "operation_failed") + assert failure["stage"] == "sender_call" + assert failure["error_type"] == "TimeoutError" + assert failure["elapsed_seconds"] >= 0 + assert sender[-1]["receiver_outcome"] == "unknown" + assert sender[-1]["retried"] is False + receiver = [json.loads(line) for line in receiver_log.path.read_text().splitlines()] + assert [item["stage"] for item in receiver] == ["receiver_offer", "local_quote_collection"] + assert all(item["event"] == "operation_started" for item in receiver) + assert "test-secret" not in sender_log.path.read_text() + assert not any(item.get("accepted") for item in sender) + finally: + release.set() + executor.shutdown(wait=True) + assert finished.wait(5) + service.shutdown() + service.server_close() + thread.join(5) + receiver = [json.loads(line) for line in receiver_log.path.read_text().splitlines()] + assert receiver[-1]["event"] == "operation_completed" + assert receiver[-1]["stage"] == "receiver_offer" diff --git a/tests/unit/test_response_authentication.py b/tests/unit/test_response_authentication.py new file mode 100644 index 0000000..5170a82 --- /dev/null +++ b/tests/unit/test_response_authentication.py @@ -0,0 +1,452 @@ +"""Response trust is scoped to one live caller and one complete request.""" + +from __future__ import annotations + +import copy +import hashlib +import hmac +import json +import pickle +import threading +from concurrent.futures import ThreadPoolExecutor +from pathlib import Path + +import pytest +from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey + +from ca2a_runtime import response as response_module +from ca2a_runtime.attestation import VerifiedPeer +from ca2a_runtime.node import PeerNode +from ca2a_runtime.policy import LocalPolicy +from ca2a_runtime.response import ( + PROFILE, + AuthenticatedPeerError, + PendingResponse, + ResponseAuthenticationFailed, + ResponseProducer, + decode, + encode, +) +from ca2a_runtime.transport import client, server +from tests.unit.test_live_call import _chain + + +def exchange(*, clock=lambda: 10.0): + private = X25519PrivateKey.generate() + peer = VerifiedPeer(private.public_key().public_bytes_raw().hex(), "none", "test") + pending = PendingResponse(peer, {"record_id": "same-id", "payload": "test"}, clock=clock) + producer = ResponseProducer(private, decode(pending.request_bytes)) + return private, pending, producer + + +def success(producer): + return producer.finish(200, {"accepted": True, "record": {"id": "r"}}) + + +def test_valid_response_and_one_use(): + _, pending, producer = exchange() + raw = encode(success(producer)) + result = pending.verify(200, raw) + assert result.body["accepted"] + assert result.authentication()["peer_assurance"] == "none" + assert result.authentication()["audience"] == "live-caller-only" + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, raw) + + +@pytest.mark.parametrize( + "field,value", + [ + ("body", {"accepted": True, "record": {"id": "substituted"}}), + ("status", 403), + ("kind", "denial"), + ("profile", "legacy"), + ("request_sha256", "00" * 32), + ("mac", "00" * 32), + ], +) +def test_changed_response_is_not_exposed(field, value): + _, pending, producer = exchange() + response = success(producer) + response[field] = value + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, encode(response)) + + +@pytest.mark.parametrize("field", ["recipient", "session", "request"]) +def test_valid_mac_for_changed_request_context_is_refused(field): + private, pending, _ = exchange() + envelope = decode(pending.request_bytes) + if field == "recipient": + envelope[field] = X25519PrivateKey.generate().public_key().public_bytes_raw().hex() + elif field == "session": + envelope[field] = "33" * 32 + else: + envelope[field]["payload"] = "another task" + producer = ResponseProducer(private, envelope) + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, encode(success(producer))) + + +def test_wrong_peer_key_and_rotation_require_new_handshake(): + _, pending, _ = exchange() + other = X25519PrivateKey.generate() + envelope = decode(pending.request_bytes) + with pytest.raises(ResponseAuthenticationFailed): + ResponseProducer(other, envelope) + envelope["peer"] = other.public_key().public_bytes_raw().hex() + substituted = ResponseProducer(other, envelope) + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, encode(success(substituted))) + fresh = PendingResponse(VerifiedPeer(envelope["peer"], "none", "new"), {"task": "new"}) + producer = ResponseProducer(other, decode(fresh.request_bytes)) + assert fresh.verify(200, encode(success(producer))).body["accepted"] + + +def test_old_response_cannot_be_consumed_after_state_loss(): + private, pending, producer = exchange() + fresh = PendingResponse( + VerifiedPeer(private.public_key().public_bytes_raw().hex(), "none", "test"), + decode(pending.request_bytes)["request"], + ) + with pytest.raises(ResponseAuthenticationFailed): + fresh.verify(200, encode(success(producer))) + + +@pytest.mark.parametrize("now", [40.0, 41.0, 9.0, float("nan"), float("inf")]) +def test_deadline_clock_rollback_and_invalid_clock_fail_closed(now): + ticks = iter([10.0, now]) + _, pending, producer = exchange(clock=lambda: next(ticks)) + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, encode(success(producer))) + + +@pytest.mark.parametrize("timeout", [0, -1, 301, float("nan"), float("inf")]) +def test_invalid_timeout(timeout): + private = X25519PrivateKey.generate() + with pytest.raises(ValueError): + PendingResponse( + VerifiedPeer(private.public_key().public_bytes_raw().hex(), "none", ""), + {}, + timeout=timeout, + ) + + +def test_concurrent_consumption_has_exactly_one_success(): + _, pending, producer = exchange() + raw = encode(success(producer)) + + def consume(_): + try: + pending.verify(200, raw) + return True + except ResponseAuthenticationFailed: + return False + + with ThreadPoolExecutor(max_workers=8) as pool: + assert sum(pool.map(consume, range(16))) == 1 + + +@pytest.mark.parametrize( + "raw", + [ + b'{"a":1,"a":2}', + b'{"outer":{"a":1,"a":2}}', + b"[]", + b'{"n":NaN}', + b'{"n":1.0}', + b'{"n":9007199254740992}', + b'{"s":"\\ud800"}', + b'{"a":' + b"[" * 70 + b"0" + b"]" * 70 + b"}", + b"x" * ((1 << 20) + 1), + ], + ids=lambda raw: f"json-{len(raw)}-{hashlib.sha256(raw).hexdigest()[:8]}", +) +def test_ambiguous_or_unbounded_json_is_rejected(raw): + with pytest.raises(ResponseAuthenticationFailed): + decode(raw) + + +def test_invalid_response_consumes_occurrence(): + _, pending, producer = exchange() + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, b"{}") + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, encode(success(producer))) + + +@pytest.mark.parametrize("mutation", ["extra", "missing", "zero-key", "uppercase", "profile"]) +def test_invalid_request_context_is_rejected_before_execution(mutation): + private, pending, _ = exchange() + envelope = decode(pending.request_bytes) + if mutation == "extra": + envelope["extra"] = True + elif mutation == "missing": + del envelope["request"] + elif mutation == "zero-key": + envelope["recipient"] = "00" * 32 + elif mutation == "uppercase": + envelope["recipient"] = "AB" * 32 + else: + envelope["profile"] = "unknown" + with pytest.raises(ResponseAuthenticationFailed): + ResponseProducer(private, envelope) + + +def test_authenticated_denial_and_unsigned_denial_are_different(): + _, pending, producer = exchange() + body = {"error": {"code": "DENIED", "message": "denied"}} + response = pending.verify(403, encode(producer.finish(403, body))) + error = AuthenticatedPeerError(response) + assert error.code == "DENIED" + assert error.http_status == 403 + assert error.response_authentication["profile"] == PROFILE + _, other, _ = exchange() + with pytest.raises(ResponseAuthenticationFailed): + other.verify(403, encode(body)) + + +def test_valid_mac_with_unsupported_confidential_output_kind_is_rejected(): + _, pending, producer = exchange() + body = { + "profile": PROFILE, + "request_sha256": producer.request_sha256, + "status": 200, + "kind": "confidential-output", + "body": {"accepted": True}, + } + # Deliberately authenticate the unsupported kind: this isolates its gate. + mac = hmac.new( + producer._key, b"ca2a/response-mac/v1/statement\x00" + encode(body), "sha256" + ).hexdigest() + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, encode({**body, "mac": mac})) + + +def test_independent_hkdf_and_mac_verifies_reference_response(): + private, pending, producer = exchange() + envelope = decode(pending.request_bytes) + from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PublicKey + + shared = private.exchange( + X25519PublicKey.from_public_bytes(bytes.fromhex(envelope["recipient"])) + ) + # RFC 5869 extract and one-block expand, independent of production HKDF. + digest = hashlib.sha256(pending.request_bytes).digest() + prk = hmac.new(digest, shared, "sha256").digest() + key = hmac.new(prk, b"ca2a/response-mac/v1/key\x01", "sha256").digest() + response = success(producer) + mac = response.pop("mac") + # ASCII-only fixture: stdlib sorting supplies independent canonical bytes. + raw = json.dumps(response, sort_keys=True, separators=(",", ":")).encode() + expected = hmac.new(key, b"ca2a/response-mac/v1/statement\x00" + raw, "sha256").hexdigest() + assert mac == expected + + +@pytest.fixture +def live_peer(): + chain, holder = _chain() + node = PeerNode(LocalPolicy.of({"read"}), trusted_root_issuers={chain[0].issuer}) + http = server.serve(node, port=0) + thread = threading.Thread(target=http.serve_forever, daemon=True) + thread.start() + try: + yield node, f"http://127.0.0.1:{http.server_port}", chain, holder + finally: + http.shutdown() + http.server_close() + thread.join(timeout=5) + + +def call(peer, capability="read", strict=True): + _, url, chain, holder = peer + return client.send_task( + url, + chain, + capability, + "same-id", + holder_key=holder, + payload=b"private-canary", + require_authenticated_response=strict, + ) + + +def test_real_http_authenticated_success_denial_and_legacy(live_peer): + result = call(live_peer) + assert result["accepted"] + assert result["response_authentication"]["peer_assurance"] == "none" + assert "private-canary" not in json.dumps(result) + with pytest.raises(AuthenticatedPeerError) as failure: + call(live_peer, "write") + assert failure.value.code == "SCOPE_NOT_PERMITTED" + assert failure.value.response_authentication["profile"] == PROFILE + assert "response_authentication" not in call(live_peer, strict=False) + + +@pytest.mark.parametrize("change", ["body", "status", "unsigned", "other-request"]) +def test_real_http_substituted_response_never_reaches_caller(live_peer, monkeypatch, change): + original = client._post_authenticated + previous = None + + def intercept(url, raw): + nonlocal previous + status, response = original(url, raw) + envelope = decode(response) + if change == "body": + envelope["body"]["record"]["record_id"] = "forged" + elif change == "status": + status = 403 + elif change == "unsigned": + envelope = envelope["body"] + elif previous is None: + previous = response + return status, response + else: + return status, previous + return status, encode(envelope) + + monkeypatch.setattr(client, "_post_authenticated", intercept) + if change == "other-request": + assert call(live_peer)["accepted"] + with pytest.raises(ResponseAuthenticationFailed, match="unknown"): + call(live_peer) + + +def test_timeout_after_execution_is_not_retried(live_peer, monkeypatch): + original = client._post_authenticated + calls = [] + + def lose_reply(url, raw): + calls.append(raw) + status, _ = original(url, raw) + assert status == 200 + raise TimeoutError("reply lost") + + monkeypatch.setattr(client, "_post_authenticated", lose_reply) + with pytest.raises(ResponseAuthenticationFailed, match="unknown"): + call(live_peer) + assert len(calls) == 1 + + +def test_snapshotting_request_avoids_mutable_aliases(): + private, pending, _ = exchange() + envelope = decode(pending.request_bytes) + producer = ResponseProducer(private, envelope) + before = copy.deepcopy(producer.envelope) + envelope["request"]["payload"] = "changed" + assert producer.envelope == before + + +@pytest.mark.parametrize("clone", [copy.copy, copy.deepcopy, pickle.dumps]) +def test_pending_state_cannot_be_copied_or_restored(clone): + _, pending, _ = exchange() + with pytest.raises(TypeError, match="cannot be copied"): + clone(pending) + + +def test_deadline_is_checked_after_verification(monkeypatch): + now = [10.0] + _, pending, producer = exchange(clock=lambda: now[0]) + raw = encode(success(producer)) + original = response_module.decode + + def slow_decode(raw): + value = original(raw) + now[0] = 40.0 + return value + + monkeypatch.setattr(response_module, "decode", slow_decode) + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, raw) + + +def test_clock_failure_is_unknown_outcome(): + def clock(): + raise RuntimeError("clock unavailable") + + _, pending, producer = exchange() + pending._clock = clock + with pytest.raises(ResponseAuthenticationFailed, match="unknown"): + pending.verify(200, encode(success(producer))) + + +def test_legacy_response_cannot_inject_local_assurance(live_peer, monkeypatch): + original = client._post_json + + def spoof(url, body): + status, result = original(url, body) + result["response_authentication"] = {"profile": PROFILE, "peer_assurance": "hardware"} + return status, result + + monkeypatch.setattr(client, "_post_json", spoof) + assert "response_authentication" not in call(live_peer, strict=False) + + +def test_real_http_redirect_is_not_followed(live_peer, monkeypatch): + calls = [] + + def redirect(handler): + handler.rfile.read(int(handler.headers.get("Content-Length", 0))) + calls.append(handler.path) + handler.send_response(307) + handler.send_header("Location", server.AUTHENTICATED_TASK_PATH) + handler.send_header("Content-Length", "2") + handler.end_headers() + handler.wfile.write(b"{}") + + monkeypatch.setattr(server._PeerHandler, "do_POST", redirect) + with pytest.raises(ResponseAuthenticationFailed): + call(live_peer) + assert calls == [server.AUTHENTICATED_TASK_PATH] + + +def test_invalid_context_precedes_handler(live_peer, monkeypatch): + node, url, _, _ = live_peer + calls = [] + monkeypatch.setattr(node, "handle", lambda message: calls.append(message)) + pending = PendingResponse(VerifiedPeer(node.channel_public_key, "none", ""), {}) + envelope = decode(pending.request_bytes) + envelope["recipient"] = "00" * 32 + status, raw = client._post_authenticated(url + server.AUTHENTICATED_TASK_PATH, encode(envelope)) + assert calls == [] + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(status, raw) + + +def vector_pending(monkeypatch, vector): + key = X25519PrivateKey.from_private_bytes(bytes.fromhex(vector["caller_private_key"])) + + class FixedKey: + @staticmethod + def generate(): + return key + + monkeypatch.setattr(response_module, "X25519PrivateKey", FixedKey) + monkeypatch.setattr( + response_module.secrets, "token_hex", lambda _: vector["request"]["session"] + ) + return PendingResponse( + VerifiedPeer(vector["request"]["peer"], "none", "synthetic"), vector["request"]["request"] + ) + + +def test_committed_language_neutral_vector_and_negative_twins(monkeypatch): + vector = json.loads((Path(__file__).parents[1] / "fixtures/response-mac-v1.json").read_bytes()) + pending = vector_pending(monkeypatch, vector) + assert pending.request_bytes.hex() == vector["request_canonical_hex"] + assert pending._key.hex() == vector["hkdf_key_hex"] + private = X25519PrivateKey.from_private_bytes(bytes.fromhex(vector["callee_private_key"])) + producer = ResponseProducer(private, vector["request"]) + assert producer.finish(200, vector["response"]["body"]) == vector["response"] + assert pending.verify(200, encode(vector["response"])).body["accepted"] + for case in vector["negative_cases"]: + valid = vector_pending(monkeypatch, vector) + assert valid.verify(200, encode(vector["response"])).body["accepted"] + pending = vector_pending(monkeypatch, vector) + changed = copy.deepcopy(vector["response"]) + current = changed + for field in case["path"][:-1]: + current = current[field] + current[case["path"][-1]] = case["value"] + with pytest.raises(ResponseAuthenticationFailed): + pending.verify(200, encode(changed))