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