use std::collections::VecDeque;
use std::future::Future;
use std::pin::Pin;
use std::sync::{Arc, Mutex, MutexGuard, Weak};
use std::time::Duration;
use tokio::sync::{watch, Notify};
use crate::cbor::{self, Value};
use crate::frame::{
self, RequestSpec, StreamEncoding, StreamFields, StreamMode, StreamRole, StreamState,
VerifiedRequest,
};
use crate::seal::KEY_ID_SIZE;
use super::admission::{Admission, SessionPlace, Verdict};
use super::authorize::authorize;
use super::confidential::{
clear_allowed, opened_request, refused_key, sealed_request, stated, unsealed, Seal, StreamSeal,
CODE_SEALED_REFUSED, CODE_SEALED_REQUIRED,
};
use super::framing::{read_frame, FrameWriter, MAX_FRAME_BYTES};
use super::serve::{
bounded_detail, without_caller, BoxFuture, Offer, StreamOffer, CODE_REQUEST_COPY,
};
use super::{frame_type_of, now_ms, Inner, Link, LinkError};
const STREAM_OPEN_BYTES: usize = 1024 * 1024;
const STREAM_OPEN_WAIT: Duration = Duration::from_secs(10);
const STREAM_INBOX: usize = 16 * 1024 * 1024;
pub const DEFAULT_STREAM_DEADLINE: Duration = Duration::from_secs(30);
const CODE_STREAM_NOT_FOUND: &str = "not_found";
const CODE_MODE_MISMATCH: &str = "mode_mismatch";
const CODE_TOO_MANY_SESSIONS: &str = "too_many_sessions";
const CODE_STREAM_HANDLER_ERROR: &str = "error";
pub type StreamHandler = Arc<dyn Fn(Stream) -> BoxFuture<Result<(), String>> + Send + Sync>;
pub fn stream_handler<F, Fut>(f: F) -> StreamHandler
where
F: Fn(Stream) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<(), String>> + Send + 'static,
{
Arc::new(move |s| Box::pin(f(s)))
}
#[derive(Debug, Clone, PartialEq)]
pub struct StreamCall {
pub realm: [u8; 32],
pub procedure: String,
pub target: [u8; 32],
pub mode: StreamMode,
pub payload: Value,
pub deadline: Duration,
pub token: Option<Vec<u8>>,
pub proofs: Vec<Vec<u8>>,
pub seal: Option<Seal>,
}
impl Default for StreamCall {
fn default() -> Self {
StreamCall {
realm: [0; 32],
procedure: String::new(),
target: [0; 32],
mode: StreamMode::ServerStream,
payload: Value::Map(Vec::new()),
deadline: Duration::ZERO,
token: None,
proofs: Vec::new(),
seal: None,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum StreamEvent {
Data {
encoding: StreamEncoding,
body: Value,
},
End {
role: StreamRole,
},
Reply {
payload: Value,
},
}
#[derive(Clone)]
pub struct Stream {
pub(super) inner: Arc<StreamInner>,
}
pub(crate) type Reseal = Box<
dyn FnOnce(
Option<[u8; KEY_ID_SIZE]>,
) -> Pin<Box<dyn Future<Output = Result<Stream, LinkError>> + Send>>
+ Send,
>;
struct Budget {
admission: Arc<Admission>,
caller: [u8; 32],
place: Option<SessionPlace>,
}
pub(super) struct StreamInner {
link: Arc<Inner>,
writer: FrameWriter,
pub(super) open: VerifiedRequest,
pub(super) caller: bool,
pub(super) sealing: Option<StreamSeal>,
send_seq: tokio::sync::Mutex<u64>,
state: Mutex<StreamSide>,
budget: Mutex<Option<Budget>>,
notify: Notify,
done_tx: watch::Sender<bool>,
reseal: Mutex<Option<Reseal>>,
}
#[derive(Default)]
pub(super) struct StreamSide {
reopening: bool,
successor: Option<Arc<StreamInner>>,
sent_end: bool,
peer_ended: bool,
inbox: VecDeque<(StreamEvent, usize)>,
held: usize,
ended: bool,
err: Option<LinkError>,
pub(super) settled: bool,
}
impl Stream {
pub fn request(&self) -> &VerifiedRequest {
&self.inner.open
}
pub fn sealed(&self) -> bool {
self.inner.sealing.is_some()
}
pub async fn send(&self, body: &[u8]) -> Result<(), LinkError> {
let at = |seq| StreamFields::Data {
seq,
encoding: StreamEncoding::Raw,
body: Value::Bytes(body.to_vec()),
};
self.sent(at, false).await.1
}
pub async fn send_value(&self, v: Value) -> Result<(), LinkError> {
let at = |seq| StreamFields::Data {
seq,
encoding: StreamEncoding::Msgpack,
body: v.clone(),
};
self.sent(at, false).await.1
}
pub async fn close_send(&self) -> Result<(), LinkError> {
let at = |seq| StreamFields::End {
seq,
role: StreamRole::Send,
};
self.sent(at, true).await.1
}
pub async fn close(&self) -> Result<(), LinkError> {
let at = |seq| StreamFields::End {
seq,
role: StreamRole::Both,
};
let (inner, sent) = self.sent(at, true).await;
StreamInner::end(&inner, None);
sent
}
pub async fn reply(&self, payload: Value) -> Result<(), LinkError> {
let at = |seq| StreamFields::Reply {
seq,
payload: payload.clone(),
};
let (inner, sent) = self.sent(at, true).await;
StreamInner::end(&inner, None);
sent
}
pub async fn abort(&self, code: &str, message: &str) -> Result<(), LinkError> {
let at = |seq| StreamFields::Error {
seq,
code: code.to_string(),
message: message.to_string(),
};
let (inner, sent) = self.sent(at, true).await;
StreamInner::end(&inner, Some(aborted(code, message)));
sent
}
pub async fn recv(&self) -> Result<StreamEvent, LinkError> {
loop {
let inner = self.live().await;
let notified = inner.notify.notified();
match inner.next_event() {
Some(outcome) => return outcome,
None => notified.await,
}
}
}
pub async fn done(&self) -> Option<LinkError> {
let mut inner = self.live().await;
while !inner.released().await {
inner = self.live().await;
}
let err = inner.side().err.clone();
err
}
async fn sent(
&self,
at: impl Fn(u64) -> StreamFields,
last: bool,
) -> (Arc<StreamInner>, Result<(), LinkError>) {
let mut inner = self.live().await;
let mut sent = inner.send_here(&at, last).await;
while sent.is_none() {
inner = self.live().await;
sent = inner.send_here(&at, last).await;
}
(inner, sent.unwrap_or(Err(LinkError::StreamClosed)))
}
async fn live(&self) -> Arc<StreamInner> {
let mut at = self.current();
while at.side().reopening {
at.reopen_awaited().await;
at = self.current();
}
at
}
pub(super) fn current(&self) -> Arc<StreamInner> {
let mut at = self.inner.clone();
loop {
let next = at.side().successor.clone();
match next {
Some(next) => at = next,
None => return at,
}
}
}
}
fn aborted(code: &str, message: &str) -> LinkError {
LinkError::Stream {
code: code.to_string(),
message: message.to_string(),
relay: false,
}
}
impl StreamInner {
fn new(
link: Arc<Inner>,
send: quinn::SendStream,
open: VerifiedRequest,
caller: bool,
sealing: Option<StreamSeal>,
) -> Arc<StreamInner> {
Arc::new(StreamInner {
link,
writer: FrameWriter::new(send),
open,
caller,
sealing,
send_seq: tokio::sync::Mutex::new(0),
state: Mutex::new(StreamSide::default()),
budget: Mutex::new(None),
notify: Notify::new(),
done_tx: watch::channel(false).0,
reseal: Mutex::new(None),
})
}
pub(super) fn side(&self) -> MutexGuard<'_, StreamSide> {
self.state.lock().unwrap_or_else(|p| p.into_inner())
}
fn budget(&self) -> MutexGuard<'_, Option<Budget>> {
self.budget.lock().unwrap_or_else(|p| p.into_inner())
}
fn next_event(&self) -> Option<Result<StreamEvent, LinkError>> {
let mut side = self.side();
if let Some((event, size)) = side.inbox.pop_front() {
side.held -= size;
drop(side);
self.release_inbox(size);
return Some(Ok(event));
}
if side.ended {
return Some(Err(side.err.clone().unwrap_or(LinkError::EndOfStream)));
}
None
}
async fn send(
self: &Arc<Self>,
at: impl FnOnce(u64) -> StreamFields,
last: bool,
) -> Result<(), LinkError> {
let seq = self.send_seq.lock().await;
self.send_at(seq, at, last).await
}
async fn send_here(
self: &Arc<Self>,
at: &impl Fn(u64) -> StreamFields,
last: bool,
) -> Option<Result<(), LinkError>> {
let seq = self.send_seq.lock().await;
if self.side().reopening || self.side().successor.is_some() {
return None;
}
Some(self.send_at(seq, at, last).await)
}
async fn send_at(
self: &Arc<Self>,
mut seq: tokio::sync::MutexGuard<'_, u64>,
at: impl FnOnce(u64) -> StreamFields,
last: bool,
) -> Result<(), LinkError> {
if self.side().sent_end {
return Err(LinkError::StreamClosed);
}
let fields = at(*seq);
let plain = match &self.sealing {
Some(sealing) => sealing.plain_of(&fields)?,
None => None,
};
let spent = plain.is_some();
let fields = match (&self.sealing, plain) {
(Some(sealing), Some(plain)) => {
*seq += 1;
let sealed = sealing.sealed(fields, plain);
sealed.inspect_err(|_| self.side().sent_end = true)?
}
_ => fields,
};
let encoded = self
.signed(&fields)
.inspect_err(|_| self.unsent_after_spending(spent))?;
if let Err(e) = self.writer.write(&encoded, MAX_FRAME_BYTES).await {
self.side().sent_end = true;
return Err(e);
}
if !spent {
*seq += 1;
}
if last {
self.sent_last().await;
}
Ok(())
}
fn unsent_after_spending(&self, spent: bool) {
if spent {
self.side().sent_end = true;
}
}
async fn sent_last(self: &Arc<Self>) {
let peer_ended = {
let mut side = self.side();
side.sent_end = true;
side.peer_ended
};
self.writer.finish().await;
if peer_ended {
StreamInner::end(self, None);
}
}
fn signed(&self, fields: &StreamFields) -> Result<Vec<u8>, LinkError> {
let signed = if self.caller {
frame::sign_caller_stream(fields, &self.open, &self.link.key)?
} else {
frame::sign_provider_stream(fields, &self.open, &self.link.key)?
};
cbor::encode(&signed)
.map_err(|e| LinkError::Frame(frame::FrameError::Payload(e.to_string())))
}
async fn abort(self: &Arc<Self>, code: &str, message: &str) -> Result<(), LinkError> {
let sent = self
.send(
|seq| StreamFields::Error {
seq,
code: code.to_string(),
message: message.to_string(),
},
true,
)
.await;
StreamInner::end(self, Some(aborted(code, message)));
sent
}
async fn released(&self) -> bool {
let mut done = self.done_tx.subscribe();
let _ = done.wait_for(|ended| *ended).await;
self.side().successor.is_none()
}
async fn reopen_awaited(&self) {
let notified = self.notify.notified();
if self.side().reopening {
notified.await;
}
}
fn reseal_slot(&self) -> MutexGuard<'_, Option<Reseal>> {
self.reseal.lock().unwrap_or_else(|p| p.into_inner())
}
async fn reopened(self: &Arc<Self>, detail: &str) -> bool {
let Some(reseal) = self.reseal_slot().take() else {
return false;
};
let seq = self.send_seq.lock().await;
if *seq != 0 || self.side().sent_end {
return false;
}
self.side().reopening = true;
drop(seq);
self.notify.notify_waiters();
let s = self.clone();
let named = refused_key(Some(detail));
tokio::spawn(async move {
let outcome = reseal(named).await;
s.adopt(outcome);
});
true
}
fn adopt(self: &Arc<Self>, outcome: Result<Stream, LinkError>) {
let reopened = {
let mut side = self.side();
side.reopening = false;
outcome.map(|next| side.successor = Some(next.inner))
};
match reopened {
Ok(()) => StreamInner::end(self, None),
Err(e) => self.peer_finished(Some(e)),
}
self.notify.notify_waiters();
}
fn deliver(&self, event: StreamEvent, size: usize) -> bool {
if !self.queue(event, size) {
return false;
}
self.notify.notify_one();
true
}
fn queue(&self, event: StreamEvent, size: usize) -> bool {
let mut side = self.side();
if side.held + size > STREAM_INBOX {
return false;
}
if !self.charge_inbox(size) {
return false;
}
side.inbox.push_back((event, size));
side.held += size;
true
}
fn charge_inbox(&self, size: usize) -> bool {
let budget = self.budget();
let Some(budget) = &*budget else {
return true;
};
budget.admission.charge_inbox(budget.caller, size)
}
fn release_inbox(&self, size: usize) {
if let Some(budget) = &*self.budget() {
budget.admission.release_inbox(budget.caller, size);
}
}
fn peer_finished(self: &Arc<Self>, err: Option<LinkError>) {
self.side().peer_ended = true;
StreamInner::end(self, err);
}
async fn fail(self: &Arc<Self>, code: &str, cause: Option<String>) {
let message = cause
.as_deref()
.map(bounded_detail)
.unwrap_or("")
.to_string();
let _ = self.abort(code, &message).await;
}
pub(super) fn end(this: &Arc<StreamInner>, err: Option<LinkError>) {
let Some(graceful) = this.mark_ended(err) else {
return;
};
if !graceful {
StreamInner::reset_sending(this);
}
if let Some(budget) = this.budget().take() {
let held = std::mem::take(&mut this.side().held);
budget.admission.release_inbox(budget.caller, held);
drop(budget.place);
}
this.link
.lock()
.streams
.retain(|w| w.strong_count() > 0 && !std::ptr::eq(w.as_ptr(), Arc::as_ptr(this)));
let _ = this.done_tx.send_replace(true);
this.notify.notify_waiters();
this.notify.notify_one();
}
}
impl StreamInner {
fn mark_ended(&self, err: Option<LinkError>) -> Option<bool> {
let mut side = self.side();
if side.ended {
return None;
}
side.ended = true;
if side.err.is_none() {
side.err = err;
}
let graceful = side.sent_end;
side.sent_end = true;
Some(graceful)
}
fn reset_sending(this: &Arc<StreamInner>) {
let released = this.clone();
let Ok(runtime) = tokio::runtime::Handle::try_current() else {
return;
};
runtime.spawn(async move { released.writer.reset().await });
}
}
fn hold_stream(inner: &Inner, s: &Arc<StreamInner>) -> bool {
let mut state = inner.lock();
if state.ended.is_some() {
return false;
}
state.streams.push(Arc::downgrade(s));
true
}
fn abandon(mut send: quinn::SendStream, mut recv: quinn::RecvStream) {
let _ = send.reset(0u32.into());
let _ = recv.stop(0u32.into());
}
impl Link {
pub async fn open_stream(&self, c: StreamCall) -> Result<Stream, LinkError> {
self.open_stream_resealing(c, None).await
}
pub(crate) async fn open_stream_resealing(
&self,
c: StreamCall,
reseal: Option<Reseal>,
) -> Result<Stream, LinkError> {
let inner = &self.inner;
stated(&c.target, &inner.station.node_id, &c.seal)?;
let deadline = if c.deadline.is_zero() {
DEFAULT_STREAM_DEADLINE
} else {
c.deadline
};
let mut request_id = [0u8; 16];
aws_lc_rs::rand::fill(&mut request_id)
.map_err(|_| LinkError::Io("no randomness".into()))?;
let deadline = (now_ms() + deadline.as_millis() as i64) as u64;
let (sealed, sealing) = sealed_open(inner, &c, request_id, deadline)?;
let signed = frame::sign_stream_open(
&RequestSpec {
request_id,
realm: c.realm,
procedure: c.procedure,
target: c.target,
deadline,
payload: c.payload,
sealed,
mode: Some(c.mode),
token: c.token,
proofs: c.proofs,
source_route: None,
retry_budget: None,
},
&inner.key,
)?;
let encoded = cbor::encode(&signed)
.map_err(|e| LinkError::Frame(frame::FrameError::Payload(e.to_string())))?;
if encoded.len() > STREAM_OPEN_BYTES {
return Err(LinkError::StreamOpenTooLarge(encoded.len()));
}
let open = frame::verify_request(&signed, inner.profile)?;
let state = frame::open_stream(&open)?;
let (send, recv) = inner
.connection
.open_bi()
.await
.map_err(|e| LinkError::Io(format!("open a stream: {e}")))?;
let s = StreamInner::new(inner.clone(), send, open, true, sealing);
let held = hold_stream(inner, &s);
let written = match held {
true => s.writer.write(&encoded, STREAM_OPEN_BYTES).await,
false => Err(inner.lock().ended.clone().unwrap_or(LinkError::Closed)),
};
if let Err(e) = written {
StreamInner::end(&s, Some(e.clone()));
let mut recv = recv;
let _ = recv.stop(0u32.into());
return Err(e);
}
*s.reseal_slot() = reseal;
tokio::spawn(read(s.clone(), recv, state));
Ok(Stream { inner: s })
}
}
fn sealed_open(
inner: &Inner,
c: &StreamCall,
request_id: [u8; 16],
deadline: u64,
) -> Result<(Option<frame::Sealed>, Option<StreamSeal>), LinkError> {
match &c.seal {
Some(Seal::To(key)) => {
let (sealed, s) = sealed_request(
inner.profile,
key,
crate::seal::FRAME_STREAM_OPEN,
c.realm,
&c.procedure,
inner.self_id,
c.target,
request_id,
deadline,
&c.payload,
)?;
Ok((Some(sealed), Some(StreamSeal::caller(&s))))
}
_ => Ok((None, None)),
}
}
pub(super) async fn accept_streams(link: Weak<Inner>) {
let Some(connection) = link.upgrade().map(|l| l.connection.clone()) else {
return;
};
while let Ok((send, recv)) = connection.accept_bi().await {
tokio::spawn(incoming(link.clone(), send, recv));
}
}
async fn incoming(link: Weak<Inner>, send: quinn::SendStream, mut recv: quinn::RecvStream) {
let Some(inner) = link.upgrade() else { return };
let Some((open, state)) = read_open(&inner, &mut recv).await else {
abandon(send, recv);
return;
};
let offer = inner
.lock()
.served
.get(&(open.realm, open.procedure.clone()))
.map(|s| s.offer.clone());
let refuse_clear = |code: &str, message: &str, send: quinn::SendStream, recv| {
let s = StreamInner::new(inner.clone(), send, open.clone(), false, None);
let (code, message) = (code.to_string(), message.to_string());
async move { refuse(&s, &code, &message, recv).await }
};
let place = match admit_stream(&inner, &open) {
Ok(place) => place,
Err(code) => return refuse_clear(code, "", send, recv).await,
};
let (session_open, sealing) = match session_open(&inner, &open, offer.as_ref()) {
Ok(opened) => opened,
Err((code, message)) => return refuse_clear(code, &message, send, recv).await,
};
let s = StreamInner::new(inner.clone(), send, session_open, false, sealing);
let Some((offer, policy)) = offer.and_then(|o| Some((o.stream?, o.policy))) else {
return refuse(&s, CODE_STREAM_NOT_FOUND, "", recv).await;
};
if let Some(code) = authorize(inner.profile, policy.as_ref(), &open) {
return refuse(&s, code, "", recv).await;
}
if Some(offer.mode) != open.mode {
return refuse(&s, CODE_MODE_MISMATCH, "", recv).await;
}
*s.budget() = Some(Budget {
admission: inner.admission.clone(),
caller: open.caller,
place: Some(place),
});
if !hold_stream(&inner, &s) {
StreamInner::end(&s, Some(LinkError::Closed));
let _ = recv.stop(0u32.into());
return;
}
tokio::spawn(read(s.clone(), recv, state));
tokio::spawn(serve(s, offer));
}
async fn read_open(
inner: &Inner,
recv: &mut quinn::RecvStream,
) -> Option<(VerifiedRequest, StreamState)> {
let payload =
match tokio::time::timeout(STREAM_OPEN_WAIT, read_frame(recv, STREAM_OPEN_BYTES)).await {
Ok(Ok(payload)) => payload,
_ => {
inner.count("stream_open_unread");
return None;
}
};
let v = match cbor::decode(&payload) {
Ok(v) if frame_type_of(&v) == "stream_open" => v,
_ => {
inner.count("stream_open_malformed");
return None;
}
};
let Ok(open) = frame::verify_request(&v, inner.profile) else {
inner.count("stream_open_unverified");
return None;
};
if open.target != inner.self_id {
inner.count("stream_for_another_node");
return None;
}
let state = frame::open_stream(&open).ok()?;
Some((open, state))
}
fn session_open(
inner: &Inner,
open: &VerifiedRequest,
offer: Option<&Offer>,
) -> Result<(VerifiedRequest, Option<StreamSeal>), (&'static str, String)> {
match &open.sealed {
Some(_) => opened_request(inner.keyring.as_deref(), open)
.map(|(payload, sealed)| {
(
VerifiedRequest {
payload: without_caller(payload),
..open.clone()
},
Some(StreamSeal::provider(&sealed)),
)
})
.map_err(|detail| (CODE_SEALED_REFUSED, detail)),
None if offer
.is_some_and(|o| !clear_allowed(o.confidential, inner.keyed_since(o), now_ms())) =>
{
let message = "this procedure takes sealed opens only";
Err((CODE_SEALED_REQUIRED, message.to_string()))
}
None => Ok((
VerifiedRequest {
payload: without_caller(open.payload.clone()),
..open.clone()
},
None,
)),
}
}
fn admit_stream(inner: &Inner, open: &VerifiedRequest) -> Result<SessionPlace, &'static str> {
match inner.admission.admit(open, &inner.share, now_ms()) {
Verdict::Refused(code) => return Err(code),
Verdict::Copy(_) => return Err(CODE_REQUEST_COPY),
Verdict::New => {}
}
inner
.admission
.open_session(open.caller)
.ok_or(CODE_TOO_MANY_SESSIONS)
}
async fn refuse(s: &Arc<StreamInner>, code: &str, message: &str, mut recv: quinn::RecvStream) {
s.link.count(&format!("stream_refused_{code}"));
let _ = s.abort(code, message).await;
let _ = recv.stop(0u32.into());
}
async fn serve(s: Arc<StreamInner>, offer: StreamOffer) {
let stream = Stream { inner: s.clone() };
let mut running = tokio::spawn((offer.handler)(stream.clone()));
let mut done = s.done_tx.subscribe();
let outcome = tokio::select! {
outcome = &mut running => outcome,
_ = done.wait_for(|ended| *ended) => {
running.abort();
return;
}
};
match outcome {
Ok(Ok(())) => {
let _ = stream.close().await;
}
Ok(Err(e)) => {
let _ = stream
.abort(CODE_STREAM_HANDLER_ERROR, bounded_detail(&e))
.await;
}
Err(panicked) => {
let _ = stream
.abort(
CODE_STREAM_HANDLER_ERROR,
bounded_detail(&panicked.to_string()),
)
.await;
}
}
}
async fn read(s: Arc<StreamInner>, mut recv: quinn::RecvStream, mut state: StreamState) {
let mut done = s.done_tx.subscribe();
loop {
let payload = tokio::select! {
_ = done.wait_for(|ended| *ended) => return,
payload = read_frame(&mut recv, MAX_FRAME_BYTES) => payload,
};
let payload = match payload {
Ok(payload) => payload,
Err(e) => return read_ended(&s, e),
};
match received(&s, &payload, &state).await {
Some(next) => state = next,
None => return,
}
}
}
fn read_ended(s: &Arc<StreamInner>, e: LinkError) {
if s.side().peer_ended {
return;
}
let err = s.link.lock().ended.clone().unwrap_or(e);
StreamInner::end(s, Some(err));
}
async fn received(
s: &Arc<StreamInner>,
payload: &[u8],
state: &StreamState,
) -> Option<StreamState> {
let v = match cbor::decode(payload) {
Ok(v) => v,
Err(e) => {
s.fail("malformed_frame", Some(e.to_string())).await;
return None;
}
};
if s.caller && v.get("relay_error").is_some() {
relay_failed(s, &v).await;
return None;
}
let verified = if s.caller {
frame::verify_provider_stream(&v, state, s.link.profile)
} else {
frame::verify_caller_stream(&v, state, s.link.profile)
};
let (verified, next) = match verified {
Ok(verified) => verified,
Err(e) => {
s.fail("malformed_frame", Some(e.to_string())).await;
return None;
}
};
let size = payload.len();
let fields = match unsealed(verified, s.sealing.as_ref()) {
Ok(fields) => fields,
Err(e) => {
StreamInner::end(s, Some(e));
return None;
}
};
if s.caller && settles(&fields, s.sealing.is_some()) {
s.side().settled = true;
}
taken(s, fields, size, next).await
}
async fn relay_failed(s: &Arc<StreamInner>, v: &Value) {
match frame::verify_relay_error(v, &s.open, s.link.profile, &s.link.station.node_id) {
Ok(relayed) => s.peer_finished(Some(LinkError::Stream {
code: relayed.code,
message: String::new(),
relay: true,
})),
Err(e) => s.fail("malformed_frame", Some(e.to_string())).await,
}
}
async fn taken(
s: &Arc<StreamInner>,
fields: StreamFields,
size: usize,
next: StreamState,
) -> Option<StreamState> {
match fields {
StreamFields::Error {
seq: 0,
ref code,
ref message,
} if s.caller && code == CODE_SEALED_REFUSED && s.reopened(message).await => None,
StreamFields::Error { code, message, .. } => {
s.peer_finished(Some(LinkError::Stream {
code,
message,
relay: false,
}));
None
}
StreamFields::Reply { payload, .. } => {
s.deliver(StreamEvent::Reply { payload }, size);
s.peer_finished(None);
None
}
StreamFields::End { role, .. } => {
s.deliver(StreamEvent::End { role }, size);
if role == StreamRole::Both {
s.peer_finished(None);
return None;
}
let mine = {
let mut side = s.side();
side.peer_ended = true;
side.sent_end
};
if mine {
StreamInner::end(s, None);
}
None
}
StreamFields::Data { encoding, body, .. } => {
if !s.deliver(StreamEvent::Data { encoding, body }, size) {
s.fail("resource_exhausted", None).await;
return None;
}
Some(next)
}
StreamFields::SealedData { .. }
| StreamFields::SealedError { .. }
| StreamFields::SealedReply { .. } => {
StreamInner::end(s, Some(LinkError::ClearAnswerToSealed));
None
}
}
}
fn settles(fields: &StreamFields, sealed: bool) -> bool {
match fields {
StreamFields::Data { .. } | StreamFields::Reply { .. } => true,
StreamFields::End { .. } => !sealed,
_ => false,
}
}