Skip to main content

ironfix_engine/
connection.rs

1/******************************************************************************
2   Author: Joaquín Béjar García
3   Email: jb@taunais.com
4   Date: 14/7/26
5******************************************************************************/
6
7//! Live-session connection handle.
8//!
9//! A [`Connection`] is returned by [`Initiator::connect`](crate::Initiator::connect)
10//! once the FIX session is established. It is a cheap-to-clone handle that
11//! provides an outbound message sink and close observation; the actual
12//! read/write reactor runs in a background task.
13
14use crate::application::SessionId;
15use crate::error::EngineError;
16use crate::outbound::OutboundMessage;
17use ironfix_session::{HeartbeatManager, SequenceManager};
18use std::sync::{Arc, Mutex};
19use tokio::sync::{mpsc, watch};
20
21/// Commands sent from a [`Connection`] handle to the session reactor.
22#[derive(Debug)]
23pub(crate) enum Command {
24    /// Send an application message.
25    Send(OutboundMessage),
26    /// Initiate a graceful logout.
27    Logout,
28}
29
30/// Session runtime state shared between the reactor and connection handles.
31#[derive(Debug)]
32pub(crate) struct SessionRuntime {
33    /// Sender/target sequence counters.
34    pub(crate) sequences: SequenceManager,
35    /// Heartbeat and TestRequest timing state.
36    pub(crate) heartbeat: Mutex<HeartbeatManager>,
37}
38
39/// Handle to a live FIX session.
40///
41/// Cloning is cheap; all clones refer to the same session. The session is
42/// closed when the transport drops, the counterparty logs out, a heartbeat
43/// timeout is detected, or [`Connection::logout`] completes.
44#[derive(Debug, Clone)]
45pub struct Connection {
46    /// Session identifier.
47    pub(crate) session_id: SessionId,
48    /// Command channel to the reactor.
49    pub(crate) commands: mpsc::Sender<Command>,
50    /// Closed-flag observation channel.
51    pub(crate) closed: watch::Receiver<bool>,
52    /// Shared session runtime state.
53    pub(crate) runtime: Arc<SessionRuntime>,
54}
55
56impl Connection {
57    /// Returns the session identifier.
58    #[must_use]
59    pub fn session_id(&self) -> &SessionId {
60        &self.session_id
61    }
62
63    /// Sends an application message on the session.
64    ///
65    /// The engine stamps the standard header (including MsgSeqNum) and
66    /// trailer; the message only needs body fields.
67    ///
68    /// # Arguments
69    /// * `message` - The application message to send
70    ///
71    /// # Errors
72    /// Returns [`EngineError::Closed`] if the session is already closed.
73    pub async fn send(&self, message: OutboundMessage) -> Result<(), EngineError> {
74        self.commands
75            .send(Command::Send(message))
76            .await
77            .map_err(|_| EngineError::Closed)
78    }
79
80    /// Initiates a graceful logout.
81    ///
82    /// A Logout message is sent to the counterparty; the session closes when
83    /// the Logout acknowledgement arrives or the logout timeout elapses.
84    /// Use [`Connection::wait_closed`] to observe completion.
85    ///
86    /// # Errors
87    /// Returns [`EngineError::Closed`] if the session is already closed.
88    pub async fn logout(&self) -> Result<(), EngineError> {
89        self.commands
90            .send(Command::Logout)
91            .await
92            .map_err(|_| EngineError::Closed)
93    }
94
95    /// Waits until the session is closed.
96    ///
97    /// Fires on transport drop, counterparty logout, heartbeat timeout, or
98    /// completion of a locally initiated logout. Returns immediately if the
99    /// session is already closed.
100    pub async fn wait_closed(&self) {
101        let mut closed = self.closed.clone();
102        // An error means the reactor is gone, which also means closed.
103        let _ = closed.wait_for(|closed| *closed).await;
104    }
105
106    /// Returns true if the session is closed.
107    #[must_use]
108    pub fn is_closed(&self) -> bool {
109        self.closed.has_changed().is_err() || *self.closed.borrow()
110    }
111
112    /// Returns true if the session has detected a heartbeat timeout
113    /// (a TestRequest went unanswered for a full heartbeat interval).
114    #[must_use]
115    pub fn is_timed_out(&self) -> bool {
116        self.runtime
117            .heartbeat
118            .lock()
119            .expect("heartbeat lock poisoned")
120            .is_timed_out()
121    }
122
123    /// Returns the next outgoing (sender) sequence number.
124    #[must_use]
125    pub fn next_sender_seq(&self) -> u64 {
126        self.runtime.sequences.next_sender_seq().value()
127    }
128
129    /// Returns the next expected incoming (target) sequence number.
130    #[must_use]
131    pub fn next_target_seq(&self) -> u64 {
132        self.runtime.sequences.next_target_seq().value()
133    }
134}