1use 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
26const WRITE_PENDING_CAP: usize = 1024 * 1024;
30use crate::slab::SessionSlab;
31use crate::spsc::{SpscReceiver, SpscSender};
32use crate::token::{Op, Token};
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq)]
36pub enum WorkerCmd {
37 Shutdown {
39 deadline_ms: u64,
41 },
42}
43
44#[derive(Debug, Clone)]
46pub struct WorkerConfig {
47 pub core: Option<usize>,
49 pub max_sessions: usize,
51 pub pool_slots: usize,
53 pub ring_entries: u32,
55 pub sqpoll: bool,
57 pub force_mio: bool,
59 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
77pub struct WorkerCtx {
79 pub id: usize,
81 pub registry: Arc<Registry>,
83 pub events: Arc<EventRing<LogEvent, { vane_observe::EVENT_RING_CAPACITY }>>,
85 pub config: WorkerConfig,
87 pub connections: vane_observe::metrics::MetricHandle,
89}
90
91pub struct WorkerHandle {
93 pub id: usize,
95 pub cmd: SpscSender<WorkerCmd, 64>,
97 done: Arc<AtomicBool>,
99 join: Option<std::thread::JoinHandle<()>>,
100}
101
102impl WorkerHandle {
103 #[must_use]
105 pub fn is_done(&self) -> bool {
106 self.done.load(Ordering::Acquire)
107 }
108
109 pub fn join(&mut self) {
111 if let Some(h) = self.join.take() {
112 let _ = h.join();
113 }
114 }
115}
116
117struct Session {
119 downstream: StreamFd,
120 peer: SocketAddr,
121 upstream: Option<StreamFd>,
122 rslot: Option<u32>,
124 wq: std::collections::VecDeque<(u32, u32, u32)>,
126 urslot: Option<u32>,
128 uwq: std::collections::VecDeque<(u32, u32, u32)>,
130 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 #[allow(dead_code)] connect_inflight: bool,
140 upstream_epoch: u8,
143 request_sent: bool,
145 downstream_eof: bool,
146 upstream_eof: bool,
147 fin_queued: bool,
149 splice: bool,
150 deadline: Option<Instant>,
151}
152
153#[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) }
166}
167impl PartialOrd for Timer {
168 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
169 Some(self.cmp(other))
170 }
171}
172
173pub(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)] mode: Mode,
186 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 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 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 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 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 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; }
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 pub(crate) fn close_session(&mut self, slot: u32, generation: u16, reason: &str) {
409 let _ = reason;
411 if !self.valid(slot, generation) {
412 return;
413 }
414 let Some(s) = self.slab.remove(slot) else {
415 return;
416 }; 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 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 }
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 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 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 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 std::mem::forget(owned);
524
525 fd
526 };
527 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 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 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 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 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 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 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 }; 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 if s.upstream.is_none() {
668 return;
669 }
670 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 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 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 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; };
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 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 && (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; }
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 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 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(_) => { }
971 Err(_) => self.close_session(slot, generation, "splice-err"),
972 },
973 }
974 }
975
976 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 }; let Some(s) = self.slab.get_mut(slot) else {
995 return;
996 };
997 match result {
998 Ok(n) if n > 0 => {
999 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 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 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; }
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 #[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 unsafe { libc::close(fd) };
1199 return;
1200 }
1201 struct FdGuard(i32);
1202 impl Drop for FdGuard {
1203 fn drop(&mut self) {
1204 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; };
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); 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
1255struct WorkerInner {
1258 state: WorkerState,
1259 handler: Box<dyn Handler>,
1260}
1261
1262impl WorkerInner {
1263 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, }
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 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#[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 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}