1use std::collections::VecDeque;
18use std::future::Future;
19use std::sync::{Arc, Mutex, MutexGuard, Weak};
20use std::time::Duration;
21
22use tokio::sync::{watch, Notify};
23
24use crate::cbor::{self, Value};
25use crate::frame::{
26 self, RequestSpec, StreamEncoding, StreamFields, StreamMode, StreamRole, StreamState,
27 VerifiedRequest,
28};
29
30use super::admission::{Admission, SessionPlace, Verdict};
31use super::confidential::{
32 clear_allowed, opened_request, sealed_request, stated, unsealed, Seal, StreamSeal,
33 CODE_SEALED_REFUSED, CODE_SEALED_REQUIRED,
34};
35use super::framing::{read_frame, FrameWriter, MAX_FRAME_BYTES};
36use super::serve::{
37 bounded_detail, without_caller, BoxFuture, Offer, StreamOffer, CODE_REQUEST_COPY,
38};
39use super::{frame_type_of, now_ms, Inner, Link, LinkError};
40
41const STREAM_OPEN_BYTES: usize = 1024 * 1024;
42const STREAM_OPEN_WAIT: Duration = Duration::from_secs(10);
43const STREAM_INBOX: usize = 16 * 1024 * 1024;
44
45pub const DEFAULT_STREAM_DEADLINE: Duration = Duration::from_secs(30);
48
49const CODE_STREAM_NOT_FOUND: &str = "not_found";
52const CODE_MODE_MISMATCH: &str = "mode_mismatch";
53const CODE_TOO_MANY_SESSIONS: &str = "too_many_sessions";
54const CODE_STREAM_HANDLER_ERROR: &str = "error";
55
56pub type StreamHandler = Arc<dyn Fn(Stream) -> BoxFuture<Result<(), String>> + Send + Sync>;
61
62pub fn stream_handler<F, Fut>(f: F) -> StreamHandler
64where
65 F: Fn(Stream) -> Fut + Send + Sync + 'static,
66 Fut: Future<Output = Result<(), String>> + Send + 'static,
67{
68 Arc::new(move |s| Box::pin(f(s)))
69}
70
71#[derive(Debug, Clone, PartialEq)]
77pub struct StreamCall {
78 pub realm: [u8; 32],
79 pub procedure: String,
80 pub target: [u8; 32],
81 pub mode: StreamMode,
82 pub payload: Value,
83 pub deadline: Duration,
84 pub token: Option<Vec<u8>>,
85 pub proofs: Vec<Vec<u8>>,
86 pub seal: Option<Seal>,
87}
88
89impl Default for StreamCall {
90 fn default() -> Self {
91 StreamCall {
92 realm: [0; 32],
93 procedure: String::new(),
94 target: [0; 32],
95 mode: StreamMode::ServerStream,
96 payload: Value::Map(Vec::new()),
97 deadline: Duration::ZERO,
98 token: None,
99 proofs: Vec::new(),
100 seal: None,
101 }
102 }
103}
104
105#[derive(Debug, Clone, PartialEq)]
109pub enum StreamEvent {
110 Data {
111 encoding: StreamEncoding,
112 body: Value,
113 },
114 End {
115 role: StreamRole,
116 },
117 Reply {
118 payload: Value,
119 },
120}
121
122#[derive(Clone)]
124pub struct Stream {
125 pub(super) inner: Arc<StreamInner>,
126}
127
128struct Budget {
131 admission: Arc<Admission>,
132 caller: [u8; 32],
133 place: Option<SessionPlace>,
134}
135
136pub(super) struct StreamInner {
137 link: Arc<Inner>,
138 writer: FrameWriter,
139 pub(super) open: VerifiedRequest,
140 pub(super) caller: bool,
141 pub(super) sealing: Option<StreamSeal>,
143 send_seq: tokio::sync::Mutex<u64>,
145 state: Mutex<StreamSide>,
146 budget: Mutex<Option<Budget>>,
147 notify: Notify,
148 done_tx: watch::Sender<bool>,
149}
150
151#[derive(Default)]
152pub(super) struct StreamSide {
153 sent_end: bool,
155 peer_ended: bool,
157 inbox: VecDeque<(StreamEvent, usize)>,
158 held: usize,
159 ended: bool,
160 err: Option<LinkError>,
161 pub(super) settled: bool,
163}
164
165impl Stream {
166 pub fn request(&self) -> &VerifiedRequest {
171 &self.inner.open
172 }
173
174 pub fn sealed(&self) -> bool {
177 self.inner.sealing.is_some()
178 }
179
180 pub async fn send(&self, body: &[u8]) -> Result<(), LinkError> {
182 self.inner
183 .send(
184 |seq| StreamFields::Data {
185 seq,
186 encoding: StreamEncoding::Raw,
187 body: Value::Bytes(body.to_vec()),
188 },
189 false,
190 )
191 .await
192 }
193
194 pub async fn send_value(&self, v: Value) -> Result<(), LinkError> {
196 self.inner
197 .send(
198 |seq| StreamFields::Data {
199 seq,
200 encoding: StreamEncoding::Msgpack,
201 body: v.clone(),
202 },
203 false,
204 )
205 .await
206 }
207
208 pub async fn close_send(&self) -> Result<(), LinkError> {
210 self.inner
211 .send(
212 |seq| StreamFields::End {
213 seq,
214 role: StreamRole::Send,
215 },
216 true,
217 )
218 .await
219 }
220
221 pub async fn close(&self) -> Result<(), LinkError> {
223 let sent = self
224 .inner
225 .send(
226 |seq| StreamFields::End {
227 seq,
228 role: StreamRole::Both,
229 },
230 true,
231 )
232 .await;
233 StreamInner::end(&self.inner, None);
234 sent
235 }
236
237 pub async fn reply(&self, payload: Value) -> Result<(), LinkError> {
239 let sent = self
240 .inner
241 .send(
242 |seq| StreamFields::Reply {
243 seq,
244 payload: payload.clone(),
245 },
246 true,
247 )
248 .await;
249 StreamInner::end(&self.inner, None);
250 sent
251 }
252
253 pub async fn abort(&self, code: &str, message: &str) -> Result<(), LinkError> {
255 self.inner.abort(code, message).await
256 }
257
258 pub async fn recv(&self) -> Result<StreamEvent, LinkError> {
262 loop {
263 let notified = self.inner.notify.notified();
264 match self.inner.next_event() {
265 Some(outcome) => return outcome,
266 None => notified.await,
267 }
268 }
269 }
270
271 pub async fn done(&self) -> Option<LinkError> {
274 let mut done = self.inner.done_tx.subscribe();
275 let _ = done.wait_for(|ended| *ended).await;
276 self.inner.side().err.clone()
277 }
278}
279
280impl StreamInner {
281 fn new(
282 link: Arc<Inner>,
283 send: quinn::SendStream,
284 open: VerifiedRequest,
285 caller: bool,
286 sealing: Option<StreamSeal>,
287 ) -> Arc<StreamInner> {
288 Arc::new(StreamInner {
289 link,
290 writer: FrameWriter::new(send),
291 open,
292 caller,
293 sealing,
294 send_seq: tokio::sync::Mutex::new(0),
295 state: Mutex::new(StreamSide::default()),
296 budget: Mutex::new(None),
297 notify: Notify::new(),
298 done_tx: watch::channel(false).0,
299 })
300 }
301
302 pub(super) fn side(&self) -> MutexGuard<'_, StreamSide> {
303 self.state.lock().unwrap_or_else(|p| p.into_inner())
304 }
305
306 fn budget(&self) -> MutexGuard<'_, Option<Budget>> {
307 self.budget.lock().unwrap_or_else(|p| p.into_inner())
308 }
309
310 fn next_event(&self) -> Option<Result<StreamEvent, LinkError>> {
314 let mut side = self.side();
315 if let Some((event, size)) = side.inbox.pop_front() {
316 side.held -= size;
317 drop(side);
318 self.release_inbox(size);
319 return Some(Ok(event));
320 }
321 if side.ended {
322 return Some(Err(side.err.clone().unwrap_or(LinkError::EndOfStream)));
323 }
324 None
325 }
326
327 async fn send(
334 self: &Arc<Self>,
335 at: impl FnOnce(u64) -> StreamFields,
336 last: bool,
337 ) -> Result<(), LinkError> {
338 let mut seq = self.send_seq.lock().await;
339 if self.side().sent_end {
340 return Err(LinkError::StreamClosed);
341 }
342 let fields = at(*seq);
343 let plain = match &self.sealing {
344 Some(sealing) => sealing.plain_of(&fields)?,
345 None => None,
346 };
347 let spent = plain.is_some();
348 let fields = match (&self.sealing, plain) {
349 (Some(sealing), Some(plain)) => {
350 *seq += 1;
351 let sealed = sealing.sealed(fields, plain);
352 sealed.inspect_err(|_| self.side().sent_end = true)?
353 }
354 _ => fields,
355 };
356 let encoded = self
357 .signed(&fields)
358 .inspect_err(|_| self.unsent_after_spending(spent))?;
359 if let Err(e) = self.writer.write(&encoded, MAX_FRAME_BYTES).await {
360 self.side().sent_end = true;
361 return Err(e);
362 }
363 if !spent {
364 *seq += 1;
365 }
366 if last {
367 self.sent_last().await;
368 }
369 Ok(())
370 }
371
372 fn unsent_after_spending(&self, spent: bool) {
375 if spent {
376 self.side().sent_end = true;
377 }
378 }
379
380 async fn sent_last(self: &Arc<Self>) {
383 let peer_ended = {
384 let mut side = self.side();
385 side.sent_end = true;
386 side.peer_ended
387 };
388 self.writer.finish().await;
389 if peer_ended {
390 StreamInner::end(self, None);
391 }
392 }
393
394 fn signed(&self, fields: &StreamFields) -> Result<Vec<u8>, LinkError> {
396 let signed = if self.caller {
397 frame::sign_caller_stream(fields, &self.open, &self.link.key)?
398 } else {
399 frame::sign_provider_stream(fields, &self.open, &self.link.key)?
400 };
401 cbor::encode(&signed)
402 .map_err(|e| LinkError::Frame(frame::FrameError::Payload(e.to_string())))
403 }
404
405 async fn abort(self: &Arc<Self>, code: &str, message: &str) -> Result<(), LinkError> {
406 let sent = self
407 .send(
408 |seq| StreamFields::Error {
409 seq,
410 code: code.to_string(),
411 message: message.to_string(),
412 },
413 true,
414 )
415 .await;
416 StreamInner::end(
417 self,
418 Some(LinkError::Stream {
419 code: code.to_string(),
420 message: message.to_string(),
421 relay: false,
422 }),
423 );
424 sent
425 }
426
427 fn deliver(&self, event: StreamEvent, size: usize) -> bool {
430 if !self.queue(event, size) {
431 return false;
432 }
433 self.notify.notify_one();
434 true
435 }
436
437 fn queue(&self, event: StreamEvent, size: usize) -> bool {
440 let mut side = self.side();
441 if side.held + size > STREAM_INBOX {
442 return false;
443 }
444 if !self.charge_inbox(size) {
445 return false;
446 }
447 side.inbox.push_back((event, size));
448 side.held += size;
449 true
450 }
451
452 fn charge_inbox(&self, size: usize) -> bool {
455 let budget = self.budget();
456 let Some(budget) = &*budget else {
457 return true;
458 };
459 budget.admission.charge_inbox(budget.caller, size)
460 }
461
462 fn release_inbox(&self, size: usize) {
463 if let Some(budget) = &*self.budget() {
464 budget.admission.release_inbox(budget.caller, size);
465 }
466 }
467
468 fn peer_finished(self: &Arc<Self>, err: Option<LinkError>) {
471 self.side().peer_ended = true;
472 StreamInner::end(self, err);
473 }
474
475 async fn fail(self: &Arc<Self>, code: &str, cause: Option<String>) {
479 let message = cause
480 .as_deref()
481 .map(bounded_detail)
482 .unwrap_or("")
483 .to_string();
484 let _ = self.abort(code, &message).await;
485 }
486
487 pub(super) fn end(this: &Arc<StreamInner>, err: Option<LinkError>) {
491 let Some(graceful) = this.mark_ended(err) else {
492 return;
493 };
494 if !graceful {
495 StreamInner::reset_sending(this);
496 }
497 if let Some(budget) = this.budget().take() {
498 let held = std::mem::take(&mut this.side().held);
499 budget.admission.release_inbox(budget.caller, held);
500 drop(budget.place);
501 }
502 this.link
503 .lock()
504 .streams
505 .retain(|w| w.strong_count() > 0 && !std::ptr::eq(w.as_ptr(), Arc::as_ptr(this)));
506 let _ = this.done_tx.send_replace(true);
507 this.notify.notify_waiters();
508 this.notify.notify_one();
509 }
510}
511
512impl StreamInner {
513 fn mark_ended(&self, err: Option<LinkError>) -> Option<bool> {
517 let mut side = self.side();
518 if side.ended {
519 return None;
520 }
521 side.ended = true;
522 if side.err.is_none() {
523 side.err = err;
524 }
525 let graceful = side.sent_end;
526 side.sent_end = true;
527 Some(graceful)
528 }
529
530 fn reset_sending(this: &Arc<StreamInner>) {
533 let released = this.clone();
534 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
535 return;
536 };
537 runtime.spawn(async move { released.writer.reset().await });
538 }
539}
540
541fn hold_stream(inner: &Inner, s: &Arc<StreamInner>) -> bool {
544 let mut state = inner.lock();
545 if state.ended.is_some() {
546 return false;
547 }
548 state.streams.push(Arc::downgrade(s));
549 true
550}
551
552fn abandon(mut send: quinn::SendStream, mut recv: quinn::RecvStream) {
554 let _ = send.reset(0u32.into());
555 let _ = recv.stop(0u32.into());
556}
557
558impl Link {
559 pub async fn open_stream(&self, c: StreamCall) -> Result<Stream, LinkError> {
563 let inner = &self.inner;
564 stated(&c.target, &inner.station.node_id, &c.seal)?;
565 let deadline = if c.deadline.is_zero() {
566 DEFAULT_STREAM_DEADLINE
567 } else {
568 c.deadline
569 };
570 let mut request_id = [0u8; 16];
571 aws_lc_rs::rand::fill(&mut request_id)
572 .map_err(|_| LinkError::Io("no randomness".into()))?;
573 let deadline = (now_ms() + deadline.as_millis() as i64) as u64;
574 let (sealed, sealing) = sealed_open(inner, &c, request_id, deadline)?;
575 let signed = frame::sign_stream_open(
576 &RequestSpec {
577 request_id,
578 realm: c.realm,
579 procedure: c.procedure,
580 target: c.target,
581 deadline,
582 payload: c.payload,
583 sealed,
584 mode: Some(c.mode),
585 token: c.token,
586 proofs: c.proofs,
587 source_route: None,
588 retry_budget: None,
589 },
590 &inner.key,
591 )?;
592 let encoded = cbor::encode(&signed)
593 .map_err(|e| LinkError::Frame(frame::FrameError::Payload(e.to_string())))?;
594 if encoded.len() > STREAM_OPEN_BYTES {
595 return Err(LinkError::StreamOpenTooLarge(encoded.len()));
596 }
597 let open = frame::verify_request(&signed, inner.profile)?;
598 let state = frame::open_stream(&open)?;
599 let (send, recv) = inner
600 .connection
601 .open_bi()
602 .await
603 .map_err(|e| LinkError::Io(format!("open a stream: {e}")))?;
604 let s = StreamInner::new(inner.clone(), send, open, true, sealing);
605 let held = hold_stream(inner, &s);
606 let written = match held {
607 true => s.writer.write(&encoded, STREAM_OPEN_BYTES).await,
608 false => Err(inner.lock().ended.clone().unwrap_or(LinkError::Closed)),
609 };
610 if let Err(e) = written {
611 StreamInner::end(&s, Some(e.clone()));
612 let mut recv = recv;
613 let _ = recv.stop(0u32.into());
614 return Err(e);
615 }
616 tokio::spawn(read(s.clone(), recv, state));
617 Ok(Stream { inner: s })
618 }
619}
620
621fn sealed_open(
624 inner: &Inner,
625 c: &StreamCall,
626 request_id: [u8; 16],
627 deadline: u64,
628) -> Result<(Option<frame::Sealed>, Option<StreamSeal>), LinkError> {
629 match &c.seal {
630 Some(Seal::To(key)) => {
631 let (sealed, s) = sealed_request(
632 inner.profile,
633 key,
634 crate::seal::FRAME_STREAM_OPEN,
635 c.realm,
636 &c.procedure,
637 inner.self_id,
638 c.target,
639 request_id,
640 deadline,
641 &c.payload,
642 )?;
643 Ok((Some(sealed), Some(StreamSeal::caller(&s))))
644 }
645 _ => Ok((None, None)),
646 }
647}
648
649pub(super) async fn accept_streams(link: Weak<Inner>) {
651 let Some(connection) = link.upgrade().map(|l| l.connection.clone()) else {
652 return;
653 };
654 while let Ok((send, recv)) = connection.accept_bi().await {
655 tokio::spawn(incoming(link.clone(), send, recv));
656 }
657}
658
659async fn incoming(link: Weak<Inner>, send: quinn::SendStream, mut recv: quinn::RecvStream) {
665 let Some(inner) = link.upgrade() else { return };
666 let Some((open, state)) = read_open(&inner, &mut recv).await else {
667 abandon(send, recv);
668 return;
669 };
670 let offer = inner
671 .lock()
672 .served
673 .get(&(open.realm, open.procedure.clone()))
674 .map(|s| s.offer.clone());
675 let refuse_clear = |code: &str, message: &str, send: quinn::SendStream, recv| {
676 let s = StreamInner::new(inner.clone(), send, open.clone(), false, None);
677 let (code, message) = (code.to_string(), message.to_string());
678 async move { refuse(&s, &code, &message, recv).await }
679 };
680 let place = match admit_stream(&inner, &open) {
683 Ok(place) => place,
684 Err(code) => return refuse_clear(code, "", send, recv).await,
685 };
686 let (session_open, sealing) = match session_open(&inner, &open, offer.as_ref()) {
687 Ok(opened) => opened,
688 Err((code, message)) => return refuse_clear(code, &message, send, recv).await,
689 };
690 let s = StreamInner::new(inner.clone(), send, session_open, false, sealing);
692 let Some(offer) = offer.and_then(|o| o.stream) else {
693 return refuse(&s, CODE_STREAM_NOT_FOUND, "", recv).await;
694 };
695 if Some(offer.mode) != open.mode {
696 return refuse(&s, CODE_MODE_MISMATCH, "", recv).await;
697 }
698 *s.budget() = Some(Budget {
699 admission: inner.admission.clone(),
700 caller: open.caller,
701 place: Some(place),
702 });
703 if !hold_stream(&inner, &s) {
704 StreamInner::end(&s, Some(LinkError::Closed));
705 let _ = recv.stop(0u32.into());
706 return;
707 }
708 tokio::spawn(read(s.clone(), recv, state));
709 tokio::spawn(serve(s, offer));
710}
711
712async fn read_open(
716 inner: &Inner,
717 recv: &mut quinn::RecvStream,
718) -> Option<(VerifiedRequest, StreamState)> {
719 let payload =
720 match tokio::time::timeout(STREAM_OPEN_WAIT, read_frame(recv, STREAM_OPEN_BYTES)).await {
721 Ok(Ok(payload)) => payload,
722 _ => {
723 inner.count("stream_open_unread");
724 return None;
725 }
726 };
727 let v = match cbor::decode(&payload) {
728 Ok(v) if frame_type_of(&v) == "stream_open" => v,
729 _ => {
730 inner.count("stream_open_malformed");
731 return None;
732 }
733 };
734 let Ok(open) = frame::verify_request(&v, inner.profile) else {
735 inner.count("stream_open_unverified");
736 return None;
737 };
738 if open.target != inner.self_id {
739 inner.count("stream_for_another_node");
740 return None;
741 }
742 let state = frame::open_stream(&open).ok()?;
743 Some((open, state))
744}
745
746fn session_open(
752 inner: &Inner,
753 open: &VerifiedRequest,
754 offer: Option<&Offer>,
755) -> Result<(VerifiedRequest, Option<StreamSeal>), (&'static str, String)> {
756 match &open.sealed {
757 Some(_) => opened_request(inner.keyring.as_deref(), open)
758 .map(|(payload, sealed)| {
759 (
760 VerifiedRequest {
761 payload: without_caller(payload),
762 ..open.clone()
763 },
764 Some(StreamSeal::provider(&sealed)),
765 )
766 })
767 .map_err(|detail| (CODE_SEALED_REFUSED, detail)),
768 None if offer
769 .is_some_and(|o| !clear_allowed(o.confidential, inner.keyed_since(o), now_ms())) =>
770 {
771 let message = "this procedure takes sealed opens only";
772 Err((CODE_SEALED_REQUIRED, message.to_string()))
773 }
774 None => Ok((
775 VerifiedRequest {
776 payload: without_caller(open.payload.clone()),
777 ..open.clone()
778 },
779 None,
780 )),
781 }
782}
783
784fn admit_stream(inner: &Inner, open: &VerifiedRequest) -> Result<SessionPlace, &'static str> {
788 match inner.admission.admit(open, &inner.share, now_ms()) {
789 Verdict::Refused(code) => return Err(code),
790 Verdict::Copy(_) => return Err(CODE_REQUEST_COPY),
791 Verdict::New => {}
792 }
793 inner
794 .admission
795 .open_session(open.caller)
796 .ok_or(CODE_TOO_MANY_SESSIONS)
797}
798
799async fn refuse(s: &Arc<StreamInner>, code: &str, message: &str, mut recv: quinn::RecvStream) {
802 s.link.count(&format!("stream_refused_{code}"));
803 let _ = s.abort(code, message).await;
804 let _ = recv.stop(0u32.into());
805}
806
807async fn serve(s: Arc<StreamInner>, offer: StreamOffer) {
811 let stream = Stream { inner: s.clone() };
812 let mut running = tokio::spawn((offer.handler)(stream.clone()));
813 let mut done = s.done_tx.subscribe();
814 let outcome = tokio::select! {
815 outcome = &mut running => outcome,
816 _ = done.wait_for(|ended| *ended) => {
817 running.abort();
818 return;
819 }
820 };
821 match outcome {
822 Ok(Ok(())) => {
823 let _ = stream.close().await;
824 }
825 Ok(Err(e)) => {
826 let _ = stream
827 .abort(CODE_STREAM_HANDLER_ERROR, bounded_detail(&e))
828 .await;
829 }
830 Err(panicked) => {
831 let _ = stream
832 .abort(
833 CODE_STREAM_HANDLER_ERROR,
834 bounded_detail(&panicked.to_string()),
835 )
836 .await;
837 }
838 }
839}
840
841async fn read(s: Arc<StreamInner>, mut recv: quinn::RecvStream, mut state: StreamState) {
845 let mut done = s.done_tx.subscribe();
846 loop {
847 let payload = tokio::select! {
848 _ = done.wait_for(|ended| *ended) => return,
849 payload = read_frame(&mut recv, MAX_FRAME_BYTES) => payload,
850 };
851 let payload = match payload {
852 Ok(payload) => payload,
853 Err(e) => return read_ended(&s, e),
854 };
855 match received(&s, &payload, &state).await {
856 Some(next) => state = next,
857 None => return,
858 }
859 }
860}
861
862fn read_ended(s: &Arc<StreamInner>, e: LinkError) {
865 if s.side().peer_ended {
866 return;
867 }
868 let err = s.link.lock().ended.clone().unwrap_or(e);
869 StreamInner::end(s, Some(err));
870}
871
872async fn received(
875 s: &Arc<StreamInner>,
876 payload: &[u8],
877 state: &StreamState,
878) -> Option<StreamState> {
879 let v = match cbor::decode(payload) {
880 Ok(v) => v,
881 Err(e) => {
882 s.fail("malformed_frame", Some(e.to_string())).await;
883 return None;
884 }
885 };
886 if s.caller && v.get("relay_error").is_some() {
887 relay_failed(s, &v).await;
888 return None;
889 }
890 let verified = if s.caller {
891 frame::verify_provider_stream(&v, state, s.link.profile)
892 } else {
893 frame::verify_caller_stream(&v, state, s.link.profile)
894 };
895 let (verified, next) = match verified {
896 Ok(verified) => verified,
897 Err(e) => {
898 s.fail("malformed_frame", Some(e.to_string())).await;
899 return None;
900 }
901 };
902 let size = payload.len();
903 let fields = match unsealed(verified, s.sealing.as_ref()) {
904 Ok(fields) => fields,
905 Err(e) => {
906 StreamInner::end(s, Some(e));
907 return None;
908 }
909 };
910 if s.caller && settles(&fields, s.sealing.is_some()) {
913 s.side().settled = true;
914 }
915 taken(s, fields, size, next).await
916}
917
918async fn relay_failed(s: &Arc<StreamInner>, v: &Value) {
921 match frame::verify_relay_error(v, &s.open, s.link.profile, &s.link.station.node_id) {
922 Ok(relayed) => s.peer_finished(Some(LinkError::Stream {
923 code: relayed.code,
924 message: String::new(),
925 relay: true,
926 })),
927 Err(e) => s.fail("malformed_frame", Some(e.to_string())).await,
928 }
929}
930
931async fn taken(
934 s: &Arc<StreamInner>,
935 fields: StreamFields,
936 size: usize,
937 next: StreamState,
938) -> Option<StreamState> {
939 match fields {
940 StreamFields::Error { code, message, .. } => {
941 s.peer_finished(Some(LinkError::Stream {
942 code,
943 message,
944 relay: false,
945 }));
946 None
947 }
948 StreamFields::Reply { payload, .. } => {
949 s.deliver(StreamEvent::Reply { payload }, size);
950 s.peer_finished(None);
951 None
952 }
953 StreamFields::End { role, .. } => {
954 s.deliver(StreamEvent::End { role }, size);
955 if role == StreamRole::Both {
956 s.peer_finished(None);
957 return None;
958 }
959 let mine = {
960 let mut side = s.side();
961 side.peer_ended = true;
962 side.sent_end
963 };
964 if mine {
965 StreamInner::end(s, None);
966 }
967 None
968 }
969 StreamFields::Data { encoding, body, .. } => {
970 if !s.deliver(StreamEvent::Data { encoding, body }, size) {
971 s.fail("resource_exhausted", None).await;
972 return None;
973 }
974 Some(next)
975 }
976 StreamFields::SealedData { .. }
978 | StreamFields::SealedError { .. }
979 | StreamFields::SealedReply { .. } => {
980 StreamInner::end(s, Some(LinkError::ClearAnswerToSealed));
981 None
982 }
983 }
984}
985
986fn settles(fields: &StreamFields, sealed: bool) -> bool {
992 match fields {
993 StreamFields::Data { .. } | StreamFields::Reply { .. } => true,
994 StreamFields::End { .. } => !sealed,
995 _ => false,
996 }
997}