numinnex commented on code in PR #4036:
URL: https://github.com/apache/iggy/pull/4036#discussion_r3932145254
##########
core/server/src/boot/threads.rs:
##########
@@ -649,6 +827,188 @@ mod tests {
);
}
+ #[test]
+ fn peer_exit_countdown_releases_the_waiter_once_every_peer_is_out() {
+ let countdown = Arc::new(PeerExitCountdown::new(2));
+ let guards: [PeerExitGuard; 2] =
+ std::array::from_fn(|_|
PeerExitGuard::new(Arc::clone(&countdown)));
+ assert_eq!(
+ countdown.wait(Duration::from_millis(10)),
+ Err(2),
+ "two live peers must hold the waiter past a short deadline"
+ );
+ let peers: Vec<_> = guards
+ .into_iter()
+ .map(|guard| {
+ thread::spawn(move || {
+ thread::sleep(Duration::from_millis(20));
+ drop(guard);
+ })
+ })
+ .collect();
+ assert_eq!(
+ countdown.wait(Duration::from_secs(30)),
+ Ok(()),
+ "the last guard drop must release the waiter"
+ );
+ for peer in peers {
+ peer.join()
+ .expect("peer thread dropped its guard without panicking");
+ }
+ }
+
+ #[test]
+ fn peer_exit_guard_counts_out_during_unwind() {
+ let countdown = Arc::new(PeerExitCountdown::new(1));
+ let guard = PeerExitGuard::new(Arc::clone(&countdown));
+ // `resume_unwind` skips the panic hook, so the unwind is silent.
+ let unwound = panic::catch_unwind(panic::AssertUnwindSafe(|| {
+ let _guard = guard;
+ panic::resume_unwind(Box::new("peer shard body panicked"));
+ }));
+ assert!(unwound.is_err());
+ assert_eq!(
+ countdown.wait(Duration::ZERO),
+ Ok(()),
+ "a guard dropped by a panic unwind must still count its peer out"
+ );
+ }
+
+ #[test]
+ fn peer_exit_wait_flips_the_flag_before_waiting() {
+ // Stands in for a peer parked on its bus token: it exits only once
+ // the shutdown flag is set, so a waiter that set the flag after
+ // waiting would sit out the whole budget.
+ let countdown = Arc::new(PeerExitCountdown::new(1));
+ let shutdown_flag = Arc::new(AtomicBool::new(false));
+ let peer = thread::spawn({
+ let guard = PeerExitGuard::new(Arc::clone(&countdown));
+ let shutdown_flag = Arc::clone(&shutdown_flag);
+ move || {
+ while !shutdown_flag.load(Ordering::Relaxed) {
+ thread::sleep(Duration::from_millis(1));
+ }
+ drop(guard);
+ }
+ });
+ let started = Instant::now();
+ drop(PeerExitWait::new(
+ countdown,
+ Arc::clone(&shutdown_flag),
+ Arc::new(ShutdownDeadline::new(Duration::from_secs(30))),
+ ));
+ assert!(
+ started.elapsed() < Duration::from_secs(5),
+ "the waiter must not sit out its budget on a peer that waits for
the flag"
+ );
+ assert!(shutdown_flag.load(Ordering::Relaxed));
+ peer.join()
+ .expect("peer thread dropped its guard without panicking");
+ }
+
+ #[test]
+ fn peer_exit_wait_gives_up_at_the_deadline() {
+ let countdown = Arc::new(PeerExitCountdown::new(1));
+ let _wedged_peer = PeerExitGuard::new(Arc::clone(&countdown));
+ let timeout = Duration::from_millis(50);
+ let started = Instant::now();
+ drop(PeerExitWait::new(
+ countdown,
+ Arc::new(AtomicBool::new(false)),
+ Arc::new(ShutdownDeadline::new(timeout)),
+ ));
+ let waited = started.elapsed();
+ assert!(
+ waited >= timeout && waited < Duration::from_secs(5),
+ "a wedged peer must be abandoned at the deadline, waited
{waited:?}"
+ );
+ }
+
+ #[test]
+ fn peer_wait_and_shard_join_share_one_shutdown_budget() {
+ // Regression: the two waits nest (shard 0 cannot start waiting for
+ // its peers until its own drain returned, already inside the join
+ // budget). With a budget each, a clean shutdown reported shard 0 as
+ // wedged and exited non-zero.
+ let deadline =
Arc::new(ShutdownDeadline::new(Duration::from_millis(200)));
+ let shutdown_flag = AtomicBool::new(true);
+ // Never finishes: stands in for the slow drain that arms and then
+ // spends the shared budget. The thread leaks into the test process.
+ let wedged_shard = thread::spawn(|| -> Result<(), ServerError> {
+ loop {
+ thread::sleep(Duration::from_secs(1));
+ }
Review Comment:
**nit** — this test spawns `thread::spawn(|| loop { thread::sleep(1s) })`
with no stop path, so the thread leaks for the life of the test binary (the
comment acknowledges it). Parked in `sleep`, so it costs stack rather than CPU.
Note `:1143`'s pre-existing `abandons_a_shard_still_running_past_the_deadline`
does the same thing, so this copied an existing pattern — worth an `AtomicBool`
stop flag in both or neither, not just the new one.
##########
core/server/src/session_manager.rs:
##########
@@ -26,13 +26,30 @@
//! transport connection and the consensus-level `(client_id, session)` pair.
use crate::cluster_meta::ClusterRoster;
+use ahash::AHashMap;
use message_bus::installer::conn_info::ClientTransportKind;
use shard::ConnectedClientInfo;
-use std::collections::HashMap;
use std::net::SocketAddr;
use std::rc::Rc;
use std::time::{Duration, Instant};
+/// What the request funnel resolves from one `connections` lookup per frame.
+///
+/// The bound consensus session, the acting user, and the transport peer
+/// address. `Default` (everything absent, no address) stands for a
+/// connection neither this map nor the bus knows.
+#[derive(Debug, Clone, Copy, Default)]
+pub struct ConnectionContext {
Review Comment:
**nit** — `ConnectionContext` is bare `pub` in `pub mod session_manager`
(`core/server/src/lib.rs:49`), in the same commit that tightened
`RequestClass`/`classify` to `pub(in crate::dispatch)`. Harmless (`server` is
`publish = false`) but inconsistent. `pub(crate)` fits.
(`RewriteStage` looks like the same issue but is not — its `pub` is forced
by `pub struct RewriteDeny { pub stage: RewriteStage }` in a `pub(crate)`
module, so narrowing it trips `private_interfaces`.)
##########
core/server/src/dispatch/mod.rs:
##########
@@ -299,8 +245,19 @@ fn enqueue_client_request<B, MJ, S, SB>(
return;
}
- let bus = shard.bus.clone();
+ let shard_handle = Rc::clone(shard_handle);
+ let sessions = Rc::clone(sessions);
+ let system_config = Arc::clone(system_config);
Review Comment:
**nit, pre-existing, follow-up material** — recording the measurement next
to the queue-cap item you deferred, since they belong in the same change.
A lockstep client pays one `compio::runtime::spawn` task allocation plus 4
`Rc::clone` and 1 `Arc::clone` (atomic RMW) per request, because
`drain_client_requests` returns whenever the queue empties (`:290-303`) so
`active.insert` wins again on the next frame and respawns. The boxed future
holds ~512 B of headers live across the consensus await in the
`ReplicatedMetadata` arm, so it is ~1 KB+ — an order of magnitude larger than
the `VecDeque` allocation the `:311` change was made to avoid.
Not introduced here and not worth acting on alone. But it is the reason
restoring the `:311` removal is the right trade: the allocation it reintroduces
is the smaller half of a pair this path already pays. Folding `active` into the
queue entry would recover both, and drop 2 of the 4 per-frame hash lookups.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]