1use 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::record::{
30 self, Authorization, ProcedureAdvertisementOptions, Reason, Record, RecordError,
31 TombstoneOptions, Trust,
32};
33use crate::seal;
34
35use super::admission::Verdict;
36use super::confidential::{
37 clear_allowed, opened_request, CallSeal, Confidentiality, CODE_SEALED_REFUSED,
38 CODE_SEALED_REQUIRED,
39};
40use super::framing::MAX_FRAME_BYTES;
41use super::stream::StreamHandler;
42use super::{now_ms, Inner, Link, LinkError};
43
44const CODE_HANDLER_ERROR: &str = "handler_error";
48const CODE_HANDLER_CRASHED: &str = "temporary_relay_failure";
49const CODE_UNKNOWN_PROCEDURE: &str = "unknown_next_peer";
50pub(super) const CODE_REQUEST_COPY: &str = "request_copy";
51const CODE_PAYLOAD_TOO_LARGE: &str = "payload_too_large";
52const CODE_UNSENDABLE: &str = "unknown_error";
53const CODE_UNAVAILABLE: &str = "unavailable";
56
57const MAX_DETAIL_BYTES: usize = 256;
59
60const MAX_ADVERTISEMENT_TTL: Duration =
62 Duration::from_millis(super::confidential::MAX_ADVERTISEMENT_TTL_MS as u64);
63const REFRESH_RETRY: Duration = Duration::from_secs(10);
66
67pub type BoxFuture<T> = Pin<Box<dyn Future<Output = T> + Send + 'static>>;
69
70pub type Handler = Arc<dyn Fn(Request) -> BoxFuture<Result<Value, String>> + Send + Sync>;
74
75pub fn handler<F, Fut>(f: F) -> Handler
77where
78 F: Fn(Request) -> Fut + Send + Sync + 'static,
79 Fut: Future<Output = Result<Value, String>> + Send + 'static,
80{
81 Arc::new(move |r| Box::pin(f(r)))
82}
83
84#[derive(Debug, Clone, PartialEq)]
88pub struct Request {
89 pub caller: [u8; 32],
90 pub realm: [u8; 32],
91 pub procedure: String,
92 pub payload: Value,
93 pub token: Option<Vec<u8>>,
94 pub proofs: Option<Vec<Vec<u8>>>,
95 pub deadline_ms: u64,
96 pub sealed: bool,
99}
100
101#[derive(Clone)]
106pub struct Offer {
107 pub realm: [u8; 32],
108 pub procedure: String,
109 pub handler: Option<Handler>,
110 pub stream: Option<StreamOffer>,
111 pub realm_key: Option<Vec<u8>>,
112 pub confidential: Confidentiality,
116 pub keyed_since_ms: Option<i64>,
121}
122
123impl Offer {
124 pub fn unary(realm: [u8; 32], procedure: &str, handler: Handler) -> Offer {
126 Offer {
127 realm,
128 procedure: procedure.to_string(),
129 handler: Some(handler),
130 stream: None,
131 realm_key: None,
132 confidential: Confidentiality::Preferred,
133 keyed_since_ms: None,
134 }
135 }
136
137 pub fn stream(
139 realm: [u8; 32],
140 procedure: &str,
141 mode: StreamMode,
142 handler: StreamHandler,
143 ) -> Offer {
144 Offer {
145 realm,
146 procedure: procedure.to_string(),
147 handler: None,
148 stream: Some(StreamOffer { mode, handler }),
149 realm_key: None,
150 confidential: Confidentiality::Preferred,
151 keyed_since_ms: None,
152 }
153 }
154}
155
156#[derive(Clone)]
160pub struct StreamOffer {
161 pub mode: StreamMode,
162 pub handler: StreamHandler,
163}
164
165pub(super) type ServedEntry = Arc<ServedInner>;
167
168pub(super) struct ServedInner {
169 link: Weak<Inner>,
170 key: ([u8; 32], String),
171 pub(super) offer: Offer,
172 latest: Mutex<Record>,
173 err: Mutex<Option<LinkError>>,
174 done_tx: watch::Sender<bool>,
175}
176
177#[derive(Clone)]
181pub struct Served {
182 inner: Arc<ServedInner>,
183}
184
185impl Link {
186 pub async fn serve(&self, mut o: Offer) -> Result<Served, LinkError> {
197 let own = record::in_own_namespace(&o.procedure);
198 if o.handler.is_some() == o.stream.is_some() || (!own && o.realm_key.is_none()) {
199 return Err(LinkError::InvalidOffer);
200 }
201 if !matches!(record::procedure_org(&o.procedure), Ok(Some(_))) {
202 return Err(LinkError::NoOrg);
203 }
204 let can_name_key = self.inner.kem_advertise && self.inner.keyring.is_some();
205 if o.confidential == Confidentiality::Required && !can_name_key {
206 return Err(LinkError::KemAdvertiseDisabled);
207 }
208 if self.inner.keyed(&o) && o.keyed_since_ms.is_none() {
209 o.keyed_since_ms = Some(now_ms());
210 }
211 let (advertisement, wire) = self.advertisement(&o, MAX_ADVERTISEMENT_TTL).await?;
212 let key = (o.realm, o.procedure.clone());
213 let (done_tx, _) = watch::channel(false);
214 let served = Arc::new(ServedInner {
215 link: Arc::downgrade(&self.inner),
216 key: key.clone(),
217 offer: o,
218 latest: Mutex::new(advertisement),
219 err: Mutex::new(None),
220 done_tx,
221 });
222 {
223 let mut state = self.inner.lock();
224 if let Some(e) = &state.ended {
225 return Err(e.clone());
226 }
227 if state.served.contains_key(&key) {
228 return Err(LinkError::AlreadyServed);
229 }
230 state.served.insert(key, served.clone());
231 }
232 if let Err(e) = self.announce(&wire).await {
233 served.end(e.clone());
234 return Err(e);
235 }
236 tokio::spawn(renew(served.clone()));
237 Ok(Served { inner: served })
238 }
239
240 async fn advertisement(
244 &self,
245 o: &Offer,
246 max_ttl: Duration,
247 ) -> Result<(Record, Vec<u8>), LinkError> {
248 let inner = &self.inner;
249 let max_ttl_ms = max_ttl.as_millis() as u64;
250 let kem_key = inner.kem_key(o);
252 let opts = if record::in_own_namespace(&o.procedure) {
253 ProcedureAdvertisementOptions {
254 authorization: Authorization::None,
255 kem_key,
256 ttl_ms: max_ttl_ms,
257 }
258 } else {
259 let org = record::procedure_org(&o.procedure)?.ok_or(LinkError::NoOrg)?;
260 let directory = self
261 .find_record(&record::org_directory_key(&o.realm, org))
262 .await?;
263 let named = record::read_org_directory(directory.record())?;
264 let delegation = self
265 .find_record(&record::procedure_delegation_key(
266 &named.org_key,
267 &inner.self_id,
268 ))
269 .await?;
270 let now = now_ms();
271 let ttl = (max_ttl_ms as i64)
272 .min(directory.record().expires_at as i64 - now)
273 .min(delegation.record().expires_at as i64 - now);
274 if ttl <= 0 {
275 return Err(RecordError::AuthorizationOutlived.into());
276 }
277 ProcedureAdvertisementOptions {
278 authorization: Authorization::Delegation {
279 org_directory: record::encode(directory.record())?,
280 procedure_delegation: record::encode(delegation.record())?,
281 },
282 kem_key,
283 ttl_ms: ttl as u64,
284 }
285 };
286 let unsigned = record::new_procedure_advertisement(
287 &inner.self_id,
288 &o.realm,
289 &o.procedure,
290 &inner.station.node_id,
291 &opts,
292 )?;
293 let signed = record::sign(&unsigned, &inner.key)?;
296 let wire = record::encode(&signed)?;
297 let now = now_ms();
298 let verified = record::verify(&wire, inner.profile, now)?;
299 record::verify_authorization(
300 &verified,
301 &Trust {
302 profile: inner.profile,
303 realm_key: o.realm_key.clone(),
304 },
305 now,
306 )?;
307 Ok((signed, wire))
308 }
309
310 async fn announce(&self, wire: &[u8]) -> Result<(), LinkError> {
314 self.inner
315 .send_control(&frame::advertise_frame(wire))
316 .await?;
317 self.put_record(wire).await
318 }
319}
320
321impl Served {
322 pub async fn stop(&self) -> Result<(), LinkError> {
327 if self.inner.is_done() {
328 return Ok(());
329 }
330 self.inner.end(LinkError::Stopped);
331 let Some(inner) = self.inner.link.upgrade() else {
332 return Ok(());
333 };
334 let latest = self
335 .inner
336 .latest
337 .lock()
338 .unwrap_or_else(|p| p.into_inner())
339 .clone();
340 let tombstone =
341 record::new_tombstone(&latest, Reason::Shutdown, &TombstoneOptions::default())?;
342 let wire = record::encode(&record::sign(&tombstone, &inner.key)?)?;
343 inner.send_control(&frame::unadvertise_frame(&wire)).await?;
344 Link { inner }.put_record(&wire).await
345 }
346
347 pub async fn done(&self) -> LinkError {
349 let mut done = self.inner.done_tx.subscribe();
350 let _ = done.wait_for(|ended| *ended).await;
351 self.error().unwrap_or(LinkError::Stopped)
352 }
353
354 pub fn error(&self) -> Option<LinkError> {
358 self.inner
359 .err
360 .lock()
361 .unwrap_or_else(|p| p.into_inner())
362 .clone()
363 }
364}
365
366impl ServedInner {
367 fn is_done(&self) -> bool {
368 *self.done_tx.borrow()
369 }
370
371 pub(super) fn end(self: &Arc<Self>, err: LinkError) {
374 {
375 let mut held = self.err.lock().unwrap_or_else(|p| p.into_inner());
376 if held.is_some() {
377 return;
378 }
379 *held = Some(err);
380 }
381 if let Some(inner) = self.link.upgrade() {
382 let mut state = inner.lock();
383 if state
384 .served
385 .get(&self.key)
386 .is_some_and(|s| Arc::ptr_eq(s, self))
387 {
388 state.served.remove(&self.key);
389 }
390 }
391 let _ = self.done_tx.send_replace(true);
392 }
393}
394
395async fn renew(served: Arc<ServedInner>) {
400 let mut current = served
401 .latest
402 .lock()
403 .unwrap_or_else(|p| p.into_inner())
404 .clone();
405 let mut wait = half_life(¤t);
406 let mut last_err: Option<LinkError> = None;
407 let mut stopped = served.done_tx.subscribe();
408 loop {
409 let Some(mut link_done) = served.link.upgrade().map(|l| l.done_rx.clone()) else {
410 return;
411 };
412 tokio::select! {
413 _ = stopped.wait_for(|ended| *ended) => return,
414 _ = link_done.wait_for(|ended| *ended) => {
415 let err = served.link.upgrade().and_then(|l| l.lock().ended.clone()).unwrap_or(LinkError::Closed);
416 served.end(err);
417 return;
418 }
419 _ = tokio::time::sleep(wait) => {}
420 }
421 if now_ms() >= current.expires_at as i64 {
422 served.end(last_err.unwrap_or(LinkError::Stopped));
423 return;
424 }
425 let Some(inner) = served.link.upgrade() else {
426 return;
427 };
428 let link = Link { inner };
429 let renewed = async {
430 let (advertisement, wire) = link
431 .advertisement(&served.offer, MAX_ADVERTISEMENT_TTL)
432 .await?;
433 link.announce(&wire).await?;
434 Ok::<_, LinkError>(advertisement)
435 }
436 .await;
437 match renewed {
438 Ok(advertisement) => {
439 *served.latest.lock().unwrap_or_else(|p| p.into_inner()) = advertisement.clone();
440 wait = half_life(&advertisement);
441 current = advertisement;
442 }
443 Err(e) => {
444 last_err = Some(e);
445 let left = (current.expires_at as i64 - now_ms()).max(0) as u64;
446 wait = REFRESH_RETRY.min(Duration::from_millis(left));
447 }
448 }
449 }
450}
451
452fn half_life(r: &Record) -> Duration {
453 Duration::from_millis(r.expires_at.saturating_sub(r.created_at) / 2)
454}
455
456pub(super) fn called(inner: &Arc<Inner>, v: &Value) {
461 let Ok(request) = frame::verify_request(v, inner.profile) else {
462 inner.count("unverified_call");
463 return;
464 };
465 if request.target != inner.self_id {
466 inner.count("call_for_another_node");
467 return;
468 }
469 let inner = inner.clone();
470 match inner.admission.admit(&request, &inner.share, now_ms()) {
471 Verdict::Refused(code) => {
472 let reply = provider_error(&inner, &request, code, None);
473 tokio::spawn(async move { send_reply(&inner, reply).await });
474 }
475 Verdict::Copy(None) => {
476 let reply = provider_error(&inner, &request, CODE_REQUEST_COPY, None);
477 tokio::spawn(async move { send_reply(&inner, reply).await });
478 }
479 Verdict::Copy(Some(stored)) => {
480 tokio::spawn(async move {
481 let _ = inner.control.write(&stored, MAX_FRAME_BYTES).await;
482 });
483 }
484 Verdict::New => {
485 tokio::spawn(answer(inner, request));
486 }
487 }
488}
489
490async fn answer(inner: Arc<Inner>, request: VerifiedRequest) {
493 let offer = inner
494 .lock()
495 .served
496 .get(&(request.realm, request.procedure.clone()))
497 .map(|s| s.offer.clone());
498 let reply = reply(&inner, &request, offer).await;
499 let Ok(encoded) = cbor::encode(&reply) else {
500 inner.count("unencodable_reply");
501 return;
502 };
503 inner.admission.store(&request, encoded.clone());
504 let _ = inner.control.write(&encoded, MAX_FRAME_BYTES).await;
505}
506
507async fn reply(inner: &Inner, request: &VerifiedRequest, offer: Option<Offer>) -> Value {
513 let (payload, sealing) = match &request.sealed {
514 Some(_) => match opened_request(inner.keyring.as_deref(), request) {
515 Ok((payload, sealing)) => (payload, Some(sealing)),
516 Err(detail) => {
517 return provider_error(inner, request, CODE_SEALED_REFUSED, Some(&detail))
518 }
519 },
520 None => {
521 let refused = offer
522 .as_ref()
523 .is_some_and(|o| !clear_allowed(o.confidential, inner.keyed_since(o), now_ms()));
524 if refused {
525 return provider_error(inner, request, CODE_SEALED_REQUIRED, None);
526 }
527 (request.payload.clone(), None)
528 }
529 };
530 let outcome = match offer.and_then(|o| o.handler) {
531 None => Outcome::Refused(CODE_UNKNOWN_PROCEDURE, None),
532 Some(handler) => handled(handler, request, payload, sealing.is_some()).await,
533 };
534 answered(inner, request, sealing.as_ref(), outcome)
535}
536
537enum Outcome {
540 Result(Value),
541 Refused(&'static str, Option<String>),
542}
543
544async fn handled(
548 handler: Handler,
549 request: &VerifiedRequest,
550 payload: Value,
551 sealed: bool,
552) -> Outcome {
553 let running = tokio::spawn(handler(Request {
554 caller: request.caller,
555 realm: request.realm,
556 procedure: request.procedure.clone(),
557 payload,
558 token: request.token.clone(),
559 proofs: request.proofs.clone(),
560 deadline_ms: request.deadline,
561 sealed,
562 }));
563 let abort = running.abort_handle();
564 let left = (request.deadline as i64 - now_ms()).max(0) as u64;
565 match tokio::time::timeout(Duration::from_millis(left), running).await {
566 Err(_) => {
567 abort.abort();
568 Outcome::Refused(
569 CODE_HANDLER_ERROR,
570 Some("the request's deadline passed".into()),
571 )
572 }
573 Ok(Err(_panicked)) => Outcome::Refused(CODE_HANDLER_CRASHED, None),
574 Ok(Ok(Err(refusal))) => Outcome::Refused(
575 CODE_HANDLER_ERROR,
576 Some(bounded_detail(&refusal).to_string()),
577 ),
578 Ok(Ok(Ok(payload))) => Outcome::Result(payload),
579 }
580}
581
582fn answered(
586 inner: &Inner,
587 request: &VerifiedRequest,
588 sealing: Option<&CallSeal>,
589 outcome: Outcome,
590) -> Value {
591 let Some(sealing) = sealing else {
592 return match outcome {
593 Outcome::Refused(code, detail) => {
594 provider_error(inner, request, code, detail.as_deref())
595 }
596 Outcome::Result(payload) => {
597 match frame::sign_result(request, &payload, None, &inner.key) {
598 Ok(signed) => signed,
599 Err(_)
600 if cbor::encode(&payload)
601 .is_ok_and(|e| e.len() > frame::MAX_FRAME_BYTES) =>
602 {
603 provider_error(inner, request, CODE_PAYLOAD_TOO_LARGE, None)
604 }
605 Err(_) => provider_error(inner, request, CODE_UNSENDABLE, None),
606 }
607 }
608 };
609 };
610 let payload = match outcome {
611 Outcome::Refused(code, detail) => {
612 return sealed_error(inner, request, sealing, code, detail.as_deref())
613 }
614 Outcome::Result(payload) => payload,
615 };
616 let plain = frame::check_payload(&payload)
617 .ok()
618 .and_then(|_| cbor::encode(&payload).ok());
619 let Some(plain) = plain else {
620 return sealed_error(inner, request, sealing, CODE_UNSENDABLE, None);
621 };
622 let signed = sealing
623 .sealed_answer(
624 seal::FRAME_RESULT,
625 &plain,
626 &request.request_hash,
627 &inner.self_id,
628 )
629 .and_then(|sealed| {
630 frame::sign_sealed_result(request, &sealed, None, &inner.key).map_err(LinkError::from)
631 });
632 match signed {
633 Ok(reply) if cbor::encode(&reply).is_ok_and(|e| e.len() <= frame::MAX_FRAME_BYTES) => reply,
634 Ok(_) => sealed_error(inner, request, sealing, CODE_PAYLOAD_TOO_LARGE, None),
635 Err(LinkError::Io(_)) => provider_error(inner, request, CODE_UNAVAILABLE, None),
636 Err(_) => sealed_error(inner, request, sealing, CODE_PAYLOAD_TOO_LARGE, None),
637 }
638}
639
640fn sealed_error(
645 inner: &Inner,
646 request: &VerifiedRequest,
647 sealing: &CallSeal,
648 code: &str,
649 detail: Option<&str>,
650) -> Value {
651 let plain = seal::error_plain(code, detail.unwrap_or(""));
652 match sealing.sealed_answer(
653 seal::FRAME_ERROR,
654 &plain,
655 &request.request_hash,
656 &inner.self_id,
657 ) {
658 Ok(sealed) => frame::sign_sealed_provider_error(request, &sealed, None, &inner.key)
659 .unwrap_or_else(|e| {
660 panic!("station_link: a sealed provider error that does not sign: {e}")
661 }),
662 Err(_) => provider_error(inner, request, CODE_UNAVAILABLE, None),
663 }
664}
665
666impl Inner {
667 pub(super) fn keyed(&self, o: &Offer) -> bool {
669 self.kem_advertise && self.keyring.is_some() && o.confidential != Confidentiality::Off
670 }
671
672 pub(super) fn keyed_since(&self, o: &Offer) -> Option<i64> {
675 o.keyed_since_ms.filter(|_| self.keyed(o))
676 }
677
678 fn kem_key(&self, o: &Offer) -> Option<Vec<u8>> {
680 let keyring = self.keyring.as_ref().filter(|_| self.keyed(o))?;
681 Some(keyring.current().public_key().carried().to_vec())
682 }
683}
684
685fn provider_error(
689 inner: &Inner,
690 request: &VerifiedRequest,
691 code: &str,
692 detail: Option<&str>,
693) -> Value {
694 frame::sign_provider_error(request, code, detail, None, &inner.key)
695 .unwrap_or_else(|e| panic!("station_link: a provider error that does not sign: {e}"))
696}
697
698async fn send_reply(inner: &Inner, reply: Value) {
699 let _ = inner.write_control(&reply).await;
700}
701
702pub(super) fn bounded_detail(text: &str) -> &str {
704 if text.len() <= MAX_DETAIL_BYTES {
705 return text;
706 }
707 let mut cut = MAX_DETAIL_BYTES;
708 while !text.is_char_boundary(cut) {
709 cut -= 1;
710 }
711 &text[..cut]
712}