Skip to main content

vane_core/
worker.rs

1//! Worker — the pinned, share-nothing thread that owns an engine, a buffer
2//! pool, a session slab, and a handler (`TH-01`, `MM-01`).
3//!
4//! State is split into [`WorkerState`] (transport) and the [`Handler`] so a
5//! handler callback can re-enter the transport through [`SessionIo`] without
6//! aliasing: `SessionIo` borrows only the transport half.
7
8use std::collections::BinaryHeap;
9use std::io;
10use std::net::{SocketAddr, TcpListener as StdTcpListener};
11use std::os::fd::{AsRawFd, RawFd};
12use std::path::PathBuf;
13use std::sync::Arc;
14use std::sync::atomic::{AtomicBool, Ordering};
15use std::time::{Duration, Instant};
16
17use vane_observe::LogEvent;
18use vane_observe::metrics::Registry;
19use vane_observe::ring::EventRing;
20
21use crate::buffer::{BufferPool, DEFAULT_BUF_SIZE, DEFAULT_POOL_SIZE};
22use crate::engine::{Engine, create_engine};
23use crate::handler::{Handler, HandlerFactory, Mode, SessionIo};
24use crate::net::{StreamFd, set_nodelay, shutdown_write};
25
26/// Per-session queued-write budget (burst cap). Steady-state flow
27/// control is read throttling (2 slots); this bounds a single handler
28/// burst (e.g. a large pre-connect body queue).
29const WRITE_PENDING_CAP: usize = 1024 * 1024;
30use crate::slab::SessionSlab;
31use crate::spsc::{SpscReceiver, SpscSender};
32use crate::token::{Op, Token};
33
34/// Commands the control plane sends a worker.
35#[derive(Debug, Clone, Copy, PartialEq, Eq)]
36pub enum WorkerCmd {
37    /// Stop accepting; drain sessions; exit by the deadline.
38    Shutdown {
39        /// Hard exit horizon in milliseconds from receipt.
40        deadline_ms: u64,
41    },
42}
43
44/// Worker configuration knobs.
45#[derive(Debug, Clone)]
46pub struct WorkerConfig {
47    /// CPU core index to pin to (physical order from `core_affinity`).
48    pub core: Option<usize>,
49    /// Max concurrent sessions.
50    pub max_sessions: usize,
51    /// Buffer pool slots (double as io_uring fixed buffers).
52    pub pool_slots: usize,
53    /// io_uring queue depth.
54    pub ring_entries: u32,
55    /// Use `IORING_SETUP_SQPOLL` (zero-syscall submission).
56    pub sqpoll: bool,
57    /// Prefer mio even when io_uring is available.
58    pub force_mio: bool,
59    /// Accept backlog per listener.
60    pub backlog: i32,
61}
62
63impl Default for WorkerConfig {
64    fn default() -> Self {
65        Self {
66            core: None,
67            max_sessions: 16_384,
68            pool_slots: DEFAULT_POOL_SIZE,
69            ring_entries: 4096,
70            sqpoll: false,
71            force_mio: false,
72            backlog: 4096,
73        }
74    }
75}
76
77/// Immutable context handed to a [`HandlerFactory::build`].
78pub struct WorkerCtx {
79    /// Worker index (0..N).
80    pub id: usize,
81    /// Shared metric registry.
82    pub registry: Arc<Registry>,
83    /// Worker-local log/event ring (drained by the control plane).
84    pub events: Arc<EventRing<LogEvent, { vane_observe::EVENT_RING_CAPACITY }>>,
85    /// Effective worker config.
86    pub config: WorkerConfig,
87    /// Live-session gauge (worker-maintained).
88    pub connections: vane_observe::metrics::MetricHandle,
89}
90
91/// Live handle to a running worker thread.
92pub struct WorkerHandle {
93    /// Worker index.
94    pub id: usize,
95    /// Command sender (shutdown).
96    pub cmd: SpscSender<WorkerCmd, 64>,
97    /// Set when the worker exits.
98    done: Arc<AtomicBool>,
99    join: Option<std::thread::JoinHandle<()>>,
100}
101
102impl WorkerHandle {
103    /// `true` once the worker thread has exited.
104    #[must_use]
105    pub fn is_done(&self) -> bool {
106        self.done.load(Ordering::Acquire)
107    }
108
109    /// Joins the worker thread.
110    pub fn join(&mut self) {
111        if let Some(h) = self.join.take() {
112            let _ = h.join();
113        }
114    }
115}
116
117/// Internal per-session state.
118struct Session {
119    downstream: StreamFd,
120    peer: SocketAddr,
121    upstream: Option<StreamFd>,
122    /// Read slot held across an in-flight downstream read.
123    rslot: Option<u32>,
124    /// Downstream write queue: `(slot, len, offset)`.
125    wq: std::collections::VecDeque<(u32, u32, u32)>,
126    /// Upstream read slot.
127    urslot: Option<u32>,
128    /// Upstream write queue.
129    uwq: std::collections::VecDeque<(u32, u32, u32)>,
130    /// Bytes queued while the pool was exhausted (runaway-guarded).
131    pending_down: Vec<u8>,
132    pending_up: Vec<u8>,
133    read_inflight: bool,
134    upstream_read_inflight: bool,
135    write_inflight: bool,
136    upstream_write_inflight: bool,
137    /// Nonblocking connect still in progress (diagnostics).
138    #[allow(dead_code)] // consumed by CQE dispatch; kept for state clarity
139    connect_inflight: bool,
140    /// Bumped whenever the upstream fd is replaced; rides the token aux so
141    /// stale completions for a discarded fd are dropped.
142    upstream_epoch: u8,
143    /// Request head already serialized to the upstream (no failover past it).
144    request_sent: bool,
145    downstream_eof: bool,
146    upstream_eof: bool,
147    /// Half-close (FIN) requested while writes are still flushing.
148    fin_queued: bool,
149    splice: bool,
150    deadline: Option<Instant>,
151}
152
153/// Timer heap entry.
154#[derive(PartialEq, Eq)]
155struct Timer {
156    at: Instant,
157    slot: u32,
158    generation: u16,
159    reason_aux: u16,
160}
161
162impl Ord for Timer {
163    fn cmp(&self, other: &Self) -> std::cmp::Ordering {
164        other.at.cmp(&self.at) // min-heap
165    }
166}
167impl PartialOrd for Timer {
168    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
169        Some(self.cmp(other))
170    }
171}
172
173/// Transport half of the worker. Every handler callback receives a
174/// [`SessionIo`] that borrows exactly this struct.
175pub(crate) struct WorkerState {
176    ctx: WorkerCtx,
177    engine: Box<dyn Engine>,
178    pool: BufferPool,
179    slab: SessionSlab<Session>,
180    listeners: Vec<StdTcpListener>,
181    timers: BinaryHeap<Timer>,
182    draining: bool,
183    drain_deadline: Option<Instant>,
184    #[allow(dead_code)] // kept for per-mode logging in follow-ups
185    mode: Mode,
186    /// Connects that completed inline (UDS always; TCP rarely) and need
187    /// their `on_upstream_connected` callback delivered on the next loop
188    /// pass — the handler is not reachable from the connect call site.
189    inline_connected: Vec<(u32, u16)>,
190}
191
192impl WorkerState {
193    fn io_for(&mut self, slot: u32, generation: u16) -> SessionIo<'_> {
194        SessionIo {
195            worker: self,
196            slot,
197            generation,
198        }
199    }
200
201    fn valid(&self, slot: u32, generation: u16) -> bool {
202        self.slab.matches(slot, generation)
203    }
204
205    // ----- SessionIo surface (called by handlers) ----------------------------
206
207    pub(crate) fn downstream_write(&mut self, slot: u32, generation: u16, bytes: &[u8]) {
208        if !self.valid(slot, generation) || bytes.is_empty() {
209            return;
210        }
211        // Direct write requires the full queue to be idle: pending_down
212        // may hold OLDER bytes (pool exhaustion defers them), and
213        // writing the new batch first would reorder the stream.
214        let can_direct = self
215            .slab
216            .get(slot)
217            .is_some_and(|s| !s.write_inflight && s.wq.is_empty() && s.pending_down.is_empty());
218        if can_direct {
219            if let Some(ws) = self.pool.take() {
220                let Some(s) = self.slab.get_mut(slot) else {
221                    return;
222                };
223                let n = bytes.len().min(self.pool.buf_size());
224                self.pool.slot_mut(ws)[..n].copy_from_slice(&bytes[..n]);
225                if bytes.len() > n {
226                    s.pending_down.extend_from_slice(&bytes[n..]);
227                }
228                s.wq.push_back((ws, n as u32, 0));
229                s.write_inflight = true;
230                let fd = s.downstream.fd();
231                let token = Token::new(Op::DownstreamWrite, slot, generation, 0);
232                let _ = self.engine.write(token, fd, ws, n, 0);
233                return;
234            }
235        }
236        // Queue path (pool exhausted or write already in flight).
237        let overflow = self
238            .slab
239            .get(slot)
240            .is_some_and(|s| s.pending_down.len() + bytes.len() > WRITE_PENDING_CAP);
241        if overflow {
242            self.close_session(slot, generation, "write-queue-overflow");
243            return;
244        }
245        if let Some(s) = self.slab.get_mut(slot) {
246            s.pending_down.extend_from_slice(bytes);
247        }
248    }
249
250    pub(crate) fn upstream_write(&mut self, slot: u32, generation: u16, bytes: &[u8]) {
251        if !self.valid(slot, generation) || bytes.is_empty() {
252            return;
253        }
254        // Same queue-jump guard as downstream_write: pending_up may
255        // hold older deferred bytes.
256        let can_direct = self.slab.get(slot).is_some_and(|s| {
257            s.upstream.is_some()
258                && !s.upstream_write_inflight
259                && s.uwq.is_empty()
260                && s.pending_up.is_empty()
261        });
262        if can_direct {
263            if let Some(ws) = self.pool.take() {
264                let Some(s) = self.slab.get_mut(slot) else {
265                    return;
266                };
267                let Some(up) = &s.upstream else { return };
268                let fd = up.fd();
269                let n = bytes.len().min(self.pool.buf_size());
270                self.pool.slot_mut(ws)[..n].copy_from_slice(&bytes[..n]);
271                if bytes.len() > n {
272                    s.pending_up.extend_from_slice(&bytes[n..]);
273                }
274                s.uwq.push_back((ws, n as u32, 0));
275                s.upstream_write_inflight = true;
276                let token = Token::new(Op::UpstreamWrite, slot, generation, 0);
277                let _ = self.engine.write(token, fd, ws, n, 0);
278                return;
279            }
280        }
281        let overflow = self
282            .slab
283            .get(slot)
284            .is_some_and(|s| s.pending_up.len() + bytes.len() > WRITE_PENDING_CAP);
285        if overflow {
286            self.close_session(slot, generation, "write-queue-overflow");
287            return;
288        }
289        if let Some(s) = self.slab.get_mut(slot) {
290            s.pending_up.extend_from_slice(bytes);
291        }
292    }
293
294    pub(crate) fn connect_upstream(
295        &mut self,
296        slot: u32,
297        generation: u16,
298        addr: SocketAddr,
299    ) -> bool {
300        if !self.valid(slot, generation) {
301            return false;
302        }
303        let token = Token::new(Op::Connect, slot, generation, 0);
304        match self.engine.connect(token, addr) {
305            Ok((fd, poll)) => {
306                let Some(s) = self.slab.get_mut(slot) else {
307                    return false;
308                };
309                crate::dbg_trace!("DIAL fd={fd} slot={slot}");
310                s.upstream = Some(StreamFd(fd));
311                if poll == crate::engine::Poll::Done(0) {
312                    self.do_upstream_connected(slot, generation);
313                    self.inline_connected.push((slot, generation));
314                }
315                true
316            }
317            Err(_) => false,
318        }
319    }
320
321    pub(crate) fn connect_upstream_unix(
322        &mut self,
323        slot: u32,
324        generation: u16,
325        path: PathBuf,
326    ) -> bool {
327        if !self.valid(slot, generation) {
328            return false;
329        }
330        let token = Token::new(Op::Connect, slot, generation, 0);
331        match self.engine.connect_unix(token, &path) {
332            Ok((fd, poll)) => {
333                let Some(s) = self.slab.get_mut(slot) else {
334                    return false;
335                };
336                crate::dbg_trace!("DIAL fd={fd} slot={slot}");
337                s.upstream = Some(StreamFd(fd));
338                if poll == crate::engine::Poll::Done(0) {
339                    self.do_upstream_connected(slot, generation);
340                    self.inline_connected.push((slot, generation));
341                }
342                true
343            }
344            Err(_) => false,
345        }
346    }
347
348    pub(crate) fn start_splice(&mut self, slot: u32, generation: u16) -> bool {
349        if !self.valid(slot, generation) {
350            return false;
351        }
352        let Some(s) = self.slab.get_mut(slot) else {
353            return false;
354        };
355        let Some(up) = &s.upstream else { return false };
356        let (a, b) = (s.downstream.fd(), up.fd());
357        // Splice bypasses the buffer pool. Slots with in-flight reads are
358        // KEPT: their CQEs will land after the pump starts and must be
359        // forwarded raw (see the splice branch in the read dispatchers) —
360        // releasing them now would drop the bytes they hold.
361        if !s.read_inflight {
362            if let Some(rs) = s.rslot.take() {
363                self.pool.release(rs);
364            }
365        }
366        if !s.upstream_read_inflight {
367            if let Some(rs) = s.urslot.take() {
368                self.pool.release(rs);
369            }
370        }
371        s.splice = true;
372        let ta = Token::new(Op::Splice, slot, generation, 1);
373        let tb = Token::new(Op::Splice, slot, generation, 2);
374        self.engine.splice_pump(ta, a, tb, b).is_ok()
375    }
376
377    pub(crate) fn shutdown_downstream_write(&mut self, slot: u32, generation: u16) {
378        if !self.valid(slot, generation) {
379            return;
380        }
381        let Some(s) = self.slab.get_mut(slot) else {
382            return;
383        };
384        if s.wq.is_empty() && s.pending_down.is_empty() && !s.write_inflight {
385            shutdown_write(s.downstream.fd());
386        } else {
387            s.fin_queued = true; // flush paths complete it
388        }
389    }
390
391    pub(crate) fn shutdown_upstream_write(&mut self, slot: u32, generation: u16) {
392        if !self.valid(slot, generation) {
393            return;
394        }
395        let Some(s) = self.slab.get_mut(slot) else {
396            return;
397        };
398        if let Some(up) = &s.upstream {
399            if s.uwq.is_empty() && s.pending_up.is_empty() && !s.upstream_write_inflight {
400                shutdown_write(up.fd());
401            }
402        }
403    }
404
405    /// Closes and drops a session. `reason` documents the call site for
406    /// crash triage (compiled out in release... kept as a parameter so the
407    /// next debug session can re-enable the trace in one line).
408    pub(crate) fn close_session(&mut self, slot: u32, generation: u16, reason: &str) {
409        // vane-core stays dependency-light: no tracing here.
410        let _ = reason;
411        if !self.valid(slot, generation) {
412            return;
413        }
414        let Some(s) = self.slab.remove(slot) else {
415            return;
416        }; // generation bumps inside
417        if let Some(rs) = s.rslot {
418            self.pool.release(rs);
419        }
420        for (ws, _, _) in s.wq {
421            self.pool.release(ws);
422        }
423        if let Some(rs) = s.urslot {
424            self.pool.release(rs);
425        }
426        for (ws, _, _) in s.uwq {
427            self.pool.release(ws);
428        }
429        // Drain unread downstream bytes before close: closing a socket
430        // with pending receive data sends RST, which wipes any frames
431        // we just wrote (e.g. GOAWAY/RST_STREAM error responses) before
432        // the peer reads them.
433        let dfd = s.downstream.fd();
434        let mut scratch = [0u8; 4096];
435        loop {
436            let n = unsafe { libc::read(dfd, scratch.as_mut_ptr().cast(), scratch.len()) };
437            if n <= 0 {
438                break;
439            }
440        }
441        self.engine.remove(s.downstream.fd());
442        if let Some(up) = &s.upstream {
443            self.engine.remove(up.fd());
444        }
445        // StreamFd Drop closes the descriptors.
446    }
447
448    pub(crate) fn set_deadline(
449        &mut self,
450        slot: u32,
451        generation: u16,
452        at: Option<Instant>,
453        reason: crate::handler::DeadlineReason,
454    ) {
455        if !self.valid(slot, generation) {
456            return;
457        }
458        let Some(s) = self.slab.get_mut(slot) else {
459            return;
460        };
461        s.deadline = at;
462        if let Some(at) = at {
463            self.timers.push(Timer {
464                at,
465                slot,
466                generation,
467                reason_aux: reason.aux(),
468            });
469        }
470    }
471
472    pub(crate) fn peer_of(&self, slot: u32, generation: u16) -> Option<SocketAddr> {
473        if self.valid(slot, generation) {
474            self.slab.get(slot).map(|s| s.peer)
475        } else {
476            None
477        }
478    }
479
480    pub(crate) fn has_upstream(&self, slot: u32, generation: u16) -> bool {
481        self.valid(slot, generation) && self.slab.get(slot).is_some_and(|s| s.upstream.is_some())
482    }
483
484    /// Discards a dead upstream: removes it from the engine, closes the fd,
485    /// releases its slots, and bumps the epoch so stale CQEs are dropped.
486    pub(crate) fn discard_upstream(&mut self, slot: u32, generation: u16) {
487        if !self.valid(slot, generation) {
488            return;
489        }
490        let Some(s) = self.slab.get_mut(slot) else {
491            return;
492        };
493        let Some(old) = s.upstream.take() else { return };
494        let fd = old.fd();
495        // StreamFd Drop closes it — ownership is simply dropped.
496        drop(old);
497        self.engine.remove(fd);
498        if let Some(rs) = s.urslot.take() {
499            self.pool.release(rs);
500        }
501        for (ws, _, _) in s.uwq.drain(..) {
502            self.pool.release(ws);
503        }
504        s.pending_up.clear();
505        s.upstream_read_inflight = false;
506        s.upstream_write_inflight = false;
507        s.upstream_eof = false;
508        s.request_sent = false;
509        s.upstream_epoch = s.upstream_epoch.wrapping_add(1);
510    }
511
512    /// Detaches the upstream fd (no close) for connection pooling.
513    pub(crate) fn detach_upstream(&mut self, slot: u32, generation: u16) -> Option<RawFd> {
514        if !self.valid(slot, generation) {
515            return None;
516        }
517        let fd = {
518            let s = self.slab.get_mut(slot)?;
519            let owned = s.upstream.take()?;
520            let fd = owned.fd();
521            // Ownership transfers to the caller (pool) — forget the guard
522            // so Drop does not close the descriptor.
523            std::mem::forget(owned);
524
525            fd
526        };
527        // Stop tracking it in the engine and return every slot it holds.
528        self.engine.remove(fd);
529        if let Some(s) = self.slab.get_mut(slot) {
530            if let Some(rs) = s.urslot.take() {
531                self.pool.release(rs);
532            }
533            s.upstream_read_inflight = false;
534            for (ws, _, _) in s.uwq.drain(..) {
535                self.pool.release(ws);
536            }
537            s.pending_up.clear();
538            s.upstream_write_inflight = false;
539            s.upstream_eof = false;
540            s.request_sent = false;
541        }
542        Some(fd)
543    }
544
545    /// Attaches a pooled fd as this session's upstream and arms its read.
546    pub(crate) fn attach_upstream(&mut self, slot: u32, generation: u16, fd: RawFd) -> bool {
547        if !self.valid(slot, generation) {
548            return false;
549        }
550        let token = Token::new(Op::UpstreamRead, slot, generation, 0);
551        if self.engine.add_stream(fd, token).is_err() {
552            return false;
553        }
554        let Some(s) = self.slab.get_mut(slot) else {
555            return false;
556        };
557        crate::dbg_trace!("DIAL fd={fd} slot={slot}");
558        s.upstream = Some(StreamFd(fd));
559        if s.urslot.is_none() {
560            if let Some(uslot) = self.pool.take() {
561                let epoch = u16::from(s.upstream_epoch);
562                s.urslot = Some(uslot);
563                s.upstream_read_inflight = true;
564                let token = Token::new(Op::UpstreamRead, slot, generation, epoch);
565                let _ = self.engine.read(token, fd, uslot);
566            }
567        }
568        true
569    }
570
571    /// Raw upstream descriptor (diagnostics).
572    pub(crate) fn upstream_fd(&self, slot: u32, generation: u16) -> Option<RawFd> {
573        if self.valid(slot, generation) {
574            self.slab
575                .get(slot)
576                .and_then(|s| s.upstream.as_ref())
577                .map(StreamFd::fd)
578        } else {
579            None
580        }
581    }
582
583    /// Request-head-written flag (failover guard).
584    pub(crate) fn request_sent(&self, slot: u32, generation: u16) -> bool {
585        self.valid(slot, generation) && self.slab.get(slot).is_some_and(|s| s.request_sent)
586    }
587
588    pub(crate) fn mark_request_sent(&mut self, slot: u32, generation: u16) {
589        if !self.valid(slot, generation) {
590            return;
591        }
592        if let Some(s) = self.slab.get_mut(slot) {
593            s.request_sent = true;
594        }
595    }
596
597    // ----- internal event handling (called with the handler split out) -------
598
599    fn do_upstream_connected(&mut self, slot: u32, generation: u16) {
600        {
601            let Some(s) = self.slab.get_mut(slot) else {
602                return;
603            };
604            if s.urslot.is_none() && !s.splice {
605                if let Some(uslot) = self.pool.take() {
606                    let epoch = u16::from(s.upstream_epoch);
607                    let fd = s.upstream.as_ref().map_or(-1, StreamFd::fd);
608                    s.urslot = Some(uslot);
609                    s.upstream_read_inflight = true;
610                    let token = Token::new(Op::UpstreamRead, slot, generation, epoch);
611                    let _ = self.engine.read(token, fd, uslot);
612                }
613            }
614        }
615        // Bytes buffered before the dial completed (early client data):
616        // kick the write queue now that the upstream exists, or they
617        // would sit until the next write-completion event (which may
618        // never come — nothing was in flight).
619        let needs_kick = self
620            .slab
621            .get(slot)
622            .is_some_and(|s| !s.pending_up.is_empty() && !s.upstream_write_inflight);
623        if needs_kick {
624            self.kick_pending_upstream(slot, generation);
625        }
626    }
627
628    fn arm_downstream_read(&mut self, slot: u32, generation: u16) {
629        // NOTE: no upstream-write backpressure pause here. Skipping the
630        // read while paused loses epoll ET edges (a readable event
631        // consumed with no parked read op is gone forever), which
632        // deadlocks streaming responses. Upstream buffering is already
633        // bounded by WRITE_PENDING_CAP at write_upstream time.
634        let Some(s) = self.slab.get_mut(slot) else {
635            return;
636        };
637        if s.read_inflight || s.splice {
638            crate::dbg_trace!(
639                "ARMDbg skip slot={slot} inflight={} splice={}",
640                s.read_inflight,
641                s.splice
642            );
643            return;
644        }
645        if s.rslot.is_none() {
646            let Some(rs) = self.pool.take() else {
647                crate::dbg_trace!("ARMDbg pool-exhausted slot={slot}");
648                return;
649            }; // backpressure
650            s.rslot = Some(rs);
651        }
652        let Some(rs) = s.rslot else { return };
653        s.read_inflight = true;
654        let fd = s.downstream.fd();
655        let token = Token::new(Op::DownstreamRead, slot, generation, 0);
656        let _ = self.engine.read(token, fd, rs);
657    }
658
659    fn arm_upstream_read(&mut self, slot: u32, generation: u16) {
660        let Some(s) = self.slab.get_mut(slot) else {
661            return;
662        };
663        // The handler may have detached (parked) or discarded the upstream
664        // during its callback — there is nothing to read anymore. Re-arming
665        // here would submit a read on fd -1 / a closed descriptor whose
666        // EBADF completion then kills the idle keep-alive session.
667        if s.upstream.is_none() {
668            return;
669        }
670        // Read throttling: pause while the client-bound write queue is
671        // backed up (resumed from the downstream write-completion path).
672        if s.pending_down.len() > 2 * self.pool.buf_size() {
673            return;
674        }
675        if s.upstream_read_inflight || s.splice {
676            return;
677        }
678        if s.urslot.is_none() {
679            let Some(rs) = self.pool.take() else { return };
680            s.urslot = Some(rs);
681        }
682        let Some(rs) = s.urslot else { return };
683        s.upstream_read_inflight = true;
684        let epoch = u16::from(s.upstream_epoch);
685        let fd = s.upstream.as_ref().map_or(-1, StreamFd::fd);
686        let token = Token::new(Op::UpstreamRead, slot, generation, epoch);
687        let _ = self.engine.read(token, fd, rs);
688    }
689
690    fn continue_downstream_write(&mut self, slot: u32, generation: u16, h: &mut dyn Handler) {
691        let next = {
692            let Some(s) = self.slab.get_mut(slot) else {
693                return;
694            };
695            if let Some((ws, _, _)) = s.wq.front() {
696                let ws = *ws;
697                s.wq.pop_front();
698                self.pool.release(ws);
699            }
700            if let Some(&(ws, len, off)) = s.wq.front() {
701                s.write_inflight = true;
702                Some((ws, len, off))
703            } else if !s.pending_down.is_empty() {
704                let Some(ws) = self.pool.take() else {
705                    s.write_inflight = false;
706                    return;
707                };
708                let n = s.pending_down.len().min(self.pool.buf_size());
709                self.pool.slot_mut(ws)[..n].copy_from_slice(&s.pending_down[..n]);
710                s.pending_down.drain(..n);
711                s.wq.push_back((ws, n as u32, 0));
712                s.write_inflight = true;
713                Some((ws, n as u32, 0))
714            } else {
715                s.write_inflight = false;
716                None
717            }
718        };
719        if let Some((ws, len, off)) = next {
720            let Some(s) = self.slab.get_mut(slot) else {
721                return;
722            };
723            let fd = s.downstream.fd();
724            let token = Token::new(Op::DownstreamWrite, slot, generation, 0);
725            let _ = self.engine.write(token, fd, ws, len as usize, off as usize);
726        } else {
727            // Queue fully flushed: complete a deferred FIN.
728            let fin = {
729                let Some(s) = self.slab.get_mut(slot) else {
730                    return;
731                };
732                if s.fin_queued {
733                    s.fin_queued = false;
734                    shutdown_write(s.downstream.fd());
735                    true
736                } else {
737                    false
738                }
739            };
740            if !fin {
741                let mut io = self.io_for(slot, generation);
742                h.on_downstream_flushed(&mut io);
743            }
744        }
745        // Queue drained: resume client reads paused by backpressure.
746        let drained = self
747            .slab
748            .get(slot)
749            .is_some_and(|s| s.pending_up.len() <= 2 * self.pool.buf_size());
750        if drained {
751            self.arm_downstream_read(slot, generation);
752        }
753    }
754
755    /// Kicks pending_up bytes into the upstream write queue (used after
756    /// a dial completes with buffered early data). Public within the
757    /// worker for `do_upstream_connected`.
758    fn kick_pending_upstream(&mut self, slot: u32, generation: u16) {
759        let next = {
760            let Some(s) = self.slab.get_mut(slot) else {
761                return;
762            };
763            if s.uwq.front().is_some() || s.pending_up.is_empty() || s.upstream_write_inflight {
764                return;
765            }
766            let Some(ws) = self.pool.take() else {
767                return; // backpressure: flushed on the next completion
768            };
769            let n = s.pending_up.len().min(self.pool.buf_size());
770            self.pool.slot_mut(ws)[..n].copy_from_slice(&s.pending_up[..n]);
771            s.pending_up.drain(..n);
772            s.uwq.push_back((ws, n as u32, 0));
773            s.upstream_write_inflight = true;
774            Some((ws, n as u32, 0))
775        };
776        if let Some((ws, len, off)) = next {
777            let Some(s) = self.slab.get_mut(slot) else {
778                return;
779            };
780            let Some(up) = &s.upstream else { return };
781            let fd = up.fd();
782            let token = Token::new(Op::UpstreamWrite, slot, generation, 0);
783            let _ = self.engine.write(token, fd, ws, len as usize, off as usize);
784        }
785    }
786
787    fn continue_upstream_write(&mut self, slot: u32, generation: u16, h: &mut dyn Handler) {
788        let next = {
789            let Some(s) = self.slab.get_mut(slot) else {
790                return;
791            };
792            if let Some((ws, _, _)) = s.uwq.front() {
793                let ws = *ws;
794                s.uwq.pop_front();
795                self.pool.release(ws);
796            }
797            if let Some(&(ws, len, off)) = s.uwq.front() {
798                s.upstream_write_inflight = true;
799                Some((ws, len, off))
800            } else if !s.pending_up.is_empty() {
801                let Some(ws) = self.pool.take() else {
802                    s.upstream_write_inflight = false;
803                    return;
804                };
805                let n = s.pending_up.len().min(self.pool.buf_size());
806                self.pool.slot_mut(ws)[..n].copy_from_slice(&s.pending_up[..n]);
807                s.pending_up.drain(..n);
808                s.uwq.push_back((ws, n as u32, 0));
809                s.upstream_write_inflight = true;
810                Some((ws, n as u32, 0))
811            } else {
812                s.upstream_write_inflight = false;
813                None
814            }
815        };
816        if let Some((ws, len, off)) = next {
817            let Some(s) = self.slab.get_mut(slot) else {
818                return;
819            };
820            let Some(up) = &s.upstream else { return };
821            let fd = up.fd();
822            let token = Token::new(Op::UpstreamWrite, slot, generation, 0);
823            let _ = self.engine.write(token, fd, ws, len as usize, off as usize);
824        } else {
825            // Queue fully drained: resume client reads that were paused
826            // by the pending_up backpressure (arm_downstream_read's
827            // pause condition). Without this the pause never lifts —
828            // the client's socket goes unread and the connection
829            // deadlocks once its send window exhausts.
830            self.arm_downstream_read(slot, generation);
831            let mut io = self.io_for(slot, generation);
832            h.on_upstream_flushed(&mut io);
833        }
834    }
835
836    fn maybe_finish(&mut self, slot: u32, generation: u16) {
837        let done = self.slab.get(slot).is_some_and(|s| {
838            s.downstream_eof
839                // A session whose upstream half never existed (client
840                // connected + dropped before any dial) must close too —
841                // otherwise it leaks its fd and its mio registration,
842                // and the REUSED fd number later collides with the
843                // stale registration (dial failures, hung sessions).
844                && (s.upstream_eof || s.upstream.is_none())
845                && s.wq.is_empty()
846                && s.pending_down.is_empty()
847                && s.uwq.is_empty()
848                && s.pending_up.is_empty()
849                && !s.write_inflight
850                && !s.upstream_write_inflight
851        });
852        if done {
853            self.close_session(slot, generation, "both-eof-flushed");
854        }
855    }
856
857    fn dispatch_cqe(&mut self, cqe: crate::engine::Cqe, h: &mut dyn Handler) {
858        let (slot, generation) = (cqe.token.slot(), cqe.token.generation());
859        if !self.valid(slot, generation) {
860            return; // stale token for a closed/reused session
861        }
862        match cqe.token.op() {
863            Op::Accept => {}
864            Op::Connect => match cqe.result {
865                Ok(_) => {
866                    self.do_upstream_connected(slot, generation);
867                    let mut io = self.io_for(slot, generation);
868                    h.on_upstream_connected(&mut io);
869                }
870                Err(e) => {
871                    let mut io = self.io_for(slot, generation);
872                    h.on_upstream_error(&mut io, e);
873                }
874            },
875            Op::DownstreamRead => {
876                {
877                    let Some(s) = self.slab.get_mut(slot) else {
878                        eprintln!("WRKDBG slot gone");
879                        return;
880                    };
881                    s.read_inflight = false;
882                    // Splice takeover: forward the in-flight bytes raw to
883                    // the upstream (the kernel pump handles everything
884                    // after) and do not re-arm.
885                    let res = match &cqe.result {
886                        Ok(n) => Ok(*n),
887                        Err(e) => Err(io::Error::from_raw_os_error(e.raw_os_error().unwrap_or(0))),
888                    };
889                    if s.splice {
890                        self.splice_forward_inflight(slot, generation, true, res);
891                        return;
892                    }
893                }
894                match cqe.result {
895                    Ok(n) if n > 0 => {
896                        let data = {
897                            let Some(s) = self.slab.get_mut(slot) else {
898                                return;
899                            };
900                            let Some(rs) = s.rslot.take() else { return };
901                            let v = self.pool.slot(rs)[..n as usize].to_vec();
902                            self.pool.release(rs);
903                            v
904                        };
905                        let mut io = self.io_for(slot, generation);
906                        h.on_downstream_data(&mut io, &data);
907                        self.arm_downstream_read(slot, generation);
908                    }
909                    Ok(0) => {
910                        let Some(s) = self.slab.get_mut(slot) else {
911                            return;
912                        };
913                        if let Some(rs) = s.rslot.take() {
914                            self.pool.release(rs);
915                        }
916                        s.downstream_eof = true;
917                        let mut io = self.io_for(slot, generation);
918                        h.on_downstream_eof(&mut io);
919                        self.maybe_finish(slot, generation);
920                    }
921                    Ok(_) => {}
922                    Err(e) => {
923                        let mut io = self.io_for(slot, generation);
924                        h.on_downstream_error(&mut io, e);
925                    }
926                }
927            }
928            Op::DownstreamWrite => match cqe.result {
929                Ok(n) if n > 0 => self.continue_downstream_write(slot, generation, h),
930                Ok(_) => {}
931                Err(e) => {
932                    let mut io = self.io_for(slot, generation);
933                    h.on_downstream_error(&mut io, e);
934                }
935            },
936            Op::UpstreamRead | Op::UpstreamWrite => {
937                // Stale completions for a discarded upstream fd (failover
938                // replaced it) must not touch the session: the epoch rides
939                // the token aux and must match.
940                let epoch = u8::try_from(cqe.token.aux()).unwrap_or(u8::MAX);
941                let epoch_ok = self
942                    .slab
943                    .get(slot)
944                    .is_some_and(|s| s.upstream_epoch == epoch);
945                if !epoch_ok {
946                    return;
947                }
948                match cqe.token.op() {
949                    Op::UpstreamRead => self.dispatch_upstream_read(cqe, slot, generation, h),
950                    _ => self.dispatch_upstream_write(cqe, slot, generation, h),
951                }
952            }
953            Op::Splice => match cqe.result {
954                Ok(0) => {
955                    let aux = cqe.token.aux();
956                    let Some(s) = self.slab.get_mut(slot) else {
957                        return;
958                    };
959                    if aux == 1 {
960                        if let Some(up) = &s.upstream {
961                            shutdown_write(up.fd());
962                        }
963                        s.downstream_eof = true;
964                    } else {
965                        shutdown_write(s.downstream.fd());
966                        s.upstream_eof = true;
967                    }
968                    self.maybe_finish(slot, generation);
969                }
970                Ok(_) => { /* bytes moved; multishot readiness continues */ }
971                Err(_) => self.close_session(slot, generation, "splice-err"),
972            },
973        }
974    }
975
976    /// Forwards the bytes held by an in-flight read slot to the opposite
977    /// side after splice takeover. `from_downstream` selects direction.
978    /// `result` is the read CQE: Ok(n) forwards, Ok(0) half-closes,
979    /// Err kills the session. The slot rides the write queue so it is
980    /// released on write completion (zero copy).
981    fn splice_forward_inflight(
982        &mut self,
983        slot: u32,
984        generation: u16,
985        from_downstream: bool,
986        result: io::Result<u32>,
987    ) {
988        let held = if from_downstream {
989            self.slab.get_mut(slot).and_then(|s| s.rslot.take())
990        } else {
991            self.slab.get_mut(slot).and_then(|s| s.urslot.take())
992        };
993        let Some(rs) = held else { return }; // no in-flight slot: pump owns it
994        let Some(s) = self.slab.get_mut(slot) else {
995            return;
996        };
997        match result {
998            Ok(n) if n > 0 => {
999                // Queue the held slot toward the opposite fd (zero copy).
1000                if from_downstream {
1001                    let Some(up) = &s.upstream else {
1002                        self.pool.release(rs);
1003                        return;
1004                    };
1005                    let fd = up.fd();
1006                    s.uwq.push_back((rs, n, 0));
1007                    s.upstream_write_inflight = true;
1008                    let token = Token::new(Op::UpstreamWrite, slot, generation, 0);
1009                    let _ = self.engine.write(token, fd, rs, n as usize, 0);
1010                } else {
1011                    let fd = s.downstream.fd();
1012                    s.wq.push_back((rs, n, 0));
1013                    s.write_inflight = true;
1014                    let token = Token::new(Op::DownstreamWrite, slot, generation, 0);
1015                    let _ = self.engine.write(token, fd, rs, n as usize, 0);
1016                }
1017            }
1018            // Only Ok(0) reaches here (n > 0 matched above): EOF —
1019            // half-close the opposite write side.
1020            Ok(_) => {
1021                if from_downstream {
1022                    if let Some(up) = &s.upstream {
1023                        shutdown_write(up.fd());
1024                    }
1025                    s.downstream_eof = true;
1026                } else {
1027                    shutdown_write(s.downstream.fd());
1028                    s.upstream_eof = true;
1029                }
1030                self.pool.release(rs);
1031                self.maybe_finish(slot, generation);
1032            }
1033            Err(_) => {
1034                self.pool.release(rs);
1035                self.close_session(slot, generation, "splice-forward-err");
1036            }
1037        }
1038    }
1039
1040    fn dispatch_upstream_read(
1041        &mut self,
1042        cqe: crate::engine::Cqe,
1043        slot: u32,
1044        generation: u16,
1045        h: &mut dyn Handler,
1046    ) {
1047        {
1048            let Some(s) = self.slab.get_mut(slot) else {
1049                return;
1050            };
1051            s.upstream_read_inflight = false;
1052            // Splice takeover: forward in-flight bytes raw to the client.
1053            let res = match &cqe.result {
1054                Ok(n) => Ok(*n),
1055                Err(e) => Err(io::Error::from_raw_os_error(e.raw_os_error().unwrap_or(0))),
1056            };
1057            if s.splice {
1058                self.splice_forward_inflight(slot, generation, false, res);
1059                return;
1060            }
1061        }
1062        match cqe.result {
1063            Ok(n) if n > 0 => {
1064                let data = {
1065                    let Some(s) = self.slab.get_mut(slot) else {
1066                        return;
1067                    };
1068                    let Some(rs) = s.urslot.take() else { return };
1069                    let v = self.pool.slot(rs)[..n as usize].to_vec();
1070                    self.pool.release(rs);
1071                    v
1072                };
1073                let mut io = self.io_for(slot, generation);
1074                h.on_upstream_data(&mut io, &data);
1075                self.arm_upstream_read(slot, generation);
1076            }
1077            Ok(0) => {
1078                let Some(s) = self.slab.get_mut(slot) else {
1079                    return;
1080                };
1081                if let Some(rs) = s.urslot.take() {
1082                    self.pool.release(rs);
1083                }
1084                s.upstream_eof = true;
1085                let mut io = self.io_for(slot, generation);
1086                h.on_upstream_eof(&mut io);
1087                self.maybe_finish(slot, generation);
1088            }
1089            Ok(_) => {}
1090            Err(e) => {
1091                let mut io = self.io_for(slot, generation);
1092                h.on_upstream_error(&mut io, e);
1093            }
1094        }
1095    }
1096
1097    fn dispatch_upstream_write(
1098        &mut self,
1099        cqe: crate::engine::Cqe,
1100        slot: u32,
1101        generation: u16,
1102        h: &mut dyn Handler,
1103    ) {
1104        match cqe.result {
1105            Ok(n) if n > 0 => self.continue_upstream_write(slot, generation, h),
1106            Ok(_) => {}
1107            Err(e) => {
1108                let mut io = self.io_for(slot, generation);
1109                h.on_upstream_error(&mut io, e);
1110            }
1111        }
1112    }
1113
1114    fn run_timers(&mut self, now: Instant, h: &mut dyn Handler) {
1115        while self.timers.peek().is_some() {
1116            if self.timers.peek().is_some_and(|t| t.at > now) {
1117                break;
1118            }
1119            #[allow(clippy::expect_used, reason = "peeked non-empty above")]
1120            let t = self.timers.pop().expect("peeked non-empty");
1121            if self.valid(t.slot, t.generation) {
1122                let Some(s) = self.slab.get(t.slot) else {
1123                    continue;
1124                };
1125                let Some(deadline) = s.deadline else { continue };
1126                if deadline > t.at {
1127                    continue; // superseded
1128                }
1129                let reason = crate::handler::DeadlineReason::from_aux(t.reason_aux);
1130                let mut io = self.io_for(t.slot, t.generation);
1131                h.on_deadline(&mut io, reason);
1132            }
1133        }
1134    }
1135
1136    fn next_timeout(&self) -> Option<Duration> {
1137        if let Some(t) = self.timers.peek() {
1138            return Some(t.at.saturating_duration_since(Instant::now()));
1139        }
1140        if self.draining {
1141            Some(Duration::from_millis(5))
1142        } else {
1143            Some(Duration::from_millis(50))
1144        }
1145    }
1146
1147    fn drain_finished(&mut self) -> bool {
1148        if !self.draining {
1149            return false;
1150        }
1151        if self.slab.live().is_empty() {
1152            return true;
1153        }
1154        if let Some(deadline) = self.drain_deadline {
1155            if Instant::now() >= deadline {
1156                for (slot, generation) in self.slab.live() {
1157                    self.close_session(slot, generation, "generic");
1158                }
1159                return true;
1160            }
1161        }
1162        false
1163    }
1164
1165    fn apply_cmd(&mut self, cmd: WorkerCmd, h: &mut dyn Handler) {
1166        match cmd {
1167            WorkerCmd::Shutdown { deadline_ms } => {
1168                self.draining = true;
1169                self.drain_deadline = Some(Instant::now() + Duration::from_millis(deadline_ms));
1170                for (slot, generation) in self.slab.live() {
1171                    let mut io = self.io_for(slot, generation);
1172                    h.on_shutdown_hint(&mut io);
1173                }
1174            }
1175        }
1176    }
1177
1178    fn accept_pending(&mut self, lidx: usize, h: &mut dyn Handler) {
1179        let Some(l) = self.listeners.get(lidx) else {
1180            return;
1181        };
1182        let lfd = l.as_raw_fd();
1183        let ltoken = Token::accept(lidx as u16);
1184        // Drain the current backlog; `Err`/`None` ends the round.
1185        #[allow(clippy::while_let_loop)]
1186        while let Ok(Some((fd, peer))) = self.engine.accept(lfd, ltoken) {
1187            self.accept_connection(fd, peer, h);
1188        }
1189    }
1190
1191    fn accept_connection(&mut self, fd: i32, peer: SocketAddr, h: &mut dyn Handler) {
1192        crate::dbg_trace!("ACCEPT fd={fd} peer={peer}");
1193        let _ = set_nodelay(fd);
1194        self.ctx.connections.inc(&self.ctx.registry);
1195        if self.draining || self.slab.live().len() >= self.ctx.config.max_sessions {
1196            // Overload / drain: refuse.
1197            // SAFETY: owned, unregistered descriptor.
1198            unsafe { libc::close(fd) };
1199            return;
1200        }
1201        struct FdGuard(i32);
1202        impl Drop for FdGuard {
1203            fn drop(&mut self) {
1204                // SAFETY: single close of an owned descriptor.
1205                unsafe { libc::close(self.0) };
1206            }
1207        }
1208        let guard = FdGuard(fd);
1209        let session = Session {
1210            downstream: StreamFd(-1),
1211            peer,
1212            upstream: None,
1213            rslot: None,
1214            wq: Default::default(),
1215            urslot: None,
1216            uwq: Default::default(),
1217            pending_down: Vec::new(),
1218            pending_up: Vec::new(),
1219            read_inflight: false,
1220            upstream_read_inflight: false,
1221            write_inflight: false,
1222            upstream_write_inflight: false,
1223            connect_inflight: false,
1224            upstream_epoch: 0,
1225            request_sent: false,
1226            downstream_eof: false,
1227            upstream_eof: false,
1228            fin_queued: false,
1229            splice: false,
1230            deadline: None,
1231        };
1232        let Ok((slot, generation)) = self.slab.insert(session) else {
1233            return; // guard closes fd
1234        };
1235        {
1236            let Some(s) = self.slab.get_mut(slot) else {
1237                return;
1238            };
1239            s.downstream = StreamFd(fd);
1240        }
1241        std::mem::forget(guard); // ownership moved into the session
1242        let token = Token::new(Op::DownstreamRead, slot, generation, 0);
1243        if self.engine.add_stream(fd, token).is_err() {
1244            self.close_session(slot, generation, "add-stream-failed");
1245            return;
1246        }
1247        {
1248            let mut io = self.io_for(slot, generation);
1249            h.on_connected(&mut io);
1250        }
1251        self.arm_downstream_read(slot, generation);
1252    }
1253}
1254
1255/// Worker thread entry — owns transport state and the handler, drives the
1256/// loop until drain completes.
1257struct WorkerInner {
1258    state: WorkerState,
1259    handler: Box<dyn Handler>,
1260}
1261
1262impl WorkerInner {
1263    /// Main loop: commands → accepts → engine CQEs → timers.
1264    pub fn run(mut self, cmd_rx: SpscReceiver<WorkerCmd, 64>) {
1265        for (i, l) in self.state.listeners.iter().enumerate() {
1266            let _ = self
1267                .state
1268                .engine
1269                .add_listener(l.as_raw_fd(), Token::accept(i as u16));
1270        }
1271        let mut cqes: Vec<crate::engine::Cqe> = Vec::with_capacity(512);
1272        loop {
1273            while let Some(cmd) = cmd_rx.recv() {
1274                let WorkerInner { state, handler } = &mut self;
1275                state.apply_cmd(cmd, handler.as_mut());
1276            }
1277            if self.state.drain_finished() {
1278                break;
1279            }
1280            for lidx in 0..self.state.listeners.len() {
1281                let WorkerInner { state, handler } = &mut self;
1282                state.accept_pending(lidx, handler.as_mut());
1283            }
1284            if self.state.drain_finished() {
1285                break;
1286            }
1287
1288            let timeout = self.state.next_timeout();
1289            match self.state.engine.poll(timeout, &mut cqes) {
1290                Ok(()) => {}
1291                Err(e) if e.kind() == io::ErrorKind::Interrupted => {}
1292                Err(_) => break, // engine failure: exit worker
1293            }
1294            for cqe in cqes.drain(..) {
1295                if cqe.token.op() == Op::Accept {
1296                    let aux = usize::from(cqe.token.aux());
1297                    let WorkerInner { state, handler } = &mut self;
1298                    state.accept_pending(aux, handler.as_mut());
1299                } else {
1300                    let WorkerInner { state, handler } = &mut self;
1301                    state.dispatch_cqe(cqe, handler.as_mut());
1302                }
1303            }
1304            // Inline-completed connects: deliver the deferred callback.
1305            if !self.state.inline_connected.is_empty() {
1306                let queued = std::mem::take(&mut self.state.inline_connected);
1307                for (slot, generation) in queued {
1308                    if !self.state.valid(slot, generation) {
1309                        continue;
1310                    }
1311                    let WorkerInner { state, handler } = &mut self;
1312                    state.do_upstream_connected(slot, generation);
1313                    let mut io = state.io_for(slot, generation);
1314                    handler.on_upstream_connected(&mut io);
1315                }
1316            }
1317            let now = Instant::now();
1318            let WorkerInner { state, handler } = &mut self;
1319            state.run_timers(now, handler.as_mut());
1320            if self.state.drain_finished() {
1321                break;
1322            }
1323        }
1324    }
1325}
1326
1327/// Spawns a worker thread (pinned when `config.core` is set) over a
1328/// pre-bound listener.
1329///
1330/// # Errors
1331/// Thread spawn or engine creation failure.
1332#[allow(
1333    clippy::expect_used,
1334    reason = "startup-only allocation failures abort the worker, never a live request"
1335)]
1336pub fn spawn(
1337    id: usize,
1338    config: WorkerConfig,
1339    listener: StdTcpListener,
1340    registry: Arc<Registry>,
1341    events: Arc<EventRing<LogEvent, { vane_observe::EVENT_RING_CAPACITY }>>,
1342    factory: &dyn HandlerFactory,
1343) -> io::Result<WorkerHandle> {
1344    let connections = registry.register(
1345        "vane_active_connections",
1346        vane_observe::metrics::MetricKind::Gauge,
1347    );
1348    let ctx = WorkerCtx {
1349        id,
1350        registry,
1351        events,
1352        config: config.clone(),
1353        connections,
1354    };
1355    let handler = factory.build(&ctx);
1356    let (cmd_tx, cmd_rx) = crate::spsc::channel::<WorkerCmd, 64>();
1357    let done = Arc::new(AtomicBool::new(false));
1358    let done2 = Arc::clone(&done);
1359    let mode = factory.mode();
1360    let core = config.core;
1361    let pool_slots = config.pool_slots;
1362    let ring_entries = config.ring_entries;
1363    let sqpoll = config.sqpoll;
1364    let force_mio = config.force_mio;
1365    let max_sessions = config.max_sessions;
1366
1367    let builder = std::thread::Builder::new().name(format!("vane-worker-{id}"));
1368    let handle = builder.spawn(move || {
1369        if let Some(core) = core {
1370            // Pin to the requested core (TH-01). Failure is non-fatal.
1371            let cores = core_affinity::get_core_ids().unwrap_or_default();
1372            if let Some(c) = cores.get(core) {
1373                let _ = core_affinity::set_for_current(*c);
1374            }
1375        }
1376        let pool = BufferPool::new(pool_slots, DEFAULT_BUF_SIZE).expect("fixed pool allocation");
1377        let engine: Box<dyn Engine> = if force_mio {
1378            Box::new(crate::engine::mio_engine::MioEngine::new(&pool).expect("mio engine init"))
1379        } else {
1380            match create_engine(ring_entries, Some(&pool), sqpoll) {
1381                Ok(e) => e,
1382                Err(_) => Box::new(
1383                    crate::engine::mio_engine::MioEngine::new(&pool).expect("mio fallback"),
1384                ),
1385            }
1386        };
1387        let slab = SessionSlab::new(max_sessions).expect("session slab");
1388        let inner = WorkerInner {
1389            state: WorkerState {
1390                ctx,
1391                engine,
1392                pool,
1393                slab,
1394                listeners: vec![listener],
1395                timers: BinaryHeap::new(),
1396                draining: false,
1397                drain_deadline: None,
1398                mode,
1399                inline_connected: Vec::new(),
1400            },
1401            handler,
1402        };
1403        inner.run(cmd_rx);
1404        done2.store(true, Ordering::Release);
1405    })?;
1406
1407    Ok(WorkerHandle {
1408        id,
1409        cmd: cmd_tx,
1410        done,
1411        join: Some(handle),
1412    })
1413}