use crate::application::SessionId;
use crate::error::EngineError;
use crate::outbound::OutboundMessage;
use ironfix_session::{HeartbeatManager, SequenceManager};
use std::sync::{Arc, Mutex};
use tokio::sync::{mpsc, watch};
#[derive(Debug)]
pub(crate) enum Command {
Send(OutboundMessage),
Logout,
}
#[derive(Debug)]
pub(crate) struct SessionRuntime {
pub(crate) sequences: SequenceManager,
pub(crate) heartbeat: Mutex<HeartbeatManager>,
}
#[derive(Debug, Clone)]
pub struct Connection {
pub(crate) session_id: SessionId,
pub(crate) commands: mpsc::Sender<Command>,
pub(crate) closed: watch::Receiver<bool>,
pub(crate) runtime: Arc<SessionRuntime>,
}
impl Connection {
#[must_use]
pub fn session_id(&self) -> &SessionId {
&self.session_id
}
pub async fn send(&self, message: OutboundMessage) -> Result<(), EngineError> {
crate::outbound::check_sendable(&message)?;
self.commands
.send(Command::Send(message))
.await
.map_err(|_| EngineError::Closed)
}
pub async fn logout(&self) -> Result<(), EngineError> {
self.commands
.send(Command::Logout)
.await
.map_err(|_| EngineError::Closed)
}
pub async fn wait_closed(&self) {
let mut closed = self.closed.clone();
let _ = closed.wait_for(|closed| *closed).await;
}
#[must_use]
pub fn is_closed(&self) -> bool {
self.closed.has_changed().is_err() || *self.closed.borrow()
}
#[must_use]
pub fn is_timed_out(&self) -> bool {
self.runtime
.heartbeat
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_timed_out()
}
#[must_use]
pub fn next_sender_seq(&self) -> u64 {
self.runtime.sequences.next_sender_seq().value()
}
#[must_use]
pub fn next_target_seq(&self) -> u64 {
self.runtime.sequences.next_target_seq().value()
}
}