use super::{
CutPoint, MAX_STEPS, OVERSIZED_MESSAGE, Outbox, SCRIPT_HANDSHAKE_TIMEOUT, TIMESTAMP, frame,
self_attestation, unframe, would_block,
};
use crate::transport::Read;
use crate::transport::handshake;
use crate::transport::mock::payload;
use crate::transport::sealing;
use crate::transport::{
Attestation, CRYPTO_DOMAIN_WIRE, CRYPTO_DOMAIN_WIRE_ARK_TO_HOST,
CRYPTO_DOMAIN_WIRE_HOST_TO_ARK, Error, Event, MAX_FRAME_SIZE, Sender,
};
use darkbio_crypto::{cbor, cose, xdsa, xhpke};
use std::collections::VecDeque;
use std::fmt;
use std::io;
use std::sync::{Arc, Mutex};
use std::time::Instant;
const PROBE_ID: u64 = u64::MAX;
#[derive(Clone, Debug, PartialEq, Eq)]
#[cfg_attr(feature = "fuzz", derive(arbitrary::Arbitrary))]
pub enum Step {
Reset,
ResetPair,
Hello,
HelloReplay,
HelloBadKey,
Ack,
AckReplay,
AckTampered,
AckBadAuth,
AckBadSigner,
AckBadPayload,
AckBadEncap,
Request(u8),
RequestReplay,
RequestTampered,
Garbage,
Junk(Vec<u8>),
Truncated(u8),
Partial,
Oversized,
Retain,
Send(u8),
SendRetained(u8),
SendOversized,
Disconnect,
Chunk(u8),
Batch(u8),
Yield,
Interrupt,
ReadTimeout,
Break,
Heal,
Cut { point: CutPoint, then_broken: bool },
Timeout(CutPoint),
}
impl Step {
fn is_action(&self) -> bool {
matches!(
self,
Self::Retain
| Self::Send(_)
| Self::SendRetained(_)
| Self::SendOversized
| Self::Disconnect
)
}
fn queues_frames(&self) -> bool {
!self.is_action()
&& !matches!(
self,
Step::Yield
| Step::Interrupt
| Step::Break
| Step::Heal
| Step::Cut { .. }
| Step::Chunk(_)
| Step::Batch(_)
| Step::Timeout(_)
| Step::ReadTimeout
)
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub enum State {
#[default]
Idle,
AwaitHello,
AwaitAck,
Established,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Summary {
pub state: State, pub dropped: usize, pub fragments: usize, pub handshakes: usize, pub delivered: usize, pub replies: usize, pub reads: usize, }
enum Frame {
Empty,
Hello(Box<Keys>),
Ack,
Request(u64),
Garbage,
Junk,
}
enum Partial {
None,
Hello(Box<Keys>),
Junk,
}
enum Emit {
Dropped,
Fragment,
ArkHello(Box<Keys>),
Reply(u64, Arc<Mutex<xhpke::Receiver>>),
}
impl fmt::Debug for Emit {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Emit::Dropped => write!(f, "Dropped"),
Emit::Fragment => write!(f, "Fragment"),
Emit::ArkHello(_) => write!(f, "ArkHello"),
Emit::Reply(id, _) => write!(f, "Reply({id})"),
}
}
}
enum Tail {
Fragment,
Body(Emit),
}
enum Payload {
Frame(Emit),
Signal,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Outcome {
Absorbed,
Message(u64),
Garbage,
Ended,
Opened,
Yield,
SendFailed,
Terminated,
}
#[derive(Clone)]
struct Keys {
signer: xdsa::SecretKey,
crypto: xhpke::SecretKey,
}
impl Keys {
fn generate() -> Self {
Self {
signer: xdsa::SecretKey::generate(),
crypto: xhpke::SecretKey::generate(),
}
}
fn hello(&self) -> Vec<u8> {
cbor::encode(&handshake::HostHello {
host_signer: self.signer.public_key(),
host_crypto: self.crypto.public_key(),
})
.unwrap()
}
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum AckFlaw {
Tampered,
Auth,
Signer,
Payload,
Encap,
}
struct Pending {
keys: Keys,
ark_crypto: xhpke::PublicKey,
receiver: xhpke::Receiver,
}
impl Pending {
fn ack(self, identity: &xdsa::PublicKey) -> (Vec<u8>, xhpke::Sender, xhpke::Receiver) {
let (sender, encap) = self
.ark_crypto
.new_sender(CRYPTO_DOMAIN_WIRE_HOST_TO_ARK)
.unwrap();
let ack = cose::seal_at(
&handshake::HostAck {
h2a_encap: encap.to_vec(),
},
&handshake::HostAckAuth {
ark_signer: identity.clone(),
ark_crypto: self.ark_crypto.clone(),
},
&self.keys.signer,
&self.ark_crypto,
CRYPTO_DOMAIN_WIRE,
TIMESTAMP,
)
.unwrap();
(ack, sender, self.receiver)
}
fn bad_ack(&self, identity: &xdsa::PublicKey, flaw: AckFlaw) -> Vec<u8> {
let (_, encap) = self
.ark_crypto
.new_sender(CRYPTO_DOMAIN_WIRE_HOST_TO_ARK)
.unwrap();
let auth = handshake::HostAckAuth {
ark_signer: identity.clone(),
ark_crypto: match flaw {
AckFlaw::Auth => xhpke::SecretKey::generate().public_key(),
_ => self.ark_crypto.clone(),
},
};
let stranger = xdsa::SecretKey::generate();
let signer = match flaw {
AckFlaw::Signer => &stranger,
_ => &self.keys.signer,
};
let mut sealed = match flaw {
AckFlaw::Payload => cose::seal_at(
&(vec![1u8], vec![2u8]),
&auth,
signer,
&self.ark_crypto,
CRYPTO_DOMAIN_WIRE,
TIMESTAMP,
),
_ => cose::seal_at(
&handshake::HostAck {
h2a_encap: match flaw {
AckFlaw::Encap => vec![0x42; 3],
_ => encap.to_vec(),
},
},
&auth,
signer,
&self.ark_crypto,
CRYPTO_DOMAIN_WIRE,
TIMESTAMP,
),
}
.unwrap();
if flaw == AckFlaw::Tampered {
*sealed.last_mut().unwrap() ^= 0xff;
}
sealed
}
}
pub struct Client {
steps: VecDeque<Step>,
action: Option<Step>, identity: xdsa::PublicKey, outbox: Outbox, bytes: Vec<u8>, offset: usize, chunk: usize, batch: usize, broken: bool, cut: Option<CutPoint>, timeout: bool, timed_out: bool,
state: State, partial: Partial, resync: bool, tail: Option<Tail>, emits: Vec<Emit>, outcome: Outcome, held: bool,
pending: Option<Pending>, sender: Option<xhpke::Sender>, receiver: Option<Arc<Mutex<xhpke::Receiver>>>,
last_hello: Option<(Vec<u8>, Keys)>, last_ack: Option<Vec<u8>>, last_request: Option<Vec<u8>>, last_valid: Option<Vec<u8>>,
summary: Summary,
}
impl Client {
fn new(steps: &[Step], identity: xdsa::PublicKey, outbox: Outbox) -> Self {
Self {
steps: steps.iter().take(MAX_STEPS).cloned().collect(),
action: None,
identity,
outbox,
bytes: Vec::new(),
offset: 0,
chunk: 0,
batch: 0,
broken: false,
cut: None,
timeout: false,
timed_out: false,
state: State::Idle,
partial: Partial::None,
resync: false,
tail: None,
emits: Vec::new(),
outcome: Outcome::Absorbed,
held: false,
pending: None,
sender: None,
receiver: None,
last_hello: None,
last_ack: None,
last_request: None,
last_valid: None,
summary: Summary::default(),
}
}
fn next_action(&mut self) -> Option<Step> {
if self.action.is_some() {
return self.action.take();
}
if self.bytes.is_empty() && self.steps.front().is_some_and(Step::is_action) {
self.sync();
return self.steps.pop_front();
}
None
}
fn execute(&mut self, step: Step) {
match step {
Step::Reset => {
self.deliver(Frame::Empty);
self.bytes.push(0x00);
}
Step::ResetPair => {
self.deliver(Frame::Empty);
self.deliver(Frame::Empty);
self.bytes.extend([0x00, 0x00]);
}
Step::Hello => {
let keys = Keys::generate();
let framed = frame(&keys.hello());
self.record(&framed);
self.last_hello = Some((framed.clone(), keys.clone()));
self.deliver(Frame::Hello(Box::new(keys)));
self.bytes.extend(framed);
}
Step::HelloReplay => {
if let Some((framed, keys)) = self.last_hello.clone() {
self.deliver(Frame::Hello(Box::new(keys)));
self.bytes.extend(framed);
}
}
Step::HelloBadKey => {
let signer = xdsa::SecretKey::generate().public_key().to_bytes().to_vec();
let hello = cbor::encode(&(signer, vec![0xffu8; xhpke::PUBLIC_KEY_SIZE])).unwrap();
self.junk(&hello);
}
Step::Ack => match self.pending.take() {
Some(pending) => {
let (ack, sender, receiver) = pending.ack(&self.identity);
let framed = frame(&ack);
self.record(&framed);
self.last_ack = Some(framed.clone());
self.sender = Some(sender);
self.receiver = Some(Arc::new(Mutex::new(receiver)));
self.deliver(Frame::Ack);
self.bytes.extend(framed);
}
None => self.junk(b"ack without a pending server hello"),
},
Step::AckReplay => {
if let Some(framed) = self.last_ack.clone() {
self.deliver(Frame::Junk);
self.bytes.extend(framed);
}
}
Step::AckTampered => self.bad_ack(AckFlaw::Tampered),
Step::AckBadAuth => self.bad_ack(AckFlaw::Auth),
Step::AckBadSigner => self.bad_ack(AckFlaw::Signer),
Step::AckBadPayload => self.bad_ack(AckFlaw::Payload),
Step::AckBadEncap => self.bad_ack(AckFlaw::Encap),
Step::Request(tag) => match self.sender.as_mut() {
Some(sender) => {
let id = tag as u64;
let packet = sealing::seal(sender, &payload(id)).unwrap();
let framed = frame(&packet);
self.record(&framed);
self.last_request = Some(framed.clone());
self.deliver(Frame::Request(id));
self.bytes.extend(framed);
}
None => self.junk(b"request without a session"),
},
Step::RequestReplay => {
if let Some(framed) = self.last_request.clone() {
self.deliver(Frame::Junk);
self.bytes.extend(framed);
}
}
Step::RequestTampered => match self.sender.as_mut() {
Some(sender) => {
let mut packet = sealing::seal(sender, &payload(0)).unwrap();
*packet.last_mut().unwrap() ^= 0xff;
self.junk(&packet);
}
None => self.junk(b"tampered request without a session"),
},
Step::Garbage => match self.sender.as_mut() {
Some(sender) => {
let packet = sender.seal(&[0x07], &[]).unwrap();
let framed = frame(&packet);
self.record(&framed);
self.deliver(Frame::Garbage);
self.bytes.extend(framed);
}
None => self.junk(b"garbage without a session"),
},
Step::Junk(mut junk) => {
for byte in junk.iter_mut() {
if *byte == 0 {
*byte = 1;
}
}
if junk.is_empty() {
junk.push(1);
}
self.deliver(Frame::Junk);
self.bytes.extend(junk);
self.bytes.push(0x00);
}
Step::Truncated(n) => {
if let Some(valid) = self.last_valid.clone() {
let keep = match valid.len() {
0..=1 => 1,
len => 1 + n as usize % (len - 1),
};
self.deliver(Frame::Junk);
self.bytes.extend(&valid[..keep]);
self.bytes.push(0x00);
}
}
Step::Partial => {
let keys = Keys::generate();
let mut framed = frame(&keys.hello());
framed.pop();
self.bytes.extend(framed);
self.partial = match self.partial {
Partial::None => Partial::Hello(Box::new(keys)),
_ => Partial::Junk,
};
}
Step::Oversized => {
self.partial = Partial::None;
self.deliver(Frame::Junk);
self.bytes.resize(self.bytes.len() + MAX_FRAME_SIZE + 1, 1);
self.bytes.push(0x00);
}
Step::Chunk(n) => self.chunk = n as usize,
Step::Batch(n) => self.batch = n as usize,
Step::Break => self.set_broken(true),
Step::Heal => self.set_broken(false),
Step::Cut { point, then_broken } => {
self.cut = Some(point);
self.timeout = false;
self.outbox.set_cut(point);
if then_broken {
self.set_broken(true);
}
}
Step::Timeout(point) => {
self.cut = Some(point);
self.timeout = true;
self.outbox.set_timeout(point);
}
Step::Yield
| Step::Interrupt
| Step::ReadTimeout
| Step::Retain
| Step::Send(_)
| Step::SendRetained(_)
| Step::SendOversized
| Step::Disconnect => {
unreachable!("control steps are handled by the reader and driver")
}
}
}
fn bad_ack(&mut self, flaw: AckFlaw) {
match self.pending.as_ref() {
Some(pending) => {
let ack = pending.bad_ack(&self.identity, flaw);
self.junk(&ack);
}
None => self.junk(b"bad ack without a pending server hello"),
}
}
fn junk(&mut self, text: &[u8]) {
self.deliver(Frame::Junk);
self.bytes.extend(frame(text));
}
fn record(&mut self, framed: &[u8]) {
self.last_valid = Some(framed[..framed.len() - 1].to_vec());
}
fn set_broken(&mut self, broken: bool) {
self.outbox.set_broken(broken);
self.broken = broken;
}
fn deliver(&mut self, frame: Frame) {
let frame = match std::mem::replace(&mut self.partial, Partial::None) {
Partial::None => frame,
Partial::Hello(keys) => match frame {
Frame::Empty => Frame::Hello(keys),
_ => Frame::Junk,
},
Partial::Junk => Frame::Junk,
};
match (self.state, frame) {
(_, Frame::Empty) => {
if std::mem::take(&mut self.held) {
self.outcome = Outcome::Ended;
}
self.forget();
self.state = State::AwaitHello;
}
(State::AwaitHello, Frame::Hello(keys)) => {
if self.send(Payload::Frame(Emit::ArkHello(keys))) {
self.state = State::AwaitAck;
} else {
self.forget();
self.state = State::Idle;
if !self.timed_out {
self.send(Payload::Signal);
}
self.outcome = Outcome::SendFailed;
}
}
(State::AwaitAck, Frame::Ack) => {
self.state = State::Established;
self.held = true;
self.outcome = Outcome::Opened;
}
(State::Established, Frame::Request(id)) => {
self.outcome = Outcome::Message(id);
}
(State::Established, Frame::Garbage) => {
self.outcome = Outcome::Garbage;
}
_ => {
self.forget();
self.state = State::Idle;
self.send(Payload::Signal);
if std::mem::take(&mut self.held) {
self.outcome = Outcome::Ended;
}
}
}
}
fn send(&mut self, payload: Payload) -> bool {
let frame = matches!(&payload, Payload::Frame(_));
let cut = match self.cut {
Some(CutPoint::Start) => self.cut.take(),
Some(CutPoint::Middle(_)) if frame => self.cut.take(),
Some(CutPoint::Delimiter | CutPoint::Flush) if frame || self.resync => self.cut.take(),
_ => None,
};
self.timed_out = cut.is_some() && std::mem::take(&mut self.timeout);
let sent = match cut {
Some(CutPoint::Start) => false,
None if self.broken => false,
_ => {
if self.resync {
self.zero_out();
}
match (payload, cut) {
(Payload::Frame(_), Some(CutPoint::Middle(n))) => {
if !self.resync || n != 0 {
self.tail = Some(Tail::Fragment);
}
false
}
(Payload::Frame(emit), Some(CutPoint::Delimiter)) => {
self.tail = Some(Tail::Body(emit));
false
}
(Payload::Frame(emit), _) => {
self.emits.push(emit);
cut != Some(CutPoint::Flush)
}
(Payload::Signal, Some(CutPoint::Delimiter)) => false,
(Payload::Signal, _) => {
self.zero_out();
cut != Some(CutPoint::Flush)
}
}
}
};
self.resync = !sent;
sent
}
fn zero_out(&mut self) {
self.emits.push(match self.tail.take() {
None => Emit::Dropped,
Some(Tail::Fragment) => Emit::Fragment,
Some(Tail::Body(emit)) => emit,
});
}
fn forget(&mut self) {
self.sender = None;
self.receiver = None;
self.pending = None;
}
fn interrupt(&mut self, outcome: Outcome) {
if matches!(self.state, State::AwaitHello | State::AwaitAck) {
self.forget();
self.state = State::Idle;
}
self.outcome = outcome;
}
fn surfaced(&mut self, outcome: Outcome) {
assert_eq!(self.outcome, outcome, "model vs server");
self.outcome = Outcome::Absorbed;
}
fn sync(&mut self) {
assert_eq!(
self.outcome,
Outcome::Absorbed,
"server read on past a step it should have surfaced"
);
let frames = self.outbox.take_frames();
let emits = std::mem::take(&mut self.emits);
assert_eq!(frames.len(), emits.len(), "model expected {emits:?}");
assert_eq!(
self.outbox.has_tail(),
self.tail.is_some(),
"unterminated frame"
);
for (frame, emit) in frames.iter().zip(emits) {
match emit {
Emit::Dropped => {
assert!(
frame.is_empty(),
"expected an empty frame, server emitted {} bytes",
frame.len()
);
self.summary.dropped += 1;
}
Emit::Fragment => {
assert!(
!frame.is_empty(),
"expected a cut frame, server emitted an empty one"
);
self.summary.fragments += 1;
}
Emit::ArkHello(keys) => {
self.receive_hello(frame, *keys);
self.summary.handshakes += 1;
}
Emit::Reply(id, receiver) => {
let opened = sealing::open(&mut receiver.lock().unwrap(), &unframe(frame))
.expect("reply failed to open");
assert_eq!(opened, payload(id));
self.summary.replies += 1;
}
}
}
}
fn receive_hello(&mut self, frame: &[u8], keys: Keys) {
let auth = handshake::ArkHelloAuth {
host_signer: keys.signer.public_key(),
host_crypto: keys.crypto.public_key(),
};
let sign1 = cose::decrypt(&unframe(frame), &auth, &keys.crypto, CRYPTO_DOMAIN_WIRE)
.expect("server hello failed to decrypt");
let hello: handshake::ArkHello =
cose::verify(&sign1, &auth, &self.identity, CRYPTO_DOMAIN_WIRE, None)
.expect("server hello signature invalid");
let encap: [u8; xhpke::ENCAP_KEY_SIZE] = hello
.a2h_encap
.try_into()
.expect("server hello encap size invalid");
let receiver = keys
.crypto
.new_receiver(&encap, CRYPTO_DOMAIN_WIRE_ARK_TO_HOST)
.unwrap();
if self.state == State::AwaitAck {
self.pending = Some(Pending {
keys,
ark_crypto: hello.ark_crypto,
receiver,
});
}
}
}
struct Feed(Arc<Mutex<Client>>);
impl Read for Feed {
fn set_read_deadline(&mut self, _deadline: Option<Instant>) -> io::Result<()> {
Ok(())
}
}
impl io::Read for Feed {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
let mut client = self.0.lock().unwrap();
if client.bytes.is_empty() {
client.sync();
client.batch = 0;
loop {
match client.steps.pop_front() {
None => {
client.interrupt(Outcome::Terminated);
return Ok(0);
}
Some(Step::Yield) => {
client.interrupt(Outcome::Yield);
return Err(would_block());
}
Some(Step::Interrupt) => return Err(io::ErrorKind::Interrupted.into()),
Some(Step::ReadTimeout) => return Err(io::ErrorKind::TimedOut.into()),
Some(step) if step.is_action() => {
client.action = Some(step);
client.interrupt(Outcome::Yield);
return Err(would_block());
}
Some(step) => client.execute(step),
}
if !client.bytes.is_empty() {
break;
}
}
while client.batch > 1
&& client.outcome == Outcome::Absorbed
&& client.steps.front().is_some_and(Step::queues_frames)
{
let step = client.steps.pop_front().unwrap();
client.execute(step);
client.batch -= 1;
}
}
let mut n = buf.len().min(client.bytes.len() - client.offset);
if client.chunk > 0 {
n = n.min(client.chunk);
}
buf[..n].copy_from_slice(&client.bytes[client.offset..client.offset + n]);
client.offset += n;
if client.offset == client.bytes.len() {
client.bytes.clear();
client.offset = 0;
}
client.summary.reads += 1;
Ok(n)
}
}
type Server = crate::transport::Server<Feed, Outbox, Attestation>;
fn check_session(sender: Option<&Sender<Outbox>>, client: &Client) {
let established = client.state == State::Established;
let refused = super::send(sender, OVERSIZED_MESSAGE);
match refused {
Err(Error::PacketTooLarge(_)) => {
assert!(established, "server has a session the model does not")
}
Err(Error::EncryptionFailed(_)) => {
assert!(!established, "server lacks the session the model has")
}
other => panic!("unexpected oversized send result: {other:?}"),
}
}
fn send(sender: Option<&Sender<Outbox>>, client: &mut Client, id: u64) {
let established = client.state == State::Established;
let expected = established.then(|| {
let receiver = client
.receiver
.clone()
.expect("established without a session");
let sent = client.send(Payload::Frame(Emit::Reply(id, receiver)));
if !sent {
client.forget();
client.state = State::Idle;
if !client.timed_out {
client.send(Payload::Signal);
}
}
sent
});
let sent = super::send(sender, &payload(id));
match (expected, sent) {
(Some(true), Ok(())) => {}
(Some(false), Err(Error::SendFailed(_))) => {}
(None, Err(Error::EncryptionFailed(_))) => {}
(expected, sent) => panic!("model expected {expected:?}, server returned {sent:?}"),
}
}
pub fn run(steps: &[Step]) -> Summary {
#[cfg(feature = "fuzz")]
super::seed::seed(super::seed::TRANSPORT_SERVER, steps);
let signer = xdsa::SecretKey::generate();
let attestation = self_attestation(&signer);
let outbox = Outbox::default();
let client = Arc::new(Mutex::new(Client::new(
steps,
signer.public_key(),
outbox.clone(),
)));
let mut server = Server::new_at(
crate::transport::Stream::new(Feed(client.clone()), outbox, || {}),
signer,
attestation,
TIMESTAMP,
)
.set_handshake_timeout(SCRIPT_HANDSHAKE_TIMEOUT);
let mut sender = None;
let mut retained = None;
let mut generation = 0u64;
let mut retained_generation = None;
loop {
let action = client.lock().unwrap().next_action();
if let Some(action) = action {
let mut client = client.lock().unwrap();
match action {
Step::Retain => {
retained = sender.clone();
retained_generation =
(client.state == State::Established).then_some(generation);
}
Step::Send(tag) => send(sender.as_ref(), &mut client, tag as u64),
Step::SendRetained(tag) => {
if client.state == State::Established && retained_generation == Some(generation)
{
send(retained.as_ref(), &mut client, tag as u64);
} else {
assert!(matches!(
super::send(retained.as_ref(), &payload(tag as u64)),
Err(Error::EncryptionFailed(_))
));
}
}
Step::SendOversized => check_session(sender.as_ref(), &client),
Step::Disconnect => {
client.forget();
if client.state != State::AwaitHello {
client.state = State::Idle;
}
client.held = false;
client.send(Payload::Signal);
server.disconnect();
}
_ => unreachable!("only owner actions reach the driver"),
}
check_session(sender.as_ref(), &client);
client.sync();
continue;
}
match server.recv() {
Ok(Event::Message(message)) => {
let mut client = client.lock().unwrap();
match client.outcome {
Outcome::Message(id) if message == payload(id) => {
client.surfaced(Outcome::Message(id));
client.summary.delivered += 1;
send(sender.as_ref(), &mut client, id);
}
_ if message == [0x07] => client.surfaced(Outcome::Garbage),
outcome => {
panic!("server delivered {message:?} where the model has {outcome:?}")
}
}
}
Ok(Event::Disconnected) => {
let mut client = client.lock().unwrap();
client.surfaced(Outcome::Ended);
check_session(sender.as_ref(), &client);
}
Ok(Event::Connected(opened)) => {
sender = Some(opened);
let mut client = client.lock().unwrap();
client.surfaced(Outcome::Opened);
generation += 1;
check_session(sender.as_ref(), &client);
}
Err(Error::RecvFailed(err)) if err.kind() == io::ErrorKind::WouldBlock => {
let mut client = client.lock().unwrap();
client.surfaced(Outcome::Yield);
check_session(sender.as_ref(), &client);
if client.action.is_none() {
send(sender.as_ref(), &mut client, PROBE_ID);
}
}
Err(Error::Terminated) => {
client.lock().unwrap().surfaced(Outcome::Terminated);
break;
}
Err(Error::SendFailed(_)) => {
client.lock().unwrap().surfaced(Outcome::SendFailed);
}
Err(err) => panic!("unexpected error from the server: {err}"),
}
}
let mut client = client.lock().unwrap();
client.sync();
check_session(sender.as_ref(), &client);
client.summary.state = client.state;
client.summary
}
#[cfg(test)]
mod tests;