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///
45/// # Dropping every handle logs out
46///
47/// The reactor stops when the last clone is dropped: it reads that as "nobody
48/// can drive this session any more" and performs a graceful Logout rather than
49/// leaving a task holding a socket forever. A consumer that wants the session
50/// to outlive its handle must keep one alive — typically by holding it until
51/// [`Connection::wait_closed`] returns.
52#[derive(Debug, Clone)]
53pub struct Connection {
54    /// Session identifier.
55    pub(crate) session_id: SessionId,
56    /// Command channel to the reactor.
57    pub(crate) commands: mpsc::Sender<Command>,
58    /// Closed-flag observation channel.
59    pub(crate) closed: watch::Receiver<bool>,
60    /// Shared session runtime state.
61    pub(crate) runtime: Arc<SessionRuntime>,
62}
63
64impl Connection {
65    /// Returns the session identifier.
66    #[must_use]
67    pub fn session_id(&self) -> &SessionId {
68        &self.session_id
69    }
70
71    /// Sends an application message on the session.
72    ///
73    /// The engine stamps the standard header (including MsgSeqNum) and
74    /// trailer; the message only needs body fields. It is checked here, before
75    /// it is queued, so a message the session layer will not carry is refused
76    /// to the caller rather than dropped later.
77    ///
78    /// # Arguments
79    /// * `message` - The application message to send
80    ///
81    /// # Errors
82    /// * [`EngineError::ReservedMsgType`] for an administrative MsgType.
83    ///   Logon, Logout, SequenceReset and the rest belong to the session state
84    ///   machine; one sent here would bypass the typestate and the engine's
85    ///   phase tracking — a Logout sent this way never arms the logout timeout.
86    /// * [`EngineError::ReservedTag`] for a body field that repeats a tag the
87    ///   engine stamps itself (see [`crate::outbound::RESERVED_TAGS`]), which
88    ///   would put two occurrences of it in the frame.
89    /// * [`EngineError::InvalidField`] for a value with no legal wire form.
90    /// * [`EngineError::Closed`] if the session is already closed.
91    ///
92    /// # Accepted is not sent
93    ///
94    /// `Ok` means the message was queued for the reactor, not that it reached
95    /// the counterparty. A message queued while a Logout is already pending is
96    /// dropped with a warning — the session is on its way out and a new
97    /// application message would arrive after the Logout the counterparty has
98    /// already seen. Use [`Connection::wait_closed`] to observe the end of the
99    /// session and `next_sender_seq` to observe progress.
100    pub async fn send(&self, message: OutboundMessage) -> Result<(), EngineError> {
101        crate::outbound::check_sendable(&message)?;
102        self.commands
103            .send(Command::Send(message))
104            .await
105            .map_err(|_| EngineError::Closed)
106    }
107
108    /// Initiates a graceful logout.
109    ///
110    /// A Logout message is sent to the counterparty; the session closes when
111    /// the Logout acknowledgement arrives or the logout timeout elapses.
112    /// Use [`Connection::wait_closed`] to observe completion.
113    ///
114    /// # Errors
115    /// Returns [`EngineError::Closed`] if the session is already closed.
116    pub async fn logout(&self) -> Result<(), EngineError> {
117        self.commands
118            .send(Command::Logout)
119            .await
120            .map_err(|_| EngineError::Closed)
121    }
122
123    /// Waits until the session is closed.
124    ///
125    /// Fires on transport drop, counterparty logout, heartbeat timeout, or
126    /// completion of a locally initiated logout. Returns immediately if the
127    /// session is already closed.
128    ///
129    /// It fires **after** the inbound application messages already queued have
130    /// been through `from_app` and `on_logout` has run, so a handler still in
131    /// flight when the session ends delays this by however long it takes to
132    /// finish, up to an internal drain bound.
133    pub async fn wait_closed(&self) {
134        let mut closed = self.closed.clone();
135        // An error means the reactor is gone, which also means closed.
136        let _ = closed.wait_for(|closed| *closed).await;
137    }
138
139    /// Returns true if the session is closed.
140    #[must_use]
141    pub fn is_closed(&self) -> bool {
142        self.closed.has_changed().is_err() || *self.closed.borrow()
143    }
144
145    /// Returns true if the session has detected a heartbeat timeout
146    /// (a TestRequest went unanswered for a full heartbeat interval).
147    #[must_use]
148    pub fn is_timed_out(&self) -> bool {
149        // A poisoned lock means a previous holder panicked. The heartbeat state
150        // is plain timestamps with no invariant a panic could leave
151        // half-applied, so the guard is recovered rather than propagated:
152        // taking the consumer's process down over a heartbeat timestamp — and
153        // under `panic = "abort"` that is what propagating means — is never the
154        // right trade.
155        self.runtime
156            .heartbeat
157            .lock()
158            .unwrap_or_else(std::sync::PoisonError::into_inner)
159            .is_timed_out()
160    }
161
162    /// Returns the next outgoing (sender) sequence number.
163    #[must_use]
164    pub fn next_sender_seq(&self) -> u64 {
165        self.runtime.sequences.next_sender_seq().value()
166    }
167
168    /// Returns the next expected incoming (target) sequence number.
169    #[must_use]
170    pub fn next_target_seq(&self) -> u64 {
171        self.runtime.sequences.next_target_seq().value()
172    }
173}