pub struct StreamEndpoint<'a> { /* private fields */ }Expand description
The per-peer streaming registry (SPEC §3): it multiplexes many concurrent streams to ONE peer,
bounds their count (MAX_CONCURRENT_STREAMS), seals every outbound frame with a fresh ephemeral,
and opens + fully verifies every inbound frame — RESETting a stream on any failed verify (gate items
DIG-Network/dig_ecosystem#1162).
One endpoint’s identity key both SEALS our outbound frames (as sender) and OPENS the peer’s inbound frames (as recipient) — the ONE BLS12-381 identity key (SPEC §5.1).
Implementations§
Source§impl<'a> StreamEndpoint<'a>
impl<'a> StreamEndpoint<'a>
Sourcepub fn new(
identity_sk: &'a SecretKey,
local_did: Bytes32,
local_epoch: u32,
peer_did: Bytes32,
peer_pub: &'a [u8; 48],
message_type: u32,
) -> Self
pub fn new( identity_sk: &'a SecretKey, local_did: Bytes32, local_epoch: u32, peer_did: Bytes32, peer_pub: &'a [u8; 48], message_type: u32, ) -> Self
A new endpoint for streams to one peer, using MAX_CONCURRENT_STREAMS as the cap.
identity_sk is our ONE BLS12-381 identity key; peer_pub is the peer’s 48-byte G1 identity
key. message_type is the SPEC §4 type every frame carries.
Sourcepub fn with_max_concurrent(self, max: usize) -> Self
pub fn with_max_concurrent(self, max: usize) -> Self
Override the per-peer concurrent-stream cap (SPEC §3), e.g. for a higher-capacity relay endpoint
or a tighter test boundary. Defaults to MAX_CONCURRENT_STREAMS.
Sourcepub fn stream_count(&self) -> usize
pub fn stream_count(&self) -> usize
The number of live streams to this peer.
Sourcepub fn session(&self, correlation_id: Bytes32) -> Option<&StreamSession>
pub fn session(&self, correlation_id: Bytes32) -> Option<&StreamSession>
Read-only access to a live session (for inspection/testing).
Sourcepub fn open(
&mut self,
correlation_id: Bytes32,
recv_window: u32,
now_ms: u64,
expires_at: u64,
) -> Result<DigMessageEnvelope>
pub fn open( &mut self, correlation_id: Bytes32, recv_window: u32, now_ms: u64, expires_at: u64, ) -> Result<DigMessageEnvelope>
Open a new outbound stream, advertising recv_window credits for the peer→us direction (SPEC §3).
Returns the sealed OPEN envelope to send.
§Errors
MessageError::StreamLimit if we already hold MAX_CONCURRENT_STREAMS streams;
MessageError::StreamProtocol if correlation_id is already in use; plus any seal error.
Sourcepub fn open_ack(
&mut self,
correlation_id: Bytes32,
recv_window: u32,
now_ms: u64,
expires_at: u64,
) -> Result<DigMessageEnvelope>
pub fn open_ack( &mut self, correlation_id: Bytes32, recv_window: u32, now_ms: u64, expires_at: u64, ) -> Result<DigMessageEnvelope>
Complete a responder handshake, advertising recv_window credits for the peer→us direction.
Returns the sealed OPEN_ACK envelope.
§Errors
MessageError::StreamProtocol for an unknown stream or an illegal handshake position; plus any
seal error.
Sourcepub fn send_data(
&mut self,
correlation_id: Bytes32,
payload: &[u8],
now_ms: u64,
expires_at: u64,
) -> Result<DigMessageEnvelope>
pub fn send_data( &mut self, correlation_id: Bytes32, payload: &[u8], now_ms: u64, expires_at: u64, ) -> Result<DigMessageEnvelope>
Send one DATA chunk on an established stream, consuming one credit (SPEC §3). Returns the sealed DATA envelope.
§Errors
MessageError::StreamProtocol for an unknown stream, a closed direction, or exhausted credit
(backpressure); plus any seal/compression error.
Sourcepub fn grant_credit(
&mut self,
correlation_id: Bytes32,
n: u32,
now_ms: u64,
expires_at: u64,
) -> Result<DigMessageEnvelope>
pub fn grant_credit( &mut self, correlation_id: Bytes32, n: u32, now_ms: u64, expires_at: u64, ) -> Result<DigMessageEnvelope>
Grant the peer n more DATA credits (SPEC §3 backpressure). Returns the sealed CREDIT envelope.
§Errors
MessageError::StreamProtocol for an unknown or not-yet-established stream; plus any seal error.
Sourcepub fn close(
&mut self,
correlation_id: Bytes32,
now_ms: u64,
expires_at: u64,
) -> Result<DigMessageEnvelope>
pub fn close( &mut self, correlation_id: Bytes32, now_ms: u64, expires_at: u64, ) -> Result<DigMessageEnvelope>
Half-close our sending direction (SPEC §3). Returns the sealed CLOSE envelope; the stream is dropped once BOTH directions are closed.
§Errors
MessageError::StreamProtocol for an unknown or already-closed direction; plus any seal error.
Sourcepub fn reset(
&mut self,
correlation_id: Bytes32,
now_ms: u64,
expires_at: u64,
) -> Result<DigMessageEnvelope>
pub fn reset( &mut self, correlation_id: Bytes32, now_ms: u64, expires_at: u64, ) -> Result<DigMessageEnvelope>
Abort a stream immediately (SPEC §3 cancel). Returns the sealed RESET envelope and drops the stream.
§Errors
MessageError::StreamProtocol for an unknown stream; plus any seal error.
Sourcepub fn accept(
&mut self,
envelope: &DigMessageEnvelope,
resolve_sender_pub: impl Fn(Bytes32, u32) -> Option<[u8; 48]>,
now_ms: u64,
) -> Result<StreamAccept>
pub fn accept( &mut self, envelope: &DigMessageEnvelope, resolve_sender_pub: impl Fn(Bytes32, u32) -> Option<[u8; 48]>, now_ms: u64, ) -> Result<StreamAccept>
Open + fully verify an inbound frame and drive the state machine (SPEC §3). On success returns a
verified StreamEvent; on ANY failed verify or protocol violation returns
StreamAccept::Dropped — a bad/unauthenticated/non-actionable frame silently discarded (a live
session is left untouched); or, ONLY for a state-machine violation by the AUTHENTICATED peer on a
KNOWN stream (or the concurrent-stream cap), StreamAccept::Reset — a sealed RESET the caller
sends. A RESET is NEVER emitted for unauthenticated or duplicate input, so the untrusted relay
(§5.4) cannot provoke a RESET reflection storm or replay-teardown (SPEC §3, gate item #1162).
resolve_sender_pub maps (peer DID, epoch) to the peer’s 48-byte G1 key (usually the endpoint’s
own peer_pub); now_ms is our wall clock for the freshness + expiry checks.
§Errors
Only a failure to SEAL the RESET response (on the authenticated-violation path) propagates as
Err; a rejected inbound frame is the non-error StreamAccept::Dropped/StreamAccept::Reset
outcome.
Auto Trait Implementations§
impl<'a> Freeze for StreamEndpoint<'a>
impl<'a> RefUnwindSafe for StreamEndpoint<'a>
impl<'a> Send for StreamEndpoint<'a>
impl<'a> Sync for StreamEndpoint<'a>
impl<'a> Unpin for StreamEndpoint<'a>
impl<'a> UnsafeUnpin for StreamEndpoint<'a>
impl<'a> UnwindSafe for StreamEndpoint<'a>
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> FmtForward for T
impl<T> FmtForward for T
Source§fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
self to use its Binary implementation when Debug-formatted.Source§fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
self to use its Display implementation when
Debug-formatted.Source§fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
self to use its LowerExp implementation when
Debug-formatted.Source§fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
self to use its LowerHex implementation when
Debug-formatted.Source§fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
self to use its Octal implementation when Debug-formatted.Source§fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
self to use its Pointer implementation when
Debug-formatted.Source§fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
self to use its UpperExp implementation when
Debug-formatted.Source§fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
self to use its UpperHex implementation when
Debug-formatted.Source§impl<T> Pipe for Twhere
T: ?Sized,
impl<T> Pipe for Twhere
T: ?Sized,
Source§fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
Source§fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
self and passes that borrow into the pipe function. Read moreSource§fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
self and passes that borrow into the pipe function. Read moreSource§fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
Source§fn pipe_borrow_mut<'a, B, R>(
&'a mut self,
func: impl FnOnce(&'a mut B) -> R,
) -> R
fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
Source§fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
self, then passes self.as_ref() into the pipe function.Source§fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
self, then passes self.as_mut() into the pipe
function.Source§fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
self, then passes self.deref() into the pipe function.impl<T> Read<Exclusive, BecauseExclusive> for Twhere
T: ?Sized,
Source§impl<T> Tap for T
impl<T> Tap for T
Source§fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
Borrow<B> of a value. Read moreSource§fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
BorrowMut<B> of a value. Read moreSource§fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
AsRef<R> view of a value. Read moreSource§fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
AsMut<R> view of a value. Read moreSource§fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
Deref::Target of a value. Read moreSource§fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
Deref::Target of a value. Read moreSource§fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
.tap() only in debug builds, and is erased in release builds.Source§fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
.tap_mut() only in debug builds, and is erased in release
builds.Source§fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
.tap_borrow() only in debug builds, and is erased in release
builds.Source§fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
.tap_borrow_mut() only in debug builds, and is erased in release
builds.Source§fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
.tap_ref() only in debug builds, and is erased in release
builds.Source§fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
.tap_ref_mut() only in debug builds, and is erased in release
builds.Source§fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
.tap_deref() only in debug builds, and is erased in release
builds.