ironfix_engine/acceptor.rs
1/******************************************************************************
2 Author: Joaquín Béjar García
3 Email: jb@taunais.com
4 Date: 22/7/26
5******************************************************************************/
6
7//! Server-side (acceptor) FIX engine: inbound connections, framing, and the
8//! acceptor half of the Logon handshake.
9//!
10//! [`Acceptor::serve`] takes an already-accepted [`TcpStream`], frames it with
11//! [`ironfix_transport::FixCodec`], drives the [`ironfix_session`] typestate
12//! machine through `accept() -> on_logon_received() -> accept_logon()`, and
13//! hands the socket to the same background session reactor the
14//! [`Initiator`](crate::Initiator) uses (see `crate::reactor`). The returned
15//! [`Connection`] handle is the outbound message sink and exposes
16//! `wait_closed()` / `is_timed_out()`. [`Acceptor::accept`] is the convenience
17//! wrapper that pulls the next connection off a [`TcpListener`] and hands it to
18//! `serve`.
19//!
20//! The accept loop and its supervision are the consumer's, exactly as
21//! reconnection is the consumer's for the initiator: an `Acceptor` is a single
22//! configured session (one expected counterparty), and each accepted
23//! connection becomes its own live session with its own sequence state. A
24//! server that fronts several counterparties runs one accept loop and matches
25//! each inbound connection to the `Acceptor` configured for it.
26//!
27//! # Handshake conformance
28//!
29//! The inbound Logon is validated in the order set out in
30//! `doc/fix_operations.md` ("Logon"): it must decode and be a Logon; its
31//! `BeginString` (8) must match this session's version and its `EncryptMethod`
32//! (98) must be 0 (None); it must carry `MsgSeqNum` (34) **in the standard
33//! header**; then counterparty identity (49/56, plus 50/57 when configured) is
34//! checked and the single admission slot is claimed — a second concurrent Logon
35//! for the same counterparty is refused rather than allowed to fork the session;
36//! then `SendingTime` (52) accuracy and the `from_admin` authentication hook;
37//! then the `HeartBtInt` (108) the initiator requested is bounded, honored, and
38//! echoed, `ResetSeqNumFlag` (141) is reconciled with the local
39//! `reset_on_logon` knob and mirrored on the ack, and finally `MsgSeqNum` is
40//! validated. A failure at any step sends a session Reject (reason 9 for
41//! identity, the `SendingTime` reason for a clock problem) and/or a Logout,
42//! drives the typestate to `reject_logon`, and drops the connection without ever
43//! reaching Active. Waiting for the inbound Logon is bounded by
44//! [`SessionConfig::logon_timeout`]. A gap in the Logon completes the handshake
45//! and immediately issues a `ResendRequest` (2).
46//!
47//! `HeartBtInt` (108) is counterparty-controlled and drives a `Duration` on the
48//! heartbeat clock, so it is bounded at the handshake: a value large enough to
49//! overflow `interval + grace` would otherwise abort the process under
50//! `panic = "abort"`. `ResetSeqNumFlag` (141) is reconciled coherently — the
51//! acceptor resets when the peer asks *or* when it is locally configured to, and
52//! whichever drives the reset the ack carries `141=Y` so the peer resets in
53//! lockstep instead of silently desyncing.
54//!
55//! Every handshake frame goes through the same peek-then-spend path the reactor
56//! uses (`send_handshake_admin`): the `to_admin` callback runs
57//! on the message before it is framed, the body is re-checked, and the sender
58//! sequence number is spent only once the frame has been built. Once the session
59//! is Active the shared session reactor runs identically for both roles. All
60//! sequence arithmetic goes through the checked
61//! [`SequenceManager::try_allocate_sender_seq`](ironfix_session::SequenceManager::try_allocate_sender_seq) /
62//! [`SequenceManager::try_increment_target_seq`](ironfix_session::SequenceManager::try_increment_target_seq);
63//! an exhausted counter refuses the session rather than wrapping.
64//!
65//! # Example
66//!
67//! ```no_run
68//! use ironfix_core::types::CompId;
69//! use ironfix_engine::application::NoOpApplication;
70//! use ironfix_engine::Acceptor;
71//! use ironfix_session::SessionConfig;
72//! use std::sync::Arc;
73//! use tokio::net::TcpListener;
74//!
75//! # async fn run() -> Result<(), Box<dyn std::error::Error>> {
76//! // Sender is the acceptor itself; target is the initiator it expects.
77//! let config = SessionConfig::new(
78//! CompId::new("VENUE").unwrap(),
79//! CompId::new("CLIENT").unwrap(),
80//! "FIX.4.4",
81//! );
82//! let acceptor = Arc::new(Acceptor::new(config, Arc::new(NoOpApplication)));
83//! let listener = TcpListener::bind("127.0.0.1:9876").await?;
84//!
85//! loop {
86//! let (stream, _) = listener.accept().await?;
87//! let acceptor = Arc::clone(&acceptor);
88//! tokio::spawn(async move {
89//! if let Ok(connection) = acceptor.serve(stream).await {
90//! connection.wait_closed().await;
91//! }
92//! });
93//! }
94//! # }
95//! ```
96
97use crate::application::{Application, NoOpApplication, RejectReason, SessionId};
98use crate::connection::{Connection, SessionRuntime};
99use crate::error::EngineError;
100use crate::reactor::{
101 DEFAULT_APP_QUEUE_CAPACITY, DEFAULT_OUTBOUND_CAPACITY, DEFAULT_WRITE_TIMEOUT, ResendState,
102 SessionParams, lock_heartbeat, send_handshake_admin, spawn_session,
103};
104use crate::wire::{self, MessageFactory, PeerIdentity, SendingTimeGuard, UnsupportedVersion};
105use futures_util::StreamExt;
106use ironfix_core::message::MsgType;
107use ironfix_core::version::FixVersion;
108use ironfix_session::sequence::SequenceResult;
109use ironfix_session::{Disconnected, HeartbeatManager, SequenceManager, Session, SessionConfig};
110use ironfix_transport::FixCodec;
111use std::collections::HashSet;
112use std::num::NonZeroU64;
113use std::sync::{Arc, Mutex, PoisonError};
114use std::time::Duration;
115use tokio::net::{TcpListener, TcpStream};
116use tokio::time::timeout;
117use tokio_util::codec::Framed;
118
119/// RAII claim on the one live session an [`Acceptor`] admits at a time for a
120/// given configured counterparty.
121///
122/// An `Acceptor` is a single configured session (one expected initiator), so
123/// two connections that both pass the handshake would otherwise each build an
124/// independent [`SequenceManager`] and both reach Active at sequence 1,
125/// silently forking the session. The guard is claimed once the inbound Logon's
126/// identity is validated and is then moved into the reactor task, so it is held
127/// for the whole life of the session and released on **every** close path when
128/// the reactor task ends — after which the counterparty may reconnect.
129struct AdmissionGuard {
130 /// The set of currently-admitted sessions, shared with the [`Acceptor`].
131 admitted: Arc<Mutex<HashSet<SessionId>>>,
132 /// The session this guard admitted.
133 session_id: SessionId,
134}
135
136impl Drop for AdmissionGuard {
137 fn drop(&mut self) {
138 // A poisoned lock means a previous holder panicked; the admission set is
139 // a plain set of identifiers with no invariant a panic could leave
140 // half-applied, so the guard is recovered rather than propagated. Under
141 // `panic = "abort"` refusing to free the slot would strand the session
142 // forever.
143 let mut admitted = self.admitted.lock().unwrap_or_else(PoisonError::into_inner);
144 admitted.remove(&self.session_id);
145 }
146}
147
148/// Server-side FIX engine.
149///
150/// Owns the session configuration and the [`Application`] callbacks. Each
151/// accepted connection ([`Acceptor::accept`] / [`Acceptor::serve`]) establishes
152/// one live session and returns a [`Connection`] handle for it. Cloning the
153/// application is cheap, so an `Acceptor` is typically wrapped in an [`Arc`]
154/// and shared across the tasks serving each connection.
155///
156/// The configuration is stated from the acceptor's point of view: its
157/// `sender_comp_id` is the acceptor's own CompID and its `target_comp_id` is
158/// the initiator it expects — the reverse of what the peer stamps on the wire,
159/// which is exactly what the inbound-identity check validates against.
160#[derive(Debug)]
161pub struct Acceptor<A: Application = NoOpApplication> {
162 /// Session configuration.
163 config: SessionConfig,
164 /// Application callbacks.
165 application: Arc<A>,
166 /// Session identifier derived from the configuration.
167 session_id: SessionId,
168 /// The configured FIX version, or why it cannot be framed. Resolved once
169 /// at construction and reported before the handshake begins.
170 version: Result<FixVersion, UnsupportedVersion>,
171 /// Initial (sender, target) sequence numbers for session continuity.
172 initial_sequences: Option<(u64, u64)>,
173 /// Capacity of the outbound command queue.
174 outbound_capacity: usize,
175 /// Sessions currently live on this acceptor, one entry per admitted
176 /// counterparty. Shared across every [`Acceptor::serve`] call so a second
177 /// concurrent Logon for a session already established is refused rather than
178 /// forking it. See [`AdmissionGuard`].
179 admitted: Arc<Mutex<HashSet<SessionId>>>,
180}
181
182impl<A: Application + 'static> Acceptor<A> {
183 /// Creates a new acceptor.
184 ///
185 /// # Arguments
186 /// * `config` - The session configuration (sender = this acceptor's CompID,
187 /// target = the initiator it expects)
188 /// * `application` - The application callback handler
189 #[must_use]
190 pub fn new(config: SessionConfig, application: Arc<A>) -> Self {
191 let session_id = wire::session_id_from_config(&config);
192 let version = wire::wire_version(&config.begin_string);
193
194 Self {
195 config,
196 application,
197 session_id,
198 version,
199 initial_sequences: None,
200 outbound_capacity: DEFAULT_OUTBOUND_CAPACITY,
201 admitted: Arc::new(Mutex::new(HashSet::new())),
202 }
203 }
204
205 /// Claims the admission slot for this acceptor's configured session.
206 ///
207 /// Returns `Some(guard)` when the slot was free — the caller then owns the
208 /// only live session for this counterparty until the guard drops — or `None`
209 /// when a session is already established, in which case the caller must
210 /// refuse the connection. The check and the claim are one atomic step under
211 /// the lock, so two concurrent Logons cannot both succeed.
212 fn try_admit(&self) -> Option<AdmissionGuard> {
213 let mut admitted = self.admitted.lock().unwrap_or_else(PoisonError::into_inner);
214 if !admitted.insert(self.session_id.clone()) {
215 return None;
216 }
217 Some(AdmissionGuard {
218 admitted: Arc::clone(&self.admitted),
219 session_id: self.session_id.clone(),
220 })
221 }
222
223 /// Seeds each accepted session with initial sequence numbers, for
224 /// continuity with a previous session. Ignored when the inbound Logon
225 /// carries `ResetSeqNumFlag` (141) = Y, which resets both counters to 1.
226 ///
227 /// # Arguments
228 /// * `sender_seq` - Next outgoing sequence number
229 /// * `target_seq` - Next expected incoming sequence number
230 #[must_use]
231 pub fn with_initial_sequences(mut self, sender_seq: u64, target_seq: u64) -> Self {
232 self.initial_sequences = Some((sender_seq, target_seq));
233 self
234 }
235
236 /// Sets the capacity of each session's outbound message queue (default
237 /// 1024).
238 #[must_use]
239 pub fn with_outbound_capacity(mut self, capacity: usize) -> Self {
240 self.outbound_capacity = capacity.max(1);
241 self
242 }
243
244 /// Returns the session configuration.
245 #[must_use]
246 pub fn config(&self) -> &SessionConfig {
247 &self.config
248 }
249
250 /// Returns the session identifier.
251 #[must_use]
252 pub fn session_id(&self) -> &SessionId {
253 &self.session_id
254 }
255
256 /// Accepts the next inbound connection on `listener` and establishes a
257 /// session on it.
258 ///
259 /// This is a thin convenience over [`Acceptor::serve`]: it awaits one TCP
260 /// connection and hands the stream to the handshake. For concurrent
261 /// handshakes, run the accept loop yourself and spawn a `serve` per
262 /// connection (see the module example).
263 ///
264 /// # Arguments
265 /// * `listener` - A bound [`TcpListener`]
266 ///
267 /// # Errors
268 /// Returns [`EngineError::Io`] if the TCP accept fails, or any error
269 /// [`Acceptor::serve`] can produce.
270 pub async fn accept(&self, listener: &TcpListener) -> Result<Connection, EngineError> {
271 let (stream, _addr) = listener.accept().await?;
272 self.serve(stream).await
273 }
274
275 /// Establishes a session on an already-accepted [`TcpStream`], completing
276 /// the acceptor-side Logon handshake and spawning the session reactor.
277 ///
278 /// On success the session is Active: `on_logon` has fired and the returned
279 /// [`Connection`] can send application messages. The reactor owns the
280 /// socket and handles heartbeats, TestRequests, sequence validation, and
281 /// admin replies until the session closes.
282 ///
283 /// # Arguments
284 /// * `stream` - An accepted TCP connection from the counterparty
285 ///
286 /// # Errors
287 /// Returns an [`EngineError`] if framing or the Logon handshake fails. That
288 /// includes [`EngineError::LogonTimeout`] when no Logon arrives within
289 /// [`SessionConfig::logon_timeout`], [`EngineError::UnexpectedMessage`] when
290 /// the first frame is not a Logon, [`EngineError::IdentityMismatch`] when
291 /// the Logon's CompIDs do not match the configured counterparty, and
292 /// [`EngineError::SequenceExhausted`] when a sequence counter has reached
293 /// `u64::MAX`. [`EngineError::LogonRejected`] covers the acceptor-side
294 /// refusals: the `from_admin` authentication hook, a `BeginString` (8) that
295 /// does not match this session's version, an unsupported `EncryptMethod`
296 /// (98), a `HeartBtInt` (108) that is missing, non-numeric, or out of range,
297 /// and a second concurrent Logon for a session already established. A
298 /// `BeginString` this engine cannot frame *at all* is refused up front with
299 /// [`EngineError::UnsupportedVersion`].
300 pub async fn serve(&self, stream: TcpStream) -> Result<Connection, EngineError> {
301 // Refuse before framing: an unsupported version cannot produce a
302 // conforming Logon reply, and guessing one would put a fabricated
303 // version on the wire.
304 let version = match &self.version {
305 Ok(version) => *version,
306 Err(err) => {
307 return Err(EngineError::UnsupportedVersion {
308 version: err.version.clone(),
309 detail: err.detail.clone(),
310 });
311 }
312 };
313
314 let session_id = self.session_id.clone();
315 self.application.on_create(&session_id).await;
316
317 // Typestate: Disconnected -> Connecting (acceptor side).
318 let session = Session::<Disconnected>::new(session_id.to_string()).accept();
319
320 let _ = stream.set_nodelay(true);
321 let codec = FixCodec::new()
322 .with_max_message_size(self.config.max_message_size)
323 .with_checksum_validation(self.config.validate_checksum);
324 let mut framed = Framed::new(stream, codec);
325
326 let sequences = match self.initial_sequences {
327 Some((sender, target)) if !self.config.reset_on_logon => {
328 // MsgSeqNum starts at 1, so a zero seed is floored to 1 rather
329 // than numbering a message 0 that every counterparty rejects.
330 SequenceManager::with_initial(
331 NonZeroU64::new(sender).unwrap_or(NonZeroU64::MIN),
332 NonZeroU64::new(target).unwrap_or(NonZeroU64::MIN),
333 )
334 }
335 _ => SequenceManager::new(),
336 };
337 let runtime = Arc::new(SessionRuntime {
338 sequences,
339 heartbeat: Mutex::new(HeartbeatManager::new(self.config.heartbeat_interval)),
340 });
341 let mut factory = MessageFactory::new(&self.config, version);
342 let identity = PeerIdentity::new(&self.config);
343 let sending_time = SendingTimeGuard::new(&self.config);
344
345 // Await the inbound Logon, bounded by logon_timeout.
346 let logon_frame = match timeout(self.config.logon_timeout, framed.next()).await {
347 Err(_) => {
348 let _ = session.disconnect();
349 return Err(EngineError::LogonTimeout(self.config.logon_timeout));
350 }
351 Ok(None) => {
352 let _ = session.disconnect();
353 return Err(EngineError::Closed);
354 }
355 Ok(Some(Err(err))) => {
356 let _ = session.disconnect();
357 return Err(err.into());
358 }
359 Ok(Some(Ok(frame))) => frame,
360 };
361
362 // Typestate: Connecting -> LogonReceived.
363 let session = session.on_logon_received();
364
365 // (gap start, high-water) of a gap detected in the inbound Logon.
366 let mut pending_resend: Option<(u64, u64)> = None;
367 // The admission claim is taken once identity is validated and moved into
368 // the reactor task on success; on any handshake failure after the claim
369 // this local drops and frees the slot. Assigned on the fall-through path
370 // out of the block below; every path before the claim returns first.
371 let admission;
372 {
373 let raw = match wire::decode_frame(&logon_frame) {
374 Ok(raw) => raw,
375 Err(err) => {
376 let _ = session.reject_logon();
377 return Err(err.into());
378 }
379 };
380
381 match raw.msg_type() {
382 MsgType::Logon => {}
383 other => {
384 let msg_type = other.as_str().to_string();
385 let _ = session.reject_logon();
386 return Err(EngineError::UnexpectedMessage { msg_type });
387 }
388 }
389
390 // BeginString (8): a peer speaking a different FIX version cannot be
391 // given a conforming ack in this session's version, so it must not
392 // reach Active in ours. Compared against the wire BeginString the
393 // header stamper uses, which is `FIXT.1.1` for every 5.0 session.
394 let inbound_begin_string = raw.begin_string().unwrap_or_default();
395 if inbound_begin_string != version.begin_string() {
396 let detail = format!(
397 "Logon BeginString (8) '{}' does not match the configured '{}'",
398 inbound_begin_string.chars().take(16).collect::<String>(),
399 version.begin_string()
400 );
401 let logout = factory.logout(Some(&detail));
402 let _ = send_handshake_admin(
403 self.application.as_ref(),
404 &session_id,
405 &mut framed,
406 &mut factory,
407 &runtime.sequences,
408 None,
409 DEFAULT_WRITE_TIMEOUT,
410 logout,
411 )
412 .await;
413 let _ = session.reject_logon();
414 return Err(EngineError::LogonRejected { reason: detail });
415 }
416
417 // EncryptMethod (98): IronFix implements no encryption, so only 0
418 // (None) can be honoured. A missing or non-zero method is refused
419 // rather than silently accepted as if it were plaintext.
420 let encrypt_method = raw.get_field_str(98);
421 if encrypt_method != Some("0") {
422 let shown = encrypt_method
423 .map_or_else(|| "<absent>".to_string(), |v| v.chars().take(16).collect());
424 let detail =
425 format!("unsupported EncryptMethod (98) '{shown}'; only 0 (None) is supported");
426 let logout = factory.logout(Some(&detail));
427 let _ = send_handshake_admin(
428 self.application.as_ref(),
429 &session_id,
430 &mut framed,
431 &mut factory,
432 &runtime.sequences,
433 None,
434 DEFAULT_WRITE_TIMEOUT,
435 logout,
436 )
437 .await;
438 let _ = session.reject_logon();
439 return Err(EngineError::LogonRejected { reason: detail });
440 }
441
442 // MsgSeqNum (34) must sit in the standard header: a body-only 34 is
443 // not the session sequence number, exactly as the identity check
444 // requires 49/56 in the header.
445 let logon_seq = match wire::header_seq_num(&raw) {
446 Some(seq) => seq,
447 None => {
448 let _ = session.reject_logon();
449 return Err(EngineError::Sequence(
450 "Logon has no valid MsgSeqNum (34) in the standard header".to_string(),
451 ));
452 }
453 };
454
455 // Identity before anything else: a cross-wired initiator must not
456 // be allowed to establish a session or move sequence state.
457 if let Err(mismatch) = identity.validate(&raw) {
458 let detail = mismatch.to_string();
459 let reason = RejectReason::new(9, detail.clone()).with_ref_tag(mismatch.tag);
460 let reject = factory.session_reject(logon_seq, MsgType::Logon.as_str(), &reason);
461 let _ = send_handshake_admin(
462 self.application.as_ref(),
463 &session_id,
464 &mut framed,
465 &mut factory,
466 &runtime.sequences,
467 None,
468 DEFAULT_WRITE_TIMEOUT,
469 reject,
470 )
471 .await;
472 let logout = factory.logout(Some(&detail));
473 let _ = send_handshake_admin(
474 self.application.as_ref(),
475 &session_id,
476 &mut framed,
477 &mut factory,
478 &runtime.sequences,
479 None,
480 DEFAULT_WRITE_TIMEOUT,
481 logout,
482 )
483 .await;
484 let _ = session.reject_logon();
485 return Err(EngineError::IdentityMismatch { detail });
486 }
487
488 // Identity is proven: claim the single admission slot. A second
489 // concurrent Logon for this same counterparty is refused here rather
490 // than allowed to fork the session into two independent sequence
491 // streams both starting at 1. Policy: the live session is preserved
492 // and the newcomer is logged out.
493 admission = match self.try_admit() {
494 Some(guard) => guard,
495 None => {
496 let detail = format!("session already active for {session_id}");
497 tracing::warn!(session = %session_id, "refusing duplicate concurrent Logon");
498 let logout = factory.logout(Some(&detail));
499 let _ = send_handshake_admin(
500 self.application.as_ref(),
501 &session_id,
502 &mut framed,
503 &mut factory,
504 &runtime.sequences,
505 None,
506 DEFAULT_WRITE_TIMEOUT,
507 logout,
508 )
509 .await;
510 let _ = session.reject_logon();
511 return Err(EngineError::LogonRejected { reason: detail });
512 }
513 };
514
515 // The clock is checked next, before the heartbeat is set: an
516 // initiator whose SendingTime is wildly skewed cannot be trusted to
517 // sequence a session, and the handshake is the cheapest place to
518 // refuse it.
519 if let Err(problem) = sending_time.validate(&raw) {
520 let detail = problem.to_string();
521 let reason = problem.reject_reason();
522 let reject = factory.session_reject(logon_seq, MsgType::Logon.as_str(), &reason);
523 let _ = send_handshake_admin(
524 self.application.as_ref(),
525 &session_id,
526 &mut framed,
527 &mut factory,
528 &runtime.sequences,
529 None,
530 DEFAULT_WRITE_TIMEOUT,
531 reject,
532 )
533 .await;
534 let logout = factory.logout(Some(&detail));
535 let _ = send_handshake_admin(
536 self.application.as_ref(),
537 &session_id,
538 &mut framed,
539 &mut factory,
540 &runtime.sequences,
541 None,
542 DEFAULT_WRITE_TIMEOUT,
543 logout,
544 )
545 .await;
546 let _ = session.reject_logon();
547 return Err(EngineError::SendingTime { detail });
548 }
549
550 // Authentication hook: the application inspects the Logon (Username
551 // 553 / Password 554 and the like) and may refuse it.
552 if let Err(reason) = self.application.from_admin(&raw, &session_id).await {
553 let logout = factory.logout(Some(&reason.text));
554 let _ = send_handshake_admin(
555 self.application.as_ref(),
556 &session_id,
557 &mut framed,
558 &mut factory,
559 &runtime.sequences,
560 None,
561 DEFAULT_WRITE_TIMEOUT,
562 logout,
563 )
564 .await;
565 let _ = session.reject_logon();
566 return Err(EngineError::LogonRejected {
567 reason: reason.text,
568 });
569 }
570
571 // Honor the heartbeat interval the initiator requested (HeartBtInt,
572 // 108): the acceptor adopts it and echoes it back on the reply. The
573 // value is counterparty-controlled and drives a Duration on the
574 // heartbeat clock, so it is bounded here — a missing, non-numeric,
575 // or over-range 108 fails the handshake rather than being defaulted
576 // or, at u64::MAX, overflowing the clock and aborting the process.
577 let requested_heartbeat_secs = match wire::parse_heartbeat_interval(&raw) {
578 Ok(secs) => secs,
579 Err(problem) => {
580 let detail = problem.to_string();
581 tracing::warn!(session = %session_id, detail = %detail, "rejecting Logon HeartBtInt");
582 let logout = factory.logout(Some(&detail));
583 let _ = send_handshake_admin(
584 self.application.as_ref(),
585 &session_id,
586 &mut framed,
587 &mut factory,
588 &runtime.sequences,
589 None,
590 DEFAULT_WRITE_TIMEOUT,
591 logout,
592 )
593 .await;
594 let _ = session.reject_logon();
595 return Err(EngineError::LogonRejected { reason: detail });
596 }
597 };
598 {
599 let mut heartbeat = lock_heartbeat(&runtime);
600 heartbeat.on_message_received(false, None);
601 if Duration::from_secs(requested_heartbeat_secs) != heartbeat.interval() {
602 tracing::info!(
603 session = %session_id,
604 heartbeat_secs = requested_heartbeat_secs,
605 "honoring HeartBtInt requested by counterparty"
606 );
607 *heartbeat =
608 HeartbeatManager::new(Duration::from_secs(requested_heartbeat_secs));
609 }
610 }
611 let reply_heartbeat_secs = requested_heartbeat_secs;
612
613 // ResetSeqNumFlag (141): the acceptor resets when the peer asks
614 // (141=Y) *or* when it is locally configured to (`reset_on_logon`).
615 // Whichever drives the reset, the ack must signal it with 141=Y so
616 // the peer resets in lockstep — a local reset that seeded fresh
617 // counters but acked without 141 would silently desync the peer.
618 // The reset is applied before MsgSeqNum is validated, otherwise the
619 // Logon's 34=1 reads as fatally too low against continuity-seeded
620 // counters.
621 let inbound_reset = raw.get_field_str(141) == Some("Y");
622 let reset = inbound_reset || self.config.reset_on_logon;
623 if inbound_reset && logon_seq != 1 {
624 // FIX requires MsgSeqNum = 1 on a Logon that itself carries
625 // ResetSeqNumFlag = Y: the reset and the number carrying it have
626 // to describe the same stream. This binds only when the *peer*
627 // declared the reset; a purely local reset places no such
628 // requirement on the number the peer chose.
629 let _ = session.reject_logon();
630 return Err(EngineError::Sequence(format!(
631 "Logon set ResetSeqNumFlag=Y but carried MsgSeqNum {logon_seq}, not 1"
632 )));
633 }
634 if reset {
635 tracing::info!(
636 session = %session_id,
637 inbound_reset,
638 "resetting sequence numbers on Logon"
639 );
640 runtime.sequences.reset();
641 }
642
643 match runtime.sequences.validate_incoming(logon_seq) {
644 SequenceResult::Ok => {
645 if let Err(err) = runtime.sequences.try_increment_target_seq() {
646 let _ = session.reject_logon();
647 return Err(err.into());
648 }
649 }
650 SequenceResult::TooLow { expected, received } => {
651 let detail = format!(
652 "logon MsgSeqNum too low: expected {expected}, received {received}"
653 );
654 let logout = factory.logout(Some(&detail));
655 let _ = send_handshake_admin(
656 self.application.as_ref(),
657 &session_id,
658 &mut framed,
659 &mut factory,
660 &runtime.sequences,
661 None,
662 DEFAULT_WRITE_TIMEOUT,
663 logout,
664 )
665 .await;
666 let _ = session.reject_logon();
667 return Err(EngineError::Sequence(detail));
668 }
669 SequenceResult::Gap { expected, received } => {
670 pending_resend = Some((expected, received));
671 }
672 }
673
674 // Reply with the Logon acknowledgement, mirroring ResetSeqNumFlag.
675 // The peek-then-spend `send_handshake_admin` frames it under the next
676 // sender sequence number (1 after a reset) and spends that number
677 // only once the frame is built.
678 let logon_ack = factory.logon(reply_heartbeat_secs, reset);
679 if let Err(err) = send_handshake_admin(
680 self.application.as_ref(),
681 &session_id,
682 &mut framed,
683 &mut factory,
684 &runtime.sequences,
685 None,
686 DEFAULT_WRITE_TIMEOUT,
687 logon_ack,
688 )
689 .await
690 {
691 let _ = session.reject_logon();
692 return Err(err);
693 }
694 lock_heartbeat(&runtime).on_message_sent();
695 }
696
697 // Typestate: LogonReceived -> Active.
698 let session = session.accept_logon();
699 self.application.on_logon(&session_id).await;
700 tracing::info!(session = %session_id, "FIX session established (acceptor)");
701
702 // A gap in the inbound Logon means we missed messages: request a resend
703 // now that the session is Active.
704 if let Some((expected, _high_water)) = pending_resend {
705 let request = factory.resend_request(expected, 0);
706 send_handshake_admin(
707 self.application.as_ref(),
708 &session_id,
709 &mut framed,
710 &mut factory,
711 &runtime.sequences,
712 None,
713 DEFAULT_WRITE_TIMEOUT,
714 request,
715 )
716 .await?;
717 lock_heartbeat(&runtime).on_message_sent();
718 }
719
720 let resend = pending_resend
721 .map(|(expected, high_water)| ResendState::first(expected, high_water, &self.config));
722 // Hand the framed socket to the shared reactor, the same one the
723 // initiator uses. The admission guard rides the reactor task and frees
724 // the slot when the session closes. The acceptor attaches no store, so
725 // resends are gap-filled.
726 let params = SessionParams {
727 framed,
728 session,
729 runtime: Arc::clone(&runtime),
730 factory,
731 identity,
732 sending_time,
733 config: self.config.clone(),
734 application: Arc::clone(&self.application),
735 session_id: session_id.clone(),
736 store: None,
737 resend,
738 write_timeout: DEFAULT_WRITE_TIMEOUT,
739 outbound_capacity: self.outbound_capacity,
740 app_queue_capacity: DEFAULT_APP_QUEUE_CAPACITY,
741 };
742 let (command_tx, closed_rx) = spawn_session(params, admission);
743
744 Ok(Connection {
745 session_id,
746 commands: command_tx,
747 closed: closed_rx,
748 runtime,
749 })
750 }
751}