Skip to main content

macula_rust/station_link/
serve.rs

1//! Serving, as macula 12's provider does it. A procedure is served under an
2//! org namespace: the realm's org directory names the org's key, the org's
3//! procedure delegation names this node, both are found in the DHT, and the
4//! signed procedure_advertisement carries them to the station in an
5//! ADVERTISE, checked against the realm key before it goes out. Or it is
6//! served in this node's own namespace, `~<node_id>/<name>` (D25 item 6): its
7//! advertisement carries no authorization, its signature alone authorizes it,
8//! and no realm key or record in the DHT is needed. The station routes CALLs
9//! for the procedure to this link; each is admitted once per (caller,
10//! request_id) and answered with a RESULT or ERROR signed by this node.
11//! UNADVERTISE carries a tombstone of the advertisement.
12//!
13//! With `kem_advertise` on, a confidential procedure's advertisement names
14//! this node's KEM key, and a request sealed to it is opened and answered
15//! sealed (confidential.rs).
16//!
17//! A procedure is served open, or gated on a UCAN policy: a request whose
18//! token the policy does not authorize is refused unauthorized
19//! (authorize.rs).
20
21use std::future::Future;
22use std::pin::Pin;
23use std::sync::{Arc, Mutex, Weak};
24use std::time::Duration;
25
26use tokio::sync::watch;
27
28use crate::cbor::{self, Value};
29use crate::frame::{self, StreamMode, VerifiedRequest};
30use crate::node_key::NodeKey;
31use crate::record::{
32    self, Authorization, ProcedureAdvertisementOptions, Reason, Record, RecordError,
33    TombstoneOptions, Trust,
34};
35use crate::seal;
36use crate::ucan::Policy;
37
38use super::admission::Verdict;
39use super::authorize::authorize;
40use super::confidential::{
41    clear_allowed, opened_request, CallSeal, Confidentiality, CODE_SEALED_REFUSED,
42    CODE_SEALED_REQUIRED,
43};
44use super::framing::MAX_FRAME_BYTES;
45use super::stream::StreamHandler;
46use super::{now_ms, Inner, Link, LinkError};
47
48/// The provider codes of a served procedure's ERRORs: a handler's own
49/// refusal, a handler that panicked, a procedure this link does not serve, a
50/// copy of a request still running, and a result the wire cannot carry.
51const CODE_HANDLER_ERROR: &str = "handler_error";
52const CODE_HANDLER_CRASHED: &str = "temporary_relay_failure";
53const CODE_UNKNOWN_PROCEDURE: &str = "unknown_next_peer";
54pub(super) const CODE_REQUEST_COPY: &str = "request_copy";
55const CODE_PAYLOAD_TOO_LARGE: &str = "payload_too_large";
56const CODE_UNSENDABLE: &str = "unknown_error";
57/// A sealed answer this node could not seal: a refusal from the closed set,
58/// which carries no application data.
59const CODE_UNAVAILABLE: &str = "unavailable";
60
61/// An ERROR's detail is at most 256 bytes, cut on a character boundary.
62const MAX_DETAIL_BYTES: usize = 256;
63
64/// macula's default and longest advertisement lifetime.
65const MAX_ADVERTISEMENT_TTL: Duration =
66    Duration::from_millis(super::confidential::MAX_ADVERTISEMENT_TTL_MS as u64);
67/// How soon a failed renewal is tried again, while the advertisement it
68/// replaces still lives.
69const REFRESH_RETRY: Duration = Duration::from_secs(10);
70
71/// A future a handler returns.
72pub type BoxFuture<T> = Pin<Box<dyn Future<Output = T> + Send + 'static>>;
73
74/// Answers a [`Request`] with a result payload, or an error whose text the
75/// caller receives as a handler_error's detail. A handler still running at
76/// the request's deadline is dropped and answered handler_error.
77pub type Handler = Arc<dyn Fn(Request) -> BoxFuture<Result<Value, String>> + Send + Sync>;
78
79/// A [`Handler`] from an async closure.
80pub fn handler<F, Fut>(f: F) -> Handler
81where
82    F: Fn(Request) -> Fut + Send + Sync + 'static,
83    Fut: Future<Output = Result<Value, String>> + Send + 'static,
84{
85    Arc::new(move |r| Box::pin(f(r)))
86}
87
88/// A CALL a served procedure answers: the caller (the key id its signature
89/// verified under), what it asked for, and its deadline in unix
90/// milliseconds.
91#[derive(Debug, Clone, PartialEq)]
92pub struct Request {
93    pub caller: [u8; 32],
94    pub realm: [u8; 32],
95    pub procedure: String,
96    pub payload: Value,
97    pub token: Option<Vec<u8>>,
98    pub proofs: Option<Vec<Vec<u8>>>,
99    pub deadline_ms: u64,
100    /// Whether the request came sealed end to end: `payload` is its opened
101    /// plaintext, and the answer goes back sealed.
102    pub sealed: bool,
103}
104
105/// A procedure to serve: its realm and name, exactly one of a unary handler
106/// and a stream offer, and, for an org procedure, the realm key the org
107/// directory must be signed with, as the realm's members pin it; a procedure
108/// in this node's own namespace needs none.
109#[derive(Clone)]
110pub struct Offer {
111    pub realm: [u8; 32],
112    pub procedure: String,
113    pub handler: Option<Handler>,
114    pub stream: Option<StreamOffer>,
115    pub realm_key: Option<Vec<u8>>,
116    /// How the procedure takes its requests: whether its advertisement names
117    /// this node's KEM key, with `kem_advertise` on, and whether it takes
118    /// clear ones.
119    pub confidential: Confidentiality,
120    /// When the procedure was first advertised naming the key (unix ms),
121    /// from which its keyless window runs; serving sets it to now when
122    /// `None`. A pool keeps it across the links it serves on, so no link
123    /// reopens the window.
124    pub keyed_since_ms: Option<i64>,
125    /// What the procedure requires of a request's UCAN: `None` serves any
126    /// caller (an open procedure ignores any token), and a [`Policy`] serves
127    /// only a request whose token it authorizes (authorize.rs).
128    pub policy: Option<Policy>,
129}
130
131impl Offer {
132    /// A unary procedure, with no realm key.
133    pub fn unary(realm: [u8; 32], procedure: &str, handler: Handler) -> Offer {
134        Offer {
135            realm,
136            procedure: procedure.to_string(),
137            handler: Some(handler),
138            stream: None,
139            realm_key: None,
140            confidential: Confidentiality::Preferred,
141            keyed_since_ms: None,
142            policy: None,
143        }
144    }
145
146    /// A streaming procedure of `mode`, with no realm key.
147    pub fn stream(
148        realm: [u8; 32],
149        procedure: &str,
150        mode: StreamMode,
151        handler: StreamHandler,
152    ) -> Offer {
153        Offer {
154            realm,
155            procedure: procedure.to_string(),
156            handler: None,
157            stream: Some(StreamOffer { mode, handler }),
158            realm_key: None,
159            confidential: Confidentiality::Preferred,
160            keyed_since_ms: None,
161            policy: None,
162        }
163    }
164}
165
166/// A streaming procedure's mode and handler. The advertisement is the one a
167/// unary procedure sends, which names no mode: a STREAM_OPEN of another mode
168/// is refused mode_mismatch.
169#[derive(Clone)]
170pub struct StreamOffer {
171    pub mode: StreamMode,
172    pub handler: StreamHandler,
173}
174
175/// A procedure a link serves, as the link holds it.
176pub(super) type ServedEntry = Arc<ServedInner>;
177
178pub(super) struct ServedInner {
179    link: Weak<Inner>,
180    key: ([u8; 32], String),
181    pub(super) offer: Offer,
182    latest: Mutex<Record>,
183    err: Mutex<Option<LinkError>>,
184    done_tx: watch::Sender<bool>,
185}
186
187/// A procedure this link serves, until [`Served::stop`] or the link ends, or
188/// until its advertisement lapses because its authorization could not be
189/// found again. Cloning it shares the serving.
190#[derive(Clone)]
191pub struct Served {
192    inner: Arc<ServedInner>,
193}
194
195impl Link {
196    /// Advertises `o`'s procedure on the link and answers its CALLs until
197    /// stopped. For an org procedure it resolves the org directory and this
198    /// node's procedure delegation from the DHT; a procedure in this node's
199    /// own namespace needs neither, and another node's namespace is refused.
200    /// It signs the advertisement with this node's identity key, naming the
201    /// connected station as the serving station and living no longer than
202    /// the records it carries nor 5 minutes, and checks its authorization
203    /// (against the realm key for an org procedure) before sending it in an
204    /// ADVERTISE and putting it in the DHT. The advertisement is renewed at
205    /// half its lifetime.
206    pub async fn serve(&self, mut o: Offer) -> Result<Served, LinkError> {
207        let own = record::in_own_namespace(&o.procedure);
208        if o.handler.is_some() == o.stream.is_some()
209            || (!own && o.realm_key.is_none())
210            || o.policy.as_ref().is_some_and(|p| !p.valid())
211        {
212            return Err(LinkError::InvalidOffer);
213        }
214        if !matches!(record::procedure_org(&o.procedure), Ok(Some(_))) {
215            return Err(LinkError::NoOrg);
216        }
217        let can_name_key = self.inner.kem_advertise && self.inner.keyring.is_some();
218        if o.confidential == Confidentiality::Required && !can_name_key {
219            return Err(LinkError::KemAdvertiseDisabled);
220        }
221        if self.inner.keyed(&o) && o.keyed_since_ms.is_none() {
222            o.keyed_since_ms = Some(now_ms());
223        }
224        let (advertisement, wire) = self.advertisement(&o, MAX_ADVERTISEMENT_TTL).await?;
225        let key = (o.realm, o.procedure.clone());
226        let (done_tx, _) = watch::channel(false);
227        let served = Arc::new(ServedInner {
228            link: Arc::downgrade(&self.inner),
229            key: key.clone(),
230            offer: o,
231            latest: Mutex::new(advertisement),
232            err: Mutex::new(None),
233            done_tx,
234        });
235        route_calls(&self.inner, key, served.clone())?;
236        if let Err(e) = self.announce(&wire).await {
237            served.end(e.clone());
238            return Err(e);
239        }
240        tokio::spawn(renew(served.clone()));
241        Ok(Served { inner: served })
242    }
243
244    /// `o`'s signed procedure_advertisement and its wire form, living at most
245    /// `max_ttl`: with the authorization resolved from the DHT for an org
246    /// procedure, with none in this node's own namespace.
247    async fn advertisement(
248        &self,
249        o: &Offer,
250        max_ttl: Duration,
251    ) -> Result<(Record, Vec<u8>), LinkError> {
252        let inner = &self.inner;
253        let max_ttl_ms = max_ttl.as_millis() as u64;
254        // Each signing reads the keyring, so a renewal names a rotated key.
255        let kem_key = inner.kem_key(o);
256        let opts = if record::in_own_namespace(&o.procedure) {
257            ProcedureAdvertisementOptions {
258                authorization: Authorization::None,
259                kem_key,
260                ttl_ms: max_ttl_ms,
261            }
262        } else {
263            self.org_advertisement_options(o, max_ttl_ms, kem_key)
264                .await?
265        };
266        let unsigned = record::new_procedure_advertisement(
267            &inner.self_id,
268            &o.realm,
269            &o.procedure,
270            &inner.station.node_id,
271            &opts,
272        )?;
273        // Signed, then checked as a caller will check it before it is sent
274        // anywhere.
275        let signed = record::sign(&unsigned, &inner.key)?;
276        let wire = record::encode(&signed)?;
277        let now = now_ms();
278        let verified = record::verify(&wire, inner.profile, now)?;
279        record::verify_authorization(
280            &verified,
281            &Trust {
282                profile: inner.profile,
283                realm_key: o.realm_key.clone(),
284            },
285            now,
286        )?;
287        Ok((signed, wire))
288    }
289
290    /// The options of an org procedure's advertisement: its org directory and
291    /// this node's procedure delegation, resolved from the DHT, and a lifetime
292    /// no longer than either of them nor `max_ttl_ms`.
293    async fn org_advertisement_options(
294        &self,
295        o: &Offer,
296        max_ttl_ms: u64,
297        kem_key: Option<Vec<u8>>,
298    ) -> Result<ProcedureAdvertisementOptions, LinkError> {
299        let inner = &self.inner;
300        let org = record::procedure_org(&o.procedure)?.ok_or(LinkError::NoOrg)?;
301        let directory = self
302            .find_record(&record::org_directory_key(&o.realm, org))
303            .await?;
304        let named = record::read_org_directory(directory.record())?;
305        let delegation = self
306            .find_record(&record::procedure_delegation_key(
307                &named.org_key,
308                &inner.self_id,
309            ))
310            .await?;
311        let now = now_ms();
312        let ttl = (max_ttl_ms as i64)
313            .min(directory.record().expires_at as i64 - now)
314            .min(delegation.record().expires_at as i64 - now);
315        if ttl <= 0 {
316            return Err(RecordError::AuthorizationOutlived.into());
317        }
318        Ok(ProcedureAdvertisementOptions {
319            authorization: Authorization::Delegation {
320                org_directory: record::encode(directory.record())?,
321                procedure_delegation: record::encode(delegation.record())?,
322            },
323            kem_key,
324            ttl_ms: ttl as u64,
325        })
326    }
327
328    /// Sends an advertisement to the station in an ADVERTISE, which routes
329    /// CALLs through it, and puts it in the DHT, where a caller resolving
330    /// the procedure finds it, as macula's advertise_direct does both.
331    async fn announce(&self, wire: &[u8]) -> Result<(), LinkError> {
332        self.inner
333            .send_control(&frame::advertise_frame(wire))
334            .await?;
335        self.put_record(wire).await
336    }
337}
338
339/// Routes the procedure's CALLs on the link to `served`, unless the link has
340/// ended or already serves the procedure.
341fn route_calls(
342    inner: &Inner,
343    key: ([u8; 32], String),
344    served: Arc<ServedInner>,
345) -> Result<(), LinkError> {
346    let mut state = inner.lock();
347    if let Some(e) = &state.ended {
348        return Err(e.clone());
349    }
350    if state.served.contains_key(&key) {
351        return Err(LinkError::AlreadyServed);
352    }
353    state.served.insert(key, served);
354    Ok(())
355}
356
357impl Served {
358    /// Withdraws the advertisement with an UNADVERTISE carrying its
359    /// tombstone, signed by this node, puts the tombstone in the
360    /// advertisement's DHT slot, and stops answering the procedure's CALLs.
361    /// Stopping a procedure no longer served does nothing.
362    pub async fn stop(&self) -> Result<(), LinkError> {
363        if self.inner.is_done() {
364            return Ok(());
365        }
366        self.inner.end(LinkError::Stopped);
367        let Some(inner) = self.inner.link.upgrade() else {
368            return Ok(());
369        };
370        let latest = self
371            .inner
372            .latest
373            .lock()
374            .unwrap_or_else(|p| p.into_inner())
375            .clone();
376        let tombstone =
377            record::new_tombstone(&latest, Reason::Shutdown, &TombstoneOptions::default())?;
378        let wire = record::encode(&record::sign(&tombstone, &inner.key)?)?;
379        inner.send_control(&frame::unadvertise_frame(&wire)).await?;
380        Link { inner }.put_record(&wire).await
381    }
382
383    /// Waits until the procedure is no longer served, and says why.
384    pub async fn done(&self) -> LinkError {
385        let mut done = self.inner.done_tx.subscribe();
386        let _ = done.wait_for(|ended| *ended).await;
387        self.error().unwrap_or(LinkError::Stopped)
388    }
389
390    /// Why the procedure is no longer served: [`LinkError::Stopped`] after
391    /// stop, the link's error when it ended, or the renewal's failure; `None`
392    /// while served.
393    pub fn error(&self) -> Option<LinkError> {
394        self.inner
395            .err
396            .lock()
397            .unwrap_or_else(|p| p.into_inner())
398            .clone()
399    }
400}
401
402impl ServedInner {
403    fn is_done(&self) -> bool {
404        *self.done_tx.borrow()
405    }
406
407    /// Ends the serving once, with `err`: the link no longer routes the
408    /// procedure's CALLs to its handler.
409    pub(super) fn end(self: &Arc<Self>, err: LinkError) {
410        if !self.hold_err(err) {
411            return;
412        }
413        self.unroute();
414        let _ = self.done_tx.send_replace(true);
415    }
416
417    /// Keeps `err` as why the serving ended, and whether it is the first;
418    /// a later one is not kept.
419    fn hold_err(&self, err: LinkError) -> bool {
420        let mut held = self.err.lock().unwrap_or_else(|p| p.into_inner());
421        if held.is_some() {
422            return false;
423        }
424        *held = Some(err);
425        true
426    }
427
428    /// Stops the link routing the procedure's CALLs here, when it still
429    /// routes them to this serving and not a later one.
430    fn unroute(self: &Arc<Self>) {
431        let Some(inner) = self.link.upgrade() else {
432            return;
433        };
434        let mut state = inner.lock();
435        if state
436            .served
437            .get(&self.key)
438            .is_some_and(|s| Arc::ptr_eq(s, self))
439        {
440            state.served.remove(&self.key);
441        }
442    }
443}
444
445/// Sends a fresh advertisement at half the current one's lifetime. A renewal
446/// that fails is tried again every [`REFRESH_RETRY`] while the current one
447/// lives; when it lapses unrenewed, the procedure is no longer served and its
448/// error says why.
449async fn renew(served: Arc<ServedInner>) {
450    let mut current = served
451        .latest
452        .lock()
453        .unwrap_or_else(|p| p.into_inner())
454        .clone();
455    let mut wait = half_life(&current);
456    let mut last_err: Option<LinkError> = None;
457    let mut stopped = served.done_tx.subscribe();
458    loop {
459        let Some(mut link_done) = served.link.upgrade().map(|l| l.done_rx.clone()) else {
460            return;
461        };
462        tokio::select! {
463            _ = stopped.wait_for(|ended| *ended) => return,
464            _ = link_done.wait_for(|ended| *ended) => {
465                let err = served.link.upgrade().and_then(|l| l.lock().ended.clone()).unwrap_or(LinkError::Closed);
466                served.end(err);
467                return;
468            }
469            _ = tokio::time::sleep(wait) => {}
470        }
471        if now_ms() >= current.expires_at as i64 {
472            served.end(last_err.unwrap_or(LinkError::Stopped));
473            return;
474        }
475        let Some(inner) = served.link.upgrade() else {
476            return;
477        };
478        let link = Link { inner };
479        let renewed = async {
480            let (advertisement, wire) = link
481                .advertisement(&served.offer, MAX_ADVERTISEMENT_TTL)
482                .await?;
483            link.announce(&wire).await?;
484            Ok::<_, LinkError>(advertisement)
485        }
486        .await;
487        match renewed {
488            Ok(advertisement) => {
489                *served.latest.lock().unwrap_or_else(|p| p.into_inner()) = advertisement.clone();
490                wait = half_life(&advertisement);
491                current = advertisement;
492            }
493            Err(e) => {
494                last_err = Some(e);
495                let left = (current.expires_at as i64 - now_ms()).max(0) as u64;
496                wait = REFRESH_RETRY.min(Duration::from_millis(left));
497            }
498        }
499    }
500}
501
502fn half_life(r: &Record) -> Duration {
503    Duration::from_millis(r.expires_at.saturating_sub(r.created_at) / 2)
504}
505
506/// Answers a CALL the station routed to this link. A request that does not
507/// verify, or targets another node, gets no reply and is counted. One that
508/// verifies is judged by the admission, then answered by its handler, or by
509/// an ERROR naming why not.
510pub(super) fn called(inner: &Arc<Inner>, v: &Value) {
511    let Ok(request) = frame::verify_request(v, inner.profile) else {
512        inner.count("unverified_call");
513        return;
514    };
515    if request.target != inner.self_id {
516        inner.count("call_for_another_node");
517        return;
518    }
519    let inner = inner.clone();
520    match inner.admission.admit(&request, &inner.share, now_ms()) {
521        Verdict::Refused(code) => refuse_call(inner, &request, code),
522        Verdict::Copy(None) => refuse_call(inner, &request, CODE_REQUEST_COPY),
523        Verdict::Copy(Some(stored)) => {
524            tokio::spawn(async move {
525                let _ = inner.control.write(&stored, MAX_FRAME_BYTES).await;
526            });
527        }
528        Verdict::New => {
529            tokio::spawn(answer(inner, request));
530        }
531    }
532}
533
534/// Sends this node's ERROR for `request` with `code`, or counts it
535/// unsignable_reply when it does not sign.
536fn refuse_call(inner: Arc<Inner>, request: &VerifiedRequest, code: &'static str) {
537    let Some(reply) = provider_error(&inner.signer(), request, code, None) else {
538        inner.count("unsignable_reply");
539        return;
540    };
541    tokio::spawn(async move { send_reply(&inner, reply).await });
542}
543
544/// Serves a request and sends its signed reply, storing it for the
545/// request's copies.
546async fn answer(inner: Arc<Inner>, request: VerifiedRequest) {
547    let offer = inner
548        .lock()
549        .served
550        .get(&(request.realm, request.procedure.clone()))
551        .map(|s| s.offer.clone());
552    let Some(reply) = reply(&inner, &request, offer).await else {
553        inner.count("unsignable_reply");
554        return;
555    };
556    let Ok(encoded) = cbor::encode(&reply) else {
557        inner.count("unencodable_reply");
558        return;
559    };
560    inner.admission.store(&request, encoded.clone());
561    let _ = inner.control.write(&encoded, MAX_FRAME_BYTES).await;
562}
563
564/// The signed reply to `request`, served by `offer` (`None` when this link
565/// serves no such procedure), in macula 13's order: a sealed request opened
566/// first, or refused sealed_refused in the clear when it does not open; a
567/// clear one to a procedure past its keyless window refused sealed_required;
568/// then the procedure, its policy and its handler, answered sealed when the
569/// request was.
570async fn reply(inner: &Inner, request: &VerifiedRequest, offer: Option<Offer>) -> Option<Value> {
571    let (payload, sealing) = match &request.sealed {
572        Some(_) => match opened_request(inner.keyring.as_deref(), request) {
573            Ok((payload, sealing)) => (payload, Some(sealing)),
574            Err(detail) => {
575                return provider_error(&inner.signer(), request, CODE_SEALED_REFUSED, Some(&detail))
576            }
577        },
578        None => {
579            let refused = offer
580                .as_ref()
581                .is_some_and(|o| !clear_allowed(o.confidential, inner.keyed_since(o), now_ms()));
582            if refused {
583                return provider_error(&inner.signer(), request, CODE_SEALED_REQUIRED, None);
584            }
585            (request.payload.clone(), None)
586        }
587    };
588    let outcome = match gated(inner.profile, offer, request) {
589        Ok(handler) => handled(handler, request, without_caller(payload), sealing.is_some()).await,
590        Err(code) => Outcome::Refused(code, None),
591    };
592    answered(&inner.signer(), request, sealing.as_ref(), outcome)
593}
594
595/// The handler that serves `request` under `offer`'s policy, or the code it
596/// is refused with: unknown_procedure when this link serves no such unary
597/// procedure, and the policy's refusal when its token does not authorize it.
598fn gated(
599    profile: crate::profile::Profile,
600    offer: Option<Offer>,
601    request: &VerifiedRequest,
602) -> Result<Handler, &'static str> {
603    let Some(Offer {
604        handler: Some(handler),
605        policy,
606        ..
607    }) = offer
608    else {
609        return Err(CODE_UNKNOWN_PROCEDURE);
610    };
611    authorize(profile, policy.as_ref(), request).map_or(Ok(handler), Err)
612}
613
614/// `payload` as a handler receives it: a map loses a text "caller" key, so
615/// the only caller a handler can learn is the verified one
616/// ([`Request::caller`]), as macula's `with_caller/2` and macula-go's
617/// `withoutCaller` arrange (macula-rust#13). Any other payload is untouched.
618pub(super) fn without_caller(payload: Value) -> Value {
619    match payload {
620        Value::Map(_) => payload.without(&["caller"]),
621        other => other,
622    }
623}
624
625/// What answers a request: a RESULT's payload, or an ERROR's code and
626/// detail.
627enum Outcome {
628    Result(Value),
629    Refused(&'static str, Option<String>),
630}
631
632/// The handler's answer to `request`, on `payload`: its result, its refusal
633/// as handler_error, a panic as temporary_relay_failure, or handler_error
634/// when it is still running at the deadline.
635async fn handled(
636    handler: Handler,
637    request: &VerifiedRequest,
638    payload: Value,
639    sealed: bool,
640) -> Outcome {
641    let running = tokio::spawn(handler(Request {
642        caller: request.caller,
643        realm: request.realm,
644        procedure: request.procedure.clone(),
645        payload,
646        token: request.token.clone(),
647        proofs: request.proofs.clone(),
648        deadline_ms: request.deadline,
649        sealed,
650    }));
651    let abort = running.abort_handle();
652    let left = (request.deadline as i64 - now_ms()).max(0) as u64;
653    match tokio::time::timeout(Duration::from_millis(left), running).await {
654        Err(_) => {
655            abort.abort();
656            Outcome::Refused(
657                CODE_HANDLER_ERROR,
658                Some("the request's deadline passed".into()),
659            )
660        }
661        Ok(Err(_panicked)) => Outcome::Refused(CODE_HANDLER_CRASHED, None),
662        Ok(Ok(Err(refusal))) => Outcome::Refused(
663            CODE_HANDLER_ERROR,
664            Some(bounded_detail(&refusal).to_string()),
665        ),
666        Ok(Ok(Ok(payload))) => Outcome::Result(payload),
667    }
668}
669
670/// What signs this node's replies: its identity key and its node id, the
671/// replies' responded_by.
672pub(super) struct Signer<'a> {
673    pub(super) key: &'a NodeKey,
674    pub(super) self_id: [u8; 32],
675}
676
677impl Inner {
678    /// This node's [`Signer`].
679    fn signer(&self) -> Signer<'_> {
680        Signer {
681            key: &self.key,
682            self_id: self.self_id,
683        }
684    }
685}
686
687/// An outcome as the clear signed reply to a clear `request`.
688fn clear_answer(signer: &Signer, request: &VerifiedRequest, outcome: Outcome) -> Option<Value> {
689    match outcome {
690        Outcome::Refused(code, detail) => provider_error(signer, request, code, detail.as_deref()),
691        Outcome::Result(payload) => clear_result(signer, request, &payload),
692    }
693}
694
695/// A result as the signed RESULT to a clear `request`, or payload_too_large
696/// or unknown_error when the wire cannot carry it.
697fn clear_result(signer: &Signer, request: &VerifiedRequest, payload: &Value) -> Option<Value> {
698    match frame::sign_result(request, payload, None, signer.key) {
699        Ok(signed) => Some(signed),
700        Err(_) if cbor::encode(payload).is_ok_and(|e| e.len() > frame::MAX_FRAME_BYTES) => {
701            provider_error(signer, request, CODE_PAYLOAD_TOO_LARGE, None)
702        }
703        Err(_) => provider_error(signer, request, CODE_UNSENDABLE, None),
704    }
705}
706
707/// An outcome as the signed reply to `request`: clear to a clear request,
708/// sealed to a sealed one. A result the wire cannot carry is answered
709/// payload_too_large or unknown_error, as macula answers it. `None` when no
710/// reply signs: the request is left unanswered (macula-rust#14).
711fn answered(
712    signer: &Signer,
713    request: &VerifiedRequest,
714    sealing: Option<&CallSeal>,
715    outcome: Outcome,
716) -> Option<Value> {
717    let Some(sealing) = sealing else {
718        return clear_answer(signer, request, outcome);
719    };
720    let payload = match outcome {
721        Outcome::Refused(code, detail) => {
722            return sealed_error(signer, request, sealing, code, detail.as_deref())
723        }
724        Outcome::Result(payload) => payload,
725    };
726    let plain = frame::check_payload(&payload)
727        .ok()
728        .and_then(|_| cbor::encode(&payload).ok());
729    let Some(plain) = plain else {
730        return sealed_error(signer, request, sealing, CODE_UNSENDABLE, None);
731    };
732    let signed = sealing
733        .sealed_answer(
734            seal::FRAME_RESULT,
735            &plain,
736            &request.request_hash,
737            &signer.self_id,
738        )
739        .and_then(|sealed| {
740            frame::sign_sealed_result(request, &sealed, None, signer.key).map_err(LinkError::from)
741        });
742    match sealed_result_refusal(signed) {
743        Ok(reply) => Some(reply),
744        Err(Refusal::Sealed(code)) => sealed_error(signer, request, sealing, code, None),
745        Err(Refusal::Clear(code)) => provider_error(signer, request, code, None),
746    }
747}
748
749/// How a sealed RESULT that cannot go is refused: sealed, or in the clear
750/// from the closed set.
751#[derive(Debug, PartialEq)]
752enum Refusal {
753    Sealed(&'static str),
754    Clear(&'static str),
755}
756
757/// A signed sealed RESULT, or how to refuse it: payload_too_large sealed when
758/// it is over the frame cap, unavailable in the clear when it could not be
759/// sealed (no randomness), unknown_error sealed when it did not sign
760/// (macula-rust#14).
761fn sealed_result_refusal(signed: Result<Value, LinkError>) -> Result<Value, Refusal> {
762    match signed {
763        Ok(reply) if cbor::encode(&reply).is_ok_and(|e| e.len() <= frame::MAX_FRAME_BYTES) => {
764            Ok(reply)
765        }
766        Ok(_) => Err(Refusal::Sealed(CODE_PAYLOAD_TOO_LARGE)),
767        Err(LinkError::Io(_)) => Err(Refusal::Clear(CODE_UNAVAILABLE)),
768        Err(_) => Err(Refusal::Sealed(CODE_UNSENDABLE)),
769    }
770}
771
772/// This node's ERROR for a sealed request, its code and detail sealed as
773/// cbor([code, detail]). One that cannot be sealed (no randomness) or does
774/// not sign is answered unavailable in the clear, from the closed set: never
775/// the code or detail in the clear (macula-rust#14).
776fn sealed_error(
777    signer: &Signer,
778    request: &VerifiedRequest,
779    sealing: &CallSeal,
780    code: &str,
781    detail: Option<&str>,
782) -> Option<Value> {
783    let plain = seal::error_plain(code, detail.unwrap_or(""));
784    sealing
785        .sealed_answer(
786            seal::FRAME_ERROR,
787            &plain,
788            &request.request_hash,
789            &signer.self_id,
790        )
791        .ok()
792        .and_then(|sealed| {
793            frame::sign_sealed_provider_error(request, &sealed, None, signer.key).ok()
794        })
795        .or_else(|| provider_error(signer, request, CODE_UNAVAILABLE, None))
796}
797
798impl Inner {
799    /// Whether `o`'s advertisements name this node's KEM key.
800    pub(super) fn keyed(&self, o: &Offer) -> bool {
801        self.kem_advertise && self.keyring.is_some() && o.confidential != Confidentiality::Off
802    }
803
804    /// When `o` was first keyed, `None` while its advertisements name no
805    /// key.
806    pub(super) fn keyed_since(&self, o: &Offer) -> Option<i64> {
807        o.keyed_since_ms.filter(|_| self.keyed(o))
808    }
809
810    /// The key `o`'s advertisement names now, as carried, `None` for none.
811    fn kem_key(&self, o: &Offer) -> Option<Vec<u8>> {
812        let keyring = self.keyring.as_ref().filter(|_| self.keyed(o))?;
813        Some(keyring.current().public_key().carried().to_vec())
814    }
815}
816
817/// This node's signed ERROR for `request`, `None` when it does not sign: a
818/// request this key cannot answer, or a key that fails to sign. Either leaves
819/// the request unanswered, counted by the caller of this, never a panic.
820fn provider_error(
821    signer: &Signer,
822    request: &VerifiedRequest,
823    code: &str,
824    detail: Option<&str>,
825) -> Option<Value> {
826    frame::sign_provider_error(request, code, detail, None, signer.key).ok()
827}
828
829async fn send_reply(inner: &Inner, reply: Value) {
830    let _ = inner.write_control(&reply).await;
831}
832
833/// `text` cut to 256 bytes on a character boundary.
834pub(super) fn bounded_detail(text: &str) -> &str {
835    if text.len() <= MAX_DETAIL_BYTES {
836        return text;
837    }
838    let mut cut = MAX_DETAIL_BYTES;
839    while !text.is_char_boundary(cut) {
840        cut -= 1;
841    }
842    &text[..cut]
843}
844
845#[cfg(test)]
846mod tests {
847    //! A provider's sealed answer on its error paths (macula-rust#14): a
848    //! RESULT is called payload_too_large only when it is over the frame cap,
849    //! and an answer that does not sign never panics.
850
851    use super::*;
852    use crate::frame::{FrameError, RequestType};
853    use crate::node_key::NodeKey;
854    use crate::profile::Profile;
855    use crate::seal::Keyring;
856    use crate::station_link::confidential::sealed_request;
857
858    const CALLER: [u8; 32] = [1; 32];
859
860    /// A call sealed to `keyring` and targeting `target`, as its provider
861    /// verified it, and the provider's keys for the answer.
862    fn sealed_call(keyring: &Keyring, target: [u8; 32]) -> (VerifiedRequest, CallSeal) {
863        let to = keyring.current().public_key().carried().to_vec();
864        let (sealed, _) = sealed_request(
865            keyring.profile(),
866            &to,
867            seal::FRAME_CALL,
868            [3; 32],
869            "~ring",
870            CALLER,
871            target,
872            [4; 16],
873            1_790_000_000_000,
874            &Value::text("secret"),
875        )
876        .unwrap();
877        let request = VerifiedRequest {
878            frame_type: RequestType::Call,
879            key: Vec::new(),
880            request_hash: [5; 48],
881            caller: CALLER,
882            request_id: [4; 16],
883            realm: [3; 32],
884            procedure: "~ring".into(),
885            target,
886            deadline: 1_790_000_000_000,
887            payload: Value::Null,
888            sealed: Some(sealed),
889            mode: None,
890            token: None,
891            proofs: None,
892        };
893        let (_, sealing) = opened_request(Some(keyring), &request).unwrap();
894        (request, sealing)
895    }
896
897    #[test]
898    fn only_a_result_over_the_frame_cap_is_payload_too_large() {
899        let oversize = Value::Bytes(vec![0; frame::MAX_FRAME_BYTES + 1]);
900        assert_eq!(
901            sealed_result_refusal(Ok(oversize)),
902            Err(Refusal::Sealed(CODE_PAYLOAD_TOO_LARGE))
903        );
904        assert_eq!(
905            sealed_result_refusal(Err(LinkError::Io("no randomness".into()))),
906            Err(Refusal::Clear(CODE_UNAVAILABLE))
907        );
908        assert_eq!(
909            sealed_result_refusal(Err(LinkError::Frame(FrameError::Unsignable))),
910            Err(Refusal::Sealed(CODE_UNSENDABLE))
911        );
912        assert_eq!(
913            sealed_result_refusal(Ok(Value::text("fits"))),
914            Ok(Value::text("fits"))
915        );
916    }
917
918    #[test]
919    fn a_sealed_answer_that_does_not_sign_does_not_panic() {
920        let profile = Profile::PqPure;
921        let keyring = Keyring::system(profile).unwrap();
922        let key = NodeKey::generate_identity(profile, 0).unwrap();
923        let signer = Signer {
924            key: &key,
925            self_id: key.key_id(),
926        };
927        // Answered as this node, the request names another node as its
928        // target, so no reply to it signs.
929        let (request, sealing) = sealed_call(&keyring, [2; 32]);
930        let refused = Outcome::Refused(CODE_HANDLER_ERROR, Some("no".into()));
931        assert_eq!(answered(&signer, &request, Some(&sealing), refused), None);
932        let result = Outcome::Result(Value::text("answer"));
933        assert_eq!(answered(&signer, &request, Some(&sealing), result), None);
934        // Targeting this node, the same answers sign, sealed.
935        let (request, sealing) = sealed_call(&keyring, key.key_id());
936        let result = Outcome::Result(Value::text("answer"));
937        let reply = answered(&signer, &request, Some(&sealing), result).unwrap();
938        assert!(reply.get("reply").is_some(), "{reply:?}");
939    }
940}