diff --git a/scripts/container-release-smoke.sh b/scripts/container-release-smoke.sh index 6eda9b62..92a3a0a6 100644 --- a/scripts/container-release-smoke.sh +++ b/scripts/container-release-smoke.sh @@ -25,6 +25,16 @@ cleanup() { } trap cleanup EXIT +emit_startup_diagnostics() { + echo "ERROR: release-smoke health check failed for ${container}" >&2 + docker inspect --format '{{json .State}}' "$container" 2>/dev/null >&2 || true + docker logs --tail 200 "$container" 2>&1 \ + | sed -E \ + -e 's/(Bearer )[A-Za-z0-9._~+\/-]+=*/\1[REDACTED]/g' \ + -e 's/((TOKEN|KEY|SECRET|PASSWORD)=)[^[:space:]]+/\1[REDACTED]/g' \ + >&2 || true +} + docker run -d --name "$container" \ -e CORTEX_API_TOKEN=release-smoke-api \ -e CORTEX_TOKEN=release-smoke-mcp \ @@ -44,11 +54,18 @@ udp_port=$(docker port "$container" 1514/udp | sed 's/.*://') # authority presented to RMCP's DNS-rebinding protection. http_host='127.0.0.1:3100' +healthy=false for _ in $(seq 1 60); do - curl -fsS -H "Host: ${http_host}" "http://127.0.0.1:${http_port}/health" >/dev/null && break + if curl -fsS -H "Host: ${http_host}" "http://127.0.0.1:${http_port}/health" >/dev/null; then + healthy=true + break + fi sleep 1 done -curl -fsS -H "Host: ${http_host}" "http://127.0.0.1:${http_port}/health" >/dev/null +if [[ "$healthy" != true ]]; then + emit_startup_diagnostics + exit 1 +fi printf '<13>Aug 29 00:00:00 release-smoke smokeapp: %stcp\n' "$marker" | nc -w 2 127.0.0.1 "$tcp_port" printf '<13>Aug 29 00:00:00 release-smoke smokeapp: %sudp\n' "$marker" | nc -u -w 2 127.0.0.1 "$udp_port" diff --git a/src/receiver.rs b/src/receiver.rs index 092ac990..1575eca6 100644 --- a/src/receiver.rs +++ b/src/receiver.rs @@ -4,6 +4,7 @@ use parking_lot::Mutex; use std::time::Duration; use anyhow::Result; +use tokio_util::sync::CancellationToken; use tracing::{error, info, warn}; use crate::config::{ReceiverConfig, StorageConfig}; @@ -64,6 +65,7 @@ async fn supervise_listener( name: &'static str, observability: Arc, set_state: fn(&RuntimeObservability, ListenerState), + shutdown: CancellationToken, make_listener: F, ) where F: Fn() -> Fut + Send + 'static, @@ -73,7 +75,22 @@ async fn supervise_listener( loop { set_state(&observability, ListenerState::Alive); let started = tokio::time::Instant::now(); - let outcome = tokio::spawn(make_listener()).await; + // A stopped supervisor must never detach its socket-owning child. + let mut listener = tokio_util::task::AbortOnDropHandle::new(tokio::spawn(make_listener())); + let outcome = tokio::select! { + biased; + _ = shutdown.cancelled() => { + // UDP/TCP listeners wait in recv/accept and do not own the + // runtime token. Abort their per-attempt task to close the + // socket before returning from this supervisor. + listener.abort(); + let _ = listener.await; + set_state(&observability, ListenerState::Down); + tracing::debug!(listener = name, "syslog listener supervisor stopped cleanly"); + return; + } + outcome = &mut listener => outcome, + }; set_state(&observability, ListenerState::Down); match outcome { Ok(Ok(())) => { @@ -99,7 +116,14 @@ async fn supervise_listener( backoff_secs = backoff.as_secs(), "restarting listener after backoff" ); - tokio::time::sleep(backoff).await; + tokio::select! { + biased; + _ = shutdown.cancelled() => { + tracing::debug!(listener = name, "syslog listener supervisor stopped during backoff"); + return; + } + _ = tokio::time::sleep(backoff) => {} + } backoff = (backoff * 2).min(LISTENER_BACKOFF_MAX); if backoff == LISTENER_BACKOFF_MAX { error!( @@ -127,6 +151,18 @@ pub(crate) async fn start_listeners( config: ReceiverConfig, ingest: ingest::IngestTx, observability: Arc, +) -> Result { + start_listeners_with_shutdown(config, ingest, observability, CancellationToken::new()).await +} + +/// Start supervised syslog listeners that terminate when `shutdown` is +/// cancelled. RuntimeCore uses its maintenance token here so graceful server +/// shutdown owns the otherwise unbounded listener loops. +pub(crate) async fn start_listeners_with_shutdown( + config: ReceiverConfig, + ingest: ingest::IngestTx, + observability: Arc, + shutdown: CancellationToken, ) -> Result { let bind_addr = config.bind_addr(); let allowed_cidrs = Arc::new(listener::parse_allowed_cidrs(&config.allowed_source_cidrs)?); @@ -139,6 +175,7 @@ pub(crate) async fn start_listeners( "udp_syslog", Arc::clone(&observability), |obs, state| obs.set_udp_listener_state(state), + shutdown.clone(), move || { let bind = udp_bind.clone(); let ingest = udp_ingest.clone(); @@ -156,6 +193,7 @@ pub(crate) async fn start_listeners( "tcp_syslog", Arc::clone(&observability), |obs, state| obs.set_tcp_listener_state(state), + shutdown, move || { let bind = tcp_bind.clone(); let ingest = tcp_ingest.clone(); diff --git a/src/receiver_tests.rs b/src/receiver_tests.rs index a938d737..0dd02f1c 100644 --- a/src/receiver_tests.rs +++ b/src/receiver_tests.rs @@ -3,6 +3,7 @@ use std::sync::atomic::{AtomicU32, Ordering}; use std::time::Duration; use parking_lot::Mutex; +use tokio_util::sync::CancellationToken; use super::*; @@ -18,6 +19,7 @@ async fn supervisor_restarts_listener_after_panic() { "test_listener", Arc::clone(&obs), |o, s| o.set_udp_listener_state(s), + CancellationToken::new(), move || { let attempts = Arc::clone(&attempts_in); async move { @@ -60,6 +62,7 @@ async fn supervisor_marks_listener_down_while_failing() { "test_listener", Arc::clone(&obs), |o, s| o.set_tcp_listener_state(s), + CancellationToken::new(), move || { let attempts = Arc::clone(&attempts_in); async move { @@ -115,6 +118,7 @@ async fn supervisor_resets_backoff_after_stable_run() { "test_listener", Arc::clone(&obs), |o, s| o.set_udp_listener_state(s), + CancellationToken::new(), move || { let attempts = Arc::clone(&attempts_in); let exits = Arc::clone(&exits_in); @@ -225,3 +229,124 @@ async fn start_listeners_wires_udp_and_tcp_supervisors_on_ephemeral_loopback_por handles.tcp.abort(); ingest.shutdown(Duration::from_secs(1)).await; } + +#[tokio::test] +async fn listener_supervisors_stop_cleanly_when_runtime_shutdown_is_cancelled() { + let dir = tempfile::tempdir().unwrap(); + let storage = StorageConfig::for_test(dir.path().join("receiver-shutdown.db")); + let pool = Arc::new(db::init_pool(&storage).unwrap()); + let storage_state = Arc::new(Mutex::new(None)); + let observability = Arc::new(RuntimeObservability::default()); + let config = ReceiverConfig { + host: "127.0.0.1".to_string(), + port: 0, + max_message_size: 1024, + max_tcp_connections: 4, + tcp_idle_timeout_secs: 1, + batch_size: 10, + flush_interval: 10, + write_channel_capacity: 16, + allowed_source_cidrs: Vec::new(), + }; + let ingest = ingest::start_writer_from_receiver_config( + &config, + storage, + pool, + storage_state, + crate::receiver::enrichment::EnrichmentConfig::default(), + Arc::clone(&observability), + ); + let shutdown = CancellationToken::new(); + let handles = start_listeners_with_shutdown( + config, + ingest.clone(), + Arc::clone(&observability), + shutdown.clone(), + ) + .await + .expect("listeners start"); + + tokio::time::timeout(Duration::from_secs(1), async { + while observability.udp_listener_state() != ListenerState::Alive + || observability.tcp_listener_state() != ListenerState::Alive + { + tokio::task::yield_now().await; + } + }) + .await + .expect("listener supervisors report alive"); + + shutdown.cancel(); + tokio::time::timeout(Duration::from_millis(250), async { + handles.udp.await.expect("UDP supervisor stops cleanly"); + handles.tcp.await.expect("TCP supervisor stops cleanly"); + }) + .await + .expect("cancellation must not wait for listener receive/accept"); + assert_eq!(observability.udp_listener_state(), ListenerState::Down); + assert_eq!(observability.tcp_listener_state(), ListenerState::Down); + ingest.shutdown(Duration::from_secs(1)).await; +} + +async fn assert_supervisor_releases_bound_socket(abort_supervisor: bool) { + let observability = Arc::new(RuntimeObservability::default()); + let shutdown = CancellationToken::new(); + let (bound_tx, bound_rx) = tokio::sync::oneshot::channel(); + let bound_tx = Arc::new(Mutex::new(Some(bound_tx))); + let supervisor = tokio::spawn(supervise_listener( + "socket_shutdown_test", + observability, + |obs, state| obs.set_udp_listener_state(state), + shutdown.clone(), + move || { + let bound_tx = Arc::clone(&bound_tx); + async move { + let socket = tokio::net::UdpSocket::bind("127.0.0.1:0").await?; + bound_tx + .lock() + .take() + .unwrap() + .send(socket.local_addr()?) + .unwrap(); + let mut buffer = [0u8; 1]; + socket.recv_from(&mut buffer).await?; + Ok(()) + } + }, + )); + let address = tokio::time::timeout(Duration::from_secs(2), bound_rx) + .await + .unwrap() + .unwrap(); + assert!(tokio::net::UdpSocket::bind(address).await.is_err()); + if abort_supervisor { + supervisor.abort(); + assert!(supervisor.await.unwrap_err().is_cancelled()); + } else { + shutdown.cancel(); + tokio::time::timeout(Duration::from_secs(2), supervisor) + .await + .unwrap() + .unwrap(); + } + tokio::time::timeout(Duration::from_secs(2), async { + loop { + if tokio::net::UdpSocket::bind(address).await.is_ok() { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("stopping the supervisor must release the actual listener socket"); +} + +#[tokio::test] +async fn supervisor_cancellation_releases_bound_socket() { + assert_supervisor_releases_bound_socket(false).await; +} + +#[tokio::test] +async fn supervisor_abort_releases_bound_socket() { + assert_supervisor_releases_bound_socket(true).await; +} diff --git a/src/runtime.rs b/src/runtime.rs index d989f01a..418e6804 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -608,36 +608,37 @@ impl RuntimeCore { /// [`MaintenanceHandles::syslog_monitor`] so it participates in the /// cooperative shutdown drain. pub async fn start_syslog(&self, handles: &mut MaintenanceHandles) -> Result<()> { - let listener_handles = receiver::start_listeners( + let listener_handles = receiver::start_listeners_with_shutdown( self.config.receiver.clone(), self.ingest.clone(), Arc::clone(&self.observability), + handles.token.clone(), ) .await?; let fatal_shutdown = self.fatal_shutdown.clone(); - let maintenance_shutdown = handles.token.clone(); + let shutdown = handles.token.clone(); let monitor = tokio::spawn(async move { let mut udp = listener_handles.udp; let mut tcp = listener_handles.tcp; let protocol = tokio::select! { - _ = maintenance_shutdown.cancelled() => { - // The listener supervisors deliberately run forever while - // serving. They are owned by this monitor, so a normal - // process shutdown must stop and join them instead of - // waiting for the monitor's "unexpected exit" branch. - // Without this branch the monitor alone consumed the - // entire maintenance shutdown budget and made every clean - // container stop look like an unclean runtime shutdown. - udp.abort(); - tcp.abort(); - let _ = tokio::join!(udp, tcp); - tracing::debug!("syslog listeners stopped for maintenance shutdown"); + biased; + _ = shutdown.cancelled() => { + // Both supervisors receive this token and abort their + // active recv/accept task before returning. Join them so + // graceful shutdown does not leave detached listeners or + // misclassify their expected exit as a fatal outage. + let (udp_result, tcp_result) = tokio::join!(udp, tcp); + for (listener, result) in [("udp", udp_result), ("tcp", tcp_result)] { + if let Err(error) = result { + tracing::warn!(listener, error = %error, + "syslog listener supervisor failed during shutdown"); + } + } + tracing::debug!("syslog listener monitor stopped cleanly"); return; } res = &mut udp => { - tcp.abort(); - let _ = tcp.await; match res { Ok(()) => tracing::error!( "syslog supervisor task (udp) exited unexpectedly — \ @@ -649,11 +650,11 @@ impl RuntimeCore { listener will not restart: {}", e ), } + tcp.abort(); + let _ = tcp.await; "udp" } res = &mut tcp => { - udp.abort(); - let _ = udp.await; match res { Ok(()) => tracing::error!( "syslog supervisor task (tcp) exited unexpectedly — \ @@ -665,6 +666,8 @@ impl RuntimeCore { listener will not restart: {}", e ), } + udp.abort(); + let _ = udp.await; "tcp" } }; diff --git a/src/runtime_tests.rs b/src/runtime_tests.rs index 2de32551..ff3cd8e5 100644 --- a/src/runtime_tests.rs +++ b/src/runtime_tests.rs @@ -467,6 +467,48 @@ async fn spawn_maintenance_tasks_constructs_expected_handles_and_shutdowns_clean runtime.shutdown(std::time::Duration::from_secs(1)).await; } +#[tokio::test] +async fn maintenance_shutdown_stops_syslog_monitor_and_listeners_cooperatively() { + let tmp = tempfile::tempdir().unwrap(); + let mut config = test_config(tmp.path(), loopback_mcp()); + // TCP and UDP can each bind their own ephemeral port; the exact port is + // irrelevant to this lifecycle assertion and avoids test-host collisions. + config.receiver.host = "127.0.0.1".into(); + config.receiver.port = 0; + let runtime = RuntimeCore::for_server(config).await.expect("runtime"); + let mut handles = runtime.spawn_maintenance_tasks(); + runtime + .start_syslog(&mut handles) + .await + .expect("syslog supervisors start"); + + tokio::time::timeout(Duration::from_secs(1), async { + while runtime.observability.udp_listener_state() + != crate::observability::ListenerState::Alive + || runtime.observability.tcp_listener_state() + != crate::observability::ListenerState::Alive + { + tokio::task::yield_now().await; + } + }) + .await + .expect("both listener supervisors report alive"); + + assert!( + handles.shutdown(Duration::from_millis(250)).await, + "graceful shutdown must not spend its timeout aborting syslog tasks" + ); + assert_eq!( + runtime.observability.udp_listener_state(), + crate::observability::ListenerState::Down + ); + assert_eq!( + runtime.observability.tcp_listener_state(), + crate::observability::ListenerState::Down + ); + runtime.shutdown(Duration::from_secs(1)).await; +} + #[tokio::test] async fn shutdown_timeout_aborts_and_joins_non_cooperative_tasks() { let progress = Arc::new(std::sync::atomic::AtomicUsize::new(0));