consumeInbound<P, R> function

Future<ConsumeOutcome<R>> consumeInbound<P, R>({
  1. required TransportHandler transport,
  2. required SpecPolicy spec,
  3. required ProofPolicy proofPolicy,
  4. required PayloadPolicy payloadPolicy,
  5. required ConsumeChecks checks,
  6. required TrustTaskDocument<P> doc,
  7. required String myVid,
  8. required DateTime now,
  9. required String newErrorId(),
  10. required Object? payloadToJson(
    1. P payload
    ),
  11. required Handler<P, R> handler,
  12. Clock? clock,
})

Run SPEC §7.2 items 4–8 against doc, then either call handler or build the routed error response per §8.1.

Implementation

Future<ConsumeOutcome<R>> consumeInbound<P, R>({
  required TransportHandler transport,
  required SpecPolicy spec,
  required ProofPolicy proofPolicy,
  required PayloadPolicy payloadPolicy,
  required ConsumeChecks checks,
  required TrustTaskDocument<P> doc,
  required String myVid,
  required DateTime now,
  required String Function() newErrorId,
  required Object? Function(P payload) payloadToJson,
  required Handler<P, R> handler,
  Clock? clock,
}) async {
  ConsumeOutcome<R> route(RejectReason reason) {
    final error = reject(transport, doc, newErrorId(), reason, clock: clock);
    return error == null ? Suppressed<R>(reason) : Rejected<R>(error);
  }

  // §7.2 item 2 — payload schema. Runs first, in the spec's own order, and
  // before anything that reasons about what the payload means: a payload that is
  // not the shape the specification declares should be refused as malformed
  // rather than interpreted.
  final schema = spec.payloadSchema;
  if (payloadPolicy is ValidatePayload && schema != null) {
    String? failure;
    try {
      failure =
          payloadPolicy.validator.validate(schema, payloadToJson(doc.payload));
    } on Object catch (e) {
      // A validator that throws has not accepted the document. Treating an
      // exception as a pass would make a broken validator indistinguishable from
      // a passing one — the failure mode this policy exists to remove.
      failure = e.toString();
    }
    if (failure != null) {
      final detail = failure.trim();
      return route(
        RejectReason(
          code: StandardCode.malformedRequest,
          message: detail.isEmpty
              ? 'payload does not conform to its schema (SPEC §7.2 item 2)'
              : 'payload does not conform to its schema (SPEC §7.2 item 2): $detail',
        ),
      );
    }
  }

  // §7.2 items 4 + 5a — expiry and wrong-recipient.
  final basic = validateBasic(doc, now, myVid);
  if (basic != null) return route(basic);

  // §7.2 item 4, the other half — the freshness bound over `issuedAt`.
  // `validateBasic` honours `expiresAt`, which is optional and which a producer
  // sets for its own reasons; on its own it leaves a document stamped years ago,
  // or years hence, indefinitely acceptable. It is also what bounds the replay
  // record below: §7.2 makes the acceptance window and the record's retention
  // one bound.
  final fresh = validateFreshness(doc, now, checks.freshness);
  if (fresh != null) return route(fresh);

  // §7.2 item 6 — in-band vs transport-derived identity cross-check.
  final resolution = resolveParties(transport, doc);
  final mismatch = resolution.error;
  if (mismatch != null) return route(identityMismatchReason(mismatch));
  final parties = resolution.parties!;

  // §7.2 item 7 clause B — the consumer's chosen proof policy.
  if (doc.proof != null) {
    switch (proofPolicy) {
      case VerifyProof(:final verifier):
        var ok = false;
        try {
          ok = await verifier.verify(doc.toJson(payloadToJson));
        } on Object {
          ok = false;
        }
        if (!ok) {
          // A constant, never the verifier's own error text. SPEC §12.4 extends
          // the §8.1 identity rule to every code, and a verifier's vocabulary
          // names DIDs it tried to resolve, whether a resolver answered, and
          // what a fetched DID document contained — a resolver-reachability
          // oracle for a sender who is, by construction, unauthenticated. Log
          // the detail; do not send it.
          return route(
            const RejectReason(
              code: StandardCode.proofInvalid,
              message: proofInvalidWireMessage,
            ),
          );
        }
      case RejectProofIfPresent():
        return route(
          const RejectReason(
            code: StandardCode.malformedRequest,
            message: proofNotAcceptedByPolicy,
          ),
        );
      case AcceptProofUnverified():
        break;
    }
  }

  // §7.2 items 5b + 7 clause A + 8 — the policy-driven checks, in one place so
  // this pipeline and any binding-specific one cannot diverge on the check set.
  final policy = enforceSpecPolicy(doc, spec);
  if (policy != null) return route(policy);

  // §7.2 item 11 — the duplicate-execution record. Deliberately **last**:
  // claiming the `id` marks the document as accepted for execution, and a
  // document some earlier check refuses was never accepted. Claiming first would
  // burn the `id` on every malformed or unauthorised arrival, so a corrected
  // resend under the same `id` would come back `idConflict` forever — and an
  // attacker could pre-burn an `id` it had merely observed.
  ReplayGuard? claimedGuard;
  String? claimedDigest;
  final replay = checks.replay;
  if (replay is GuardedReplay) {
    final guard = replay.guard;
    final digest = documentDigest(doc, payloadToJson);

    // §7.2 (*Bounding the record*): "A consumer that can establish neither an
    // `expiresAt` nor an age for a document has no window in which to place it,
    // and MUST NOT execute a consequential Trust Task on it." A guard asked to
    // retain a record forever is not a guard, so refuse rather than pretend.
    final retainUntil = recordExpiry(doc, checks.freshness, now);
    if (retainUntil == null) {
      return route(
        const RejectReason(
          code: StandardCode.expired,
          message: staleWireMessage,
        ),
      );
    }

    ReplayVerdict verdict;
    try {
      verdict = await guard.claim(doc.id, digest, retainUntil, now);
    } on Object {
      // Fail closed. A consumer that cannot consult its record has not satisfied
      // item 11, and executing anyway is exactly the double execution the rule
      // forbids. `unavailable` is retryable, which is the honest answer: the
      // producer's bit-for-bit resend will be absorbed correctly once the store
      // is back. The thrown detail — a hostname, a connection string — stays out
      // of the message, per §10.4.
      return route(
        const RejectReason(
          code: StandardCode.unavailable,
          message: replayRecordUnavailable,
          retryable: true,
        ),
      );
    }

    switch (verdict) {
      case Duplicate(:final priorResponse, :final inFlight):
        return DuplicateOutcome<R>(
          priorResponse: priorResponse,
          inFlight: inFlight,
        );
      case Conflict():
        return route(
          const RejectReason(
            code: StandardCode.idConflict,
            message: idConflictWireMessage,
          ),
        );
      case Fresh():
        claimedGuard = guard;
        claimedDigest = digest;
    }
  }

  /// Record the response for a completed execution, best-effort.
  ///
  /// The effect has already happened, so a guard that cannot cache the response
  /// cannot un-happen it. The record of the *claim* is what item 11 needs and it
  /// is already written; all that is lost is the ability to hand the same
  /// response back, and the duplicate is still absorbed.
  Future<void> settle(Object? response) async {
    if (claimedGuard == null) return;
    try {
      await claimedGuard.recordResponse(doc.id, response);
    } on Object {
      /* see above */
    }
  }

  TrustTaskDocument<R>? response;
  try {
    response = await handler(doc, parties);
  } on Refusal catch (refusal) {
    // §8.4: a retryable refusal has just invited the producer to re-send this
    // document bit-for-bit. Holding the claim would answer that invited retry
    // with the cached failure forever. A non-retryable refusal is final, so the
    // record stands and a replay is answered with the same determination.
    if (refusal.response.payload.retryable) {
      if (claimedGuard != null) {
        try {
          await claimedGuard.release(doc.id, claimedDigest!);
        } on Object {
          /* best-effort */
        }
      }
    } else {
      await settle(refusal.response.toJson((p) => p.toJson()));
    }
    return Rejected<R>(refusal.response);
  }

  if (response == null) {
    // Fire-and-forget: nothing to cache, but the claim stands — the effect
    // happened, and item 11 is about the effect, not about the response.
    await settle(null);
    return Accepted<R>();
  }

  await settle(response);
  return Handled<R>(response);
}