consumeInbound<P, R> function
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 newErrorId(),
- required Object? payloadToJson(
- P payload
- required Handler<
P, R> handler, - 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);
}