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}