Skip to main content

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}