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}