use super::session::SessionInner;
use super::{Error, Message, Promise, schema};
use std::sync::Weak;
use std::time::Instant;
#[derive(Debug)]
pub struct Responder {
session: Weak<SessionInner>,
id: Option<u64>,
}
impl Responder {
pub fn reply(
self,
response: impl Into<Message>,
deadline: Instant,
) -> Result<Promise<()>, Error> {
self.enqueue(Ok(response.into()), deadline)
}
pub fn fail(self, error: schema::Error, deadline: Instant) -> Result<Promise<()>, Error> {
self.enqueue(Err(error), deadline)
}
fn enqueue(
mut self,
result: Result<Message, schema::Error>,
deadline: Instant,
) -> Result<Promise<()>, Error> {
let promise = self.session.upgrade().ok_or(Error::Closed)?.reply(
self.id.expect("reply obligation present"),
result,
deadline,
)?;
self.id = None; Ok(promise)
}
pub(super) fn new(session: Weak<SessionInner>, id: u64) -> Self {
Self {
session,
id: Some(id),
}
}
}
impl Drop for Responder {
fn drop(&mut self) {
if let Some(id) = self.id.take()
&& let Some(session) = self.session.upgrade()
{
session.reply_unanswered(id);
}
}
}
#[cfg(test)]
#[cfg_attr(coverage_nightly, coverage(off))]
mod tests {
use crate::protocol::schema::DeviceInfoResponse;
use crate::protocol::{Error, Message, Promise, Responder, Session, schema};
use std::fmt::Debug;
use std::time::Instant;
#[allow(dead_code)]
fn receive_on_host(session: &mut Session, deadline: Instant) -> Result<(), Error> {
let (request, responder): (Message, Responder) = session.recv()?;
let _ = request;
responder
.fail(
schema::Error::reserved(schema::ReservedErrors::Unspecified, "refused"),
deadline,
)?
.wait()
}
#[allow(dead_code)]
fn receive_on_server(session: &mut Session, deadline: Instant) -> Result<(), Error> {
let (request, responder): (Message, Responder) = session.recv()?;
match request {
Message::DeviceInfoRequest(_) => {
let written: Promise<()> =
responder.reply(DeviceInfoResponse::default(), deadline)?;
written.wait()?;
}
Message::Develop(bytes) => {
let requester = session.requester();
std::thread::spawn(move || -> Result<(), Error> {
let answer: Vec<u8> = requester.request(bytes, deadline)?.wait()?;
drop(responder.reply(answer, deadline)?);
Ok(())
});
}
_ => drop(responder), }
Ok(())
}
#[test]
fn test_thread_capabilities() {
fn movable<T: Debug + Send + 'static>() {}
movable::<Responder>();
}
}