use super::vector::{Event, ReadError, Vector};
use super::{
CutPoint, MAX_STEPS, Outbox, Recorder, TIMESTAMP, cloud_attestation, frame, self_attestation,
trace, unframe, would_block,
};
use crate::client::MAX_STALE_FRAMES;
use crate::handshake;
use crate::protocol::{ArkToHost, HostToArk, host_to_ark};
use crate::{
Attestation, CRYPTO_DOMAIN_WIRE, CRYPTO_DOMAIN_WIRE_ARK_TO_HOST,
CRYPTO_DOMAIN_WIRE_HOST_TO_ARK, Error, MAX_FRAME_SIZE, MAX_MESSAGE_SIZE,
};
use darkbio_cobs as cobs;
use darkbio_crypto::{cbor, cose, xdsa, xhpke};
use prost::Message;
use std::cell::RefCell;
use std::collections::VecDeque;
use std::io::{self, Read, Write};
use std::rc::Rc;
const PROBE_ID: u64 = u64::MAX;
#[derive(Clone, Debug, PartialEq, Eq)]
#[cfg_attr(feature = "fuzz", derive(arbitrary::Arbitrary))]
pub enum Step {
Handshake,
Send(u8),
Recv,
Hello,
HelloStale,
HelloTampered,
HelloBadAuth,
HelloBadSigner,
HelloBadPayload,
HelloBadKey,
HelloBadEncap,
HelloBadAttest,
Reply(u8),
ReplyReplay,
ReplyTampered,
Garbage,
Dropped,
Junk(Vec<u8>),
Undecodable,
Truncated(u8),
Partial,
Oversized,
Yield,
Interrupt,
Break,
Heal,
Cut { point: CutPoint, then_broken: bool },
Chunk(u8),
Batch(u8),
}
impl Step {
fn is_call(&self) -> bool {
matches!(self, Step::Handshake | Step::Send(_) | Step::Recv)
}
fn queues_frame(&self) -> bool {
!matches!(
self,
Step::Handshake
| Step::Send(_)
| Step::Recv
| Step::Yield
| Step::Interrupt
| Step::Break
| Step::Heal
| Step::Cut { .. }
| Step::Chunk(_)
| Step::Batch(_)
)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Kind {
PacketDecoding,
FrameDecoding,
Send,
Recv,
Terminated,
SessionReset,
InvalidAttestation,
Handshake,
Encryption,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Summary {
pub established: bool, pub handshakes: usize, pub messages: usize, pub resets: usize, pub failures: usize, pub reads: usize, }
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Flaw {
None,
Stale,
Tampered,
Auth,
Signer,
Payload,
Key,
Encap,
Attest,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Frame {
ArkHello { generation: u64, flaw: Flaw },
Sealed {
session: u64,
seq: u64,
tag: u8,
garbage: bool,
},
Dropped,
Junk,
Undecodable,
}
enum Partial {
None,
Hello(u64, Vec<u8>),
Junk(Vec<u8>),
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Call {
None,
Handshake { generation: u64, stale: usize },
Recv,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Expect {
Ok(Option<u64>),
Err(Kind),
}
struct Outstanding {
generation: u64,
crypto: xhpke::SecretKey,
sender: xhpke::Sender,
host_signer: xdsa::PublicKey,
}
struct ServerSession {
id: u64, sender: xhpke::Sender,
receiver: xhpke::Receiver,
seq: u64, }
pub struct Server {
steps: VecDeque<Step>,
identity: xdsa::SecretKey,
attestation: Attestation,
outbox: Outbox, recorder: Recorder, queue: VecDeque<(Vec<u8>, Option<Frame>)>, bytes: Vec<u8>, chunk: usize, batch: usize, broken: bool, cut: Option<CutPoint>, fragment: bool, partial: Partial,
call: Call, expect: Option<Expect>, client_session: Option<u64>, client_seq: u64,
generation: u64, latest_hello: Option<(xdsa::PublicKey, xhpke::PublicKey)>, outstanding: Vec<Outstanding>, session: Option<ServerSession>, last_reply: Option<(Vec<u8>, Frame)>, last_valid: Option<Vec<u8>>,
summary: Summary,
}
impl Server {
fn new(steps: &[Step], outbox: Outbox, recorder: Recorder) -> Self {
let identity = xdsa::SecretKey::generate();
let attestation = self_attestation(&identity);
Self {
steps: steps.iter().take(MAX_STEPS).cloned().collect(),
identity,
attestation,
outbox,
recorder,
queue: VecDeque::new(),
bytes: Vec::new(),
chunk: 0,
batch: 0,
broken: false,
cut: None,
fragment: false,
partial: Partial::None,
call: Call::None,
expect: None,
client_session: None,
client_seq: 0,
generation: 0,
latest_hello: None,
outstanding: Vec::new(),
session: None,
last_reply: None,
last_valid: None,
summary: Summary::default(),
}
}
fn execute(&mut self, step: Step) {
self.ingest();
let produced = match step {
Step::Hello => self.ark_hello(Flaw::None),
Step::HelloStale => self.ark_hello(Flaw::Stale),
Step::HelloTampered => self.ark_hello(Flaw::Tampered),
Step::HelloBadAuth => self.ark_hello(Flaw::Auth),
Step::HelloBadSigner => self.ark_hello(Flaw::Signer),
Step::HelloBadPayload => self.ark_hello(Flaw::Payload),
Step::HelloBadKey => self.ark_hello(Flaw::Key),
Step::HelloBadEncap => self.ark_hello(Flaw::Encap),
Step::HelloBadAttest => self.ark_hello(Flaw::Attest),
Step::Reply(tag) => self.reply(Some(tag)),
Step::ReplyReplay => match self.last_reply.clone() {
Some(replay) => replay,
None => return,
},
Step::ReplyTampered => self.tampered_reply(),
Step::Garbage => self.reply(None),
Step::Dropped => (vec![0x00], Frame::Dropped),
Step::Junk(bytes) => (frame(&bytes), Frame::Junk),
Step::Undecodable => (vec![0xff, 0x01, 0x00], Frame::Undecodable),
Step::Truncated(n) => match self.last_valid.clone() {
Some(valid) => {
let keep = match valid.len() {
0..=1 => 1,
len => 1 + n as usize % (len - 1),
};
let mut bytes = valid[..keep].to_vec();
bytes.push(0x00);
(bytes, classify(&valid[..keep]))
}
None => return,
},
Step::Partial => {
let (mut bytes, frame) = self.ark_hello(Flaw::None);
bytes.pop();
self.partial = match std::mem::replace(&mut self.partial, Partial::None) {
Partial::None => match frame {
Frame::ArkHello { generation, .. } => {
Partial::Hello(generation, bytes.clone())
}
_ => Partial::Junk(bytes.clone()),
},
Partial::Hello(_, mut prior) | Partial::Junk(mut prior) => {
prior.extend_from_slice(&bytes);
Partial::Junk(prior)
}
};
self.queue.push_back((bytes, None));
return;
}
Step::Oversized => {
self.partial = Partial::None;
let mut bytes = vec![1u8; MAX_FRAME_SIZE + 1];
bytes.push(0x00);
self.queue.push_back((bytes, None));
return;
}
Step::Break => {
self.set_broken(true);
return;
}
Step::Heal => {
self.set_broken(false);
return;
}
Step::Cut { point, then_broken } => {
self.cut = Some(point);
self.outbox.set_cut(point);
if then_broken {
self.set_broken(true);
}
return;
}
Step::Chunk(n) => {
self.chunk = n as usize;
return;
}
Step::Batch(n) => {
self.batch = n as usize;
return;
}
Step::Handshake | Step::Send(_) | Step::Recv | Step::Yield | Step::Interrupt => {
unreachable!(
"calls, yields and interrupts are handled by the driver and the reader"
)
}
};
let (bytes, frame) = produced;
let frame = match std::mem::replace(&mut self.partial, Partial::None) {
Partial::None => frame,
Partial::Hello(generation, _) if frame == Frame::Dropped => Frame::ArkHello {
generation,
flaw: Flaw::None,
},
Partial::Hello(_, prior) | Partial::Junk(prior) => {
let mut merged = prior;
merged.extend_from_slice(&bytes[..bytes.len() - 1]);
classify(&merged)
}
};
self.queue.push_back((bytes, Some(frame)));
}
fn ark_hello(&mut self, flaw: Flaw) -> (Vec<u8>, Frame) {
let Some((host_signer, host_crypto)) = self.latest_hello.clone() else {
return self.junk(b"server hello without a client hello");
};
let crypto = xhpke::SecretKey::generate();
if let Some(vector) = self.recorder.borrow_mut().as_mut() {
vector.server_key(crypto.to_bytes().to_vec());
}
let (sender, encap) = host_crypto
.new_sender(CRYPTO_DOMAIN_WIRE_ARK_TO_HOST)
.unwrap();
let payload = handshake::ArkHello {
ark_attest: match flaw {
Flaw::Attest => cloud_attestation(&self.identity),
_ => self.attestation.as_bytes().to_vec(),
},
ark_crypto: crypto.public_key(),
a2h_encap: match flaw {
Flaw::Encap => vec![0x42; 3],
_ => encap.to_vec(),
},
};
let auth = handshake::ArkHelloAuth {
host_signer: match flaw {
Flaw::Auth => xdsa::SecretKey::generate().public_key(),
_ => host_signer.clone(),
},
host_crypto: host_crypto.clone(),
};
let stranger_signer = xdsa::SecretKey::generate();
let signer = match flaw {
Flaw::Signer => &stranger_signer,
_ => &self.identity,
};
let stranger_crypto = xhpke::SecretKey::generate().public_key();
let recipient = match flaw {
Flaw::Stale => &stranger_crypto,
_ => &host_crypto,
};
let mut sealed = match flaw {
Flaw::Payload => cose::seal_at(
&handshake::HostAck {
h2a_encap: vec![1, 2, 3],
},
&auth,
signer,
recipient,
CRYPTO_DOMAIN_WIRE,
TIMESTAMP,
),
Flaw::Key => cose::seal_at(
&(
self.attestation.as_bytes().to_vec(),
vec![0xffu8; xhpke::PUBLIC_KEY_SIZE],
encap.to_vec(),
),
&auth,
signer,
recipient,
CRYPTO_DOMAIN_WIRE,
TIMESTAMP,
),
_ => cose::seal_at(
&payload,
&auth,
signer,
recipient,
CRYPTO_DOMAIN_WIRE,
TIMESTAMP,
),
}
.unwrap();
if flaw == Flaw::Tampered {
*sealed.last_mut().unwrap() ^= 0xff;
}
let generation = match flaw {
Flaw::Stale => 0,
_ => self.generation,
};
if flaw == Flaw::None {
self.outstanding.push(Outstanding {
generation,
crypto,
sender,
host_signer,
});
}
let framed = frame(&sealed);
if flaw == Flaw::None {
self.record(&framed);
}
(framed, Frame::ArkHello { generation, flaw })
}
fn reply(&mut self, tag: Option<u8>) -> (Vec<u8>, Frame) {
let Some(session) = self.session.as_mut() else {
return self.junk(b"reply without a session");
};
let plaintext = match tag {
Some(tag) => ArkToHost {
id: Some(tag as u64),
err: None,
content: None,
}
.encode_to_vec(),
None => vec![0x07],
};
let packet = session.sender.seal(&plaintext, &[]).unwrap();
let sealed = Frame::Sealed {
session: session.id,
seq: session.seq,
tag: tag.unwrap_or_default(),
garbage: tag.is_none(),
};
session.seq += 1;
let produced = (frame(&packet), sealed);
self.record(&produced.0);
self.last_reply = Some(produced.clone());
produced
}
fn record(&mut self, framed: &[u8]) {
self.last_valid = Some(framed[..framed.len() - 1].to_vec());
}
fn tampered_reply(&mut self) -> (Vec<u8>, Frame) {
let Some(session) = self.session.as_mut() else {
return self.junk(b"tampered reply without a session");
};
let plaintext = ArkToHost {
id: Some(0),
err: None,
content: None,
}
.encode_to_vec();
let mut packet = session.sender.seal(&plaintext, &[]).unwrap();
session.seq += 1;
*packet.last_mut().unwrap() ^= 0xff;
(frame(&packet), Frame::Junk)
}
fn junk(&self, text: &[u8]) -> (Vec<u8>, Frame) {
(frame(text), Frame::Junk)
}
fn trace(&self, event: impl FnOnce() -> Event) {
trace(&self.recorder, event);
}
fn set_broken(&mut self, broken: bool) {
self.outbox.set_broken(broken);
self.broken = broken;
}
fn ingest(&mut self) {
for framed in self.outbox.take_frames() {
if framed.is_empty() {
continue;
}
if std::mem::take(&mut self.fragment) {
continue;
}
let packet = unframe(&framed);
if let Ok(hello) = cbor::decode::<handshake::HostHello>(&packet) {
self.generation += 1;
self.latest_hello = Some((hello.host_signer, hello.host_crypto));
continue;
}
if self.ack(&packet) {
continue;
}
if let Some(session) = self.session.as_mut()
&& let Ok(plain) = session.receiver.open(&packet, &[])
{
HostToArk::decode(&plain[..]).expect("client request undecodable");
continue;
}
panic!(
"client wrote a frame the server cannot interpret ({} bytes)",
packet.len()
);
}
}
fn ack(&mut self, packet: &[u8]) -> bool {
for i in (0..self.outstanding.len()).rev() {
let out = &self.outstanding[i];
let auth = handshake::HostAckAuth {
ark_signer: self.identity.public_key(),
ark_crypto: out.crypto.public_key(),
};
let Ok(ack) = cose::open::<handshake::HostAck, _>(
packet,
&auth,
&out.crypto,
&out.host_signer,
CRYPTO_DOMAIN_WIRE,
None,
) else {
continue;
};
let encap: [u8; xhpke::ENCAP_KEY_SIZE] = ack
.h2a_encap
.try_into()
.expect("client ack encap size invalid");
let receiver = out
.crypto
.new_receiver(&encap, CRYPTO_DOMAIN_WIRE_HOST_TO_ARK)
.unwrap();
let out = self.outstanding.remove(i);
self.session = Some(ServerSession {
id: out.generation,
sender: out.sender,
receiver,
seq: 0,
});
return true;
}
false
}
fn consume(&mut self, frame: Frame) {
match self.call {
Call::Handshake { generation, stale } => {
let result = match frame {
Frame::ArkHello {
generation: answered,
flaw,
} if answered == generation => match flaw {
Flaw::None if self.broken || self.cut.is_some() => {
self.apply_cut();
Expect::Err(Kind::Send)
}
Flaw::None => Expect::Ok(None),
Flaw::Tampered
| Flaw::Auth
| Flaw::Signer
| Flaw::Payload
| Flaw::Key
| Flaw::Encap => Expect::Err(Kind::Handshake),
Flaw::Attest => Expect::Err(Kind::InvalidAttestation),
Flaw::Stale => unreachable!("stale hellos answer no generation"),
},
_ => {
if stale < MAX_STALE_FRAMES {
self.call = Call::Handshake {
generation,
stale: stale + 1,
};
return;
}
Expect::Err(Kind::Handshake)
}
};
if result == Expect::Ok(None) {
self.client_session = Some(generation);
self.client_seq = 0;
}
self.settle(result);
}
Call::Recv => {
let result = match (self.client_session, frame) {
(_, Frame::Dropped) => {
self.client_session = None;
Expect::Err(Kind::SessionReset)
}
(_, Frame::Undecodable) => {
self.client_session = None;
Expect::Err(Kind::FrameDecoding)
}
(None, _) => Expect::Err(Kind::Encryption),
(
Some(id),
Frame::Sealed {
session,
seq,
tag,
garbage,
},
) if session == id && seq == self.client_seq => {
self.client_seq += 1;
if garbage {
Expect::Err(Kind::PacketDecoding)
} else {
Expect::Ok(Some(tag as u64))
}
}
(Some(_), _) => {
self.client_session = None;
Expect::Err(Kind::Encryption)
}
};
self.settle(result);
}
Call::None => panic!("client read a frame outside a call"),
}
}
fn apply_cut(&mut self) {
if let Some(CutPoint::Middle(_)) = self.cut.take() {
self.fragment = true;
}
}
fn interrupt(&mut self, kind: Kind) {
match self.call {
Call::Handshake { .. } => self.settle(Expect::Err(kind)),
Call::Recv => {
self.client_session = None;
self.settle(Expect::Err(kind));
}
Call::None => panic!("client read outside a call"),
}
}
fn settle(&mut self, result: Expect) {
assert!(self.expect.is_none(), "call settled twice");
self.expect = Some(result);
self.call = Call::None;
}
fn finish(&mut self, result: Result<Option<u64>, Kind>) {
let expected = self
.expect
.take()
.expect("call returned before the model settled it");
let actual = match result {
Ok(id) => Expect::Ok(id),
Err(kind) => Expect::Err(kind),
};
assert_eq!(actual, expected);
self.call = Call::None;
self.ingest();
}
fn next_call(&mut self) -> Option<Step> {
loop {
match self.steps.pop_front() {
None => return None,
Some(step) if step.is_call() => return Some(step),
Some(Step::Yield) | Some(Step::Interrupt) => {}
Some(step) => self.execute(step),
}
}
}
}
struct Feed(Rc<RefCell<Server>>);
impl Read for Feed {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
let mut server = self.0.borrow_mut();
loop {
if !server.bytes.is_empty() {
let mut n = buf.len().min(server.bytes.len());
if server.chunk > 0 {
n = n.min(server.chunk);
}
buf[..n].copy_from_slice(&server.bytes[..n]);
server.bytes.drain(..n);
server.summary.reads += 1;
return Ok(n);
}
if let Some((bytes, frame)) = server.queue.pop_front() {
if let Some(frame) = frame {
server.consume(frame);
}
server.bytes = bytes;
while server.batch > 1 && server.expect.is_none() {
if server.queue.is_empty() {
match server.steps.front() {
Some(step) if step.queues_frame() => {
let step = server.steps.pop_front().unwrap();
server.execute(step);
continue;
}
_ => break,
}
}
let (bytes, frame) = server.queue.pop_front().unwrap();
if let Some(frame) = frame {
server.consume(frame);
server.batch -= 1;
}
server.bytes.extend(bytes);
}
server.batch = 0;
server.trace(|| Event::Read {
bytes: server.bytes.clone(),
chunk: server.chunk,
});
continue;
}
match server.steps.front() {
None => {
server.interrupt(Kind::Terminated);
server.trace(|| Event::ReadFailed {
error: ReadError::Eof,
});
return Ok(0);
}
Some(step) if step.is_call() => {
server.interrupt(Kind::Recv);
server.trace(|| Event::ReadFailed {
error: ReadError::Failed,
});
return Err(would_block());
}
Some(Step::Yield) => {
server.steps.pop_front();
server.interrupt(Kind::Recv);
server.trace(|| Event::ReadFailed {
error: ReadError::Failed,
});
return Err(would_block());
}
Some(Step::Interrupt) => {
server.steps.pop_front();
server.trace(|| Event::ReadFailed {
error: ReadError::Interrupted,
});
return Err(io::ErrorKind::Interrupted.into());
}
Some(_) => {
let step = server.steps.pop_front().unwrap();
server.execute(step);
}
}
}
}
}
type Client = crate::Client<Feed, Outbox>;
fn classify(bytes: &[u8]) -> Frame {
let mut buf = vec![0u8; cobs::decode_buffer(bytes.len())];
match cobs::decode(bytes, &mut buf) {
Ok(_) => Frame::Junk,
Err(_) => Frame::Undecodable,
}
}
fn kind(err: Error) -> Kind {
match err {
Error::PacketDecodingFailed(_) => Kind::PacketDecoding,
Error::FrameDecodingFailed(_) => Kind::FrameDecoding,
Error::SendFailed(_) => Kind::Send,
Error::RecvFailed(_) => Kind::Recv,
Error::Terminated => Kind::Terminated,
Error::SessionReset => Kind::SessionReset,
Error::InvalidAttestation => Kind::InvalidAttestation,
Error::HandshakeFailed(_) => Kind::Handshake,
Error::EncryptionFailed(_) => Kind::Encryption,
err => panic!("unexpected error from the client: {err}"),
}
}
pub(super) fn check_session<R: Read, W: Write>(
client: &mut crate::Client<R, W>,
established: bool,
) {
let oversized = vec![0x42; MAX_MESSAGE_SIZE + 1];
let refused = client.send_message(HostToArk {
id: Some(PROBE_ID),
content: Some(host_to_ark::Content::Develop(oversized)),
});
match refused {
Err(Error::PacketTooLarge(_)) => {
assert!(established, "client has a session the model does not")
}
Err(Error::EncryptionFailed(_)) => {
assert!(!established, "client lacks the session the model has")
}
other => panic!("unexpected oversized send result: {other:?}"),
}
}
pub fn run(steps: &[Step]) -> Summary {
#[cfg(feature = "fuzz")]
super::seed::seed(super::seed::CLIENT_PROTOCOL, steps);
let scenario = super::vector::scenario();
#[cfg(all(test, feature = "fuzz", getrandom_backend = "custom"))]
if let Some(scenario) = &scenario {
super::random::reseed(scenario);
}
let recorder = Recorder::default();
let outbox = Outbox {
recorder: recorder.clone(),
..Outbox::default()
};
let server = Rc::new(RefCell::new(Server::new(
steps,
outbox.clone(),
recorder.clone(),
)));
let identity = server.borrow().identity.public_key();
let mut client = Client::new(Feed(server.clone()), outbox);
let steps = &steps[..steps.len().min(MAX_STEPS)];
let write_failures = steps
.iter()
.any(|step| matches!(step, Step::Break | Step::Cut { .. }));
*recorder.borrow_mut() = Vector::open(
scenario,
steps,
write_failures,
identity.to_bytes().to_vec(),
server.borrow().attestation.as_bytes().to_vec(),
);
loop {
let call = server.borrow_mut().next_call();
let Some(call) = call else { break };
match call {
Step::Handshake => {
let signer = xdsa::SecretKey::generate();
let crypto = xhpke::SecretKey::generate();
trace(&recorder, || Event::Handshake {
xdsa: signer.to_bytes().to_vec(),
xhpke: crypto.to_bytes().to_vec(),
});
let failing = {
let server = server.borrow();
server.broken || server.cut.is_some()
};
if failing {
let result = client
.handshake_with_keys(&identity, signer, crypto, TIMESTAMP)
.map(|_| ());
assert!(matches!(result, Err(Error::SendFailed(_))), "{result:?}");
trace(&recorder, || Event::Err {
kind: "SendFailed".into(),
});
let mut server = server.borrow_mut();
server.client_session = None;
server.summary.failures += 1;
server.apply_cut();
continue;
}
{
let mut server = server.borrow_mut();
server.client_session = None;
server.call = Call::Handshake {
generation: server.generation + 1,
stale: 0,
};
}
let result = client.handshake_with_keys(&identity, signer, crypto, TIMESTAMP);
trace(&recorder, || match &result {
Ok(_) => Event::Ok { message: None },
Err(err) => Event::Err {
kind: <&str>::from(err).into(),
},
});
let result = result.map(|_| None).map_err(kind);
let mut server = server.borrow_mut();
match result {
Ok(_) => server.summary.handshakes += 1,
Err(_) => server.summary.failures += 1,
}
server.finish(result);
}
Step::Send(tag) => {
let established = server.borrow().client_session.is_some();
trace(&recorder, || Event::Session { established });
check_session(&mut client, established);
let expected: Result<Option<u64>, Kind> = {
let mut server = server.borrow_mut();
match (established, server.broken || server.cut.is_some()) {
(true, false) => Ok(None),
(true, true) => {
server.client_session = None;
server.apply_cut();
Err(Kind::Send)
}
(false, _) => Err(Kind::Encryption),
}
};
let request = HostToArk {
id: Some(tag as u64),
content: None,
};
trace(&recorder, || Event::Send {
message: request.encode_to_vec(),
});
let result = client.send_message(request);
trace(&recorder, || match &result {
Ok(_) => Event::Ok { message: None },
Err(err) => Event::Err {
kind: <&str>::from(err).into(),
},
});
let result = result.map(|_| None).map_err(kind);
assert_eq!(result, expected);
server.borrow_mut().ingest();
}
Step::Recv => {
server.borrow_mut().call = Call::Recv;
trace(&recorder, || Event::Recv);
let result = client.next_message();
trace(&recorder, || match &result {
Ok(msg) => Event::Ok {
message: Some(msg.encode_to_vec()),
},
Err(err) => Event::Err {
kind: <&str>::from(err).into(),
},
});
let result = result.map(|msg| msg.id).map_err(kind);
let mut server = server.borrow_mut();
match result {
Ok(_) => server.summary.messages += 1,
Err(Kind::SessionReset) => server.summary.resets += 1,
Err(_) => server.summary.failures += 1,
}
server.finish(result);
}
_ => unreachable!("only calls reach the driver"),
}
}
let mut server = server.borrow_mut();
server.ingest();
let established = server.client_session.is_some();
trace(&recorder, || Event::Session { established });
check_session(&mut client, established);
server.summary.established = established;
if let Some(vector) = recorder.borrow().as_ref() {
vector.write();
#[cfg(test)]
{
assert!(
super::vector::replay::parse(&vector.json()) == *vector,
"transcript does not survive its encoding"
);
super::vector::replay::run(vector);
}
}
server.summary
}
#[cfg(test)]
mod tests;