A cluster rehearsal moved a Log tracer to another directory: the write was reported x:settingsReload applied:true, but the tracer kept writing to the old file until a restart. Telemetry::update only refreshed each running tracer's events, level and lossiness; a tracer's own settings (path, prefix, rotation, format, endpoint, headers, ...) stayed as built. Each tracer now carries a hash of the registry object it was built from, less the fields that change in place. The reload compares it with the running tracer's: unchanged ones are updated in place as before, changed ones are started over, new ones started and removed ones stopped. Only tracers this server started are removed; upstream removed every subscriber not in the settings, which also cut off live-tracing streams on each reload. Starting over is a swap in the collector, so no event is lost or written twice: a subscriber registered under a running one's id replaces it between two collection passes. The old one's batch is sent first (what its full channel can't take moves to the new one), and dropping it closes its channel, so its task writes what is queued and ends. Per tracer kind: - Log: a tracer started over on the same files (rotation or format changed) waits for the old one to finish, so lines don't interleave. - Webhook: the task held a sender of its own channel for retries, so it never ended; retries now use a weak sender, and pending events are posted when the channel closes. - OpenTelemetry: pending logs and spans are exported when the channel closes instead of dropped, and a span that was open across the swap is exported by the new tracer with the events it saw. - Console and journal: nothing kept between batches. - Trace history: built from the tracing store, which takes a restart, so it is never started over. No kind needs a restart, so x:settingsReload doesn't gain one. system::tracer_reload::tracer_reload_tests (new): a Log tracer created over JMAP writes to its directory; its path is changed over JMAP while 2000 numbered events are emitted; after the reload, events land in the new file and not the old one, each numbered event is in exactly one of the two files, and a destroyed tracer writes nothing. On main the new file never appears.
191 lines
6.3 KiB
Rust
191 lines
6.3 KiB
Rust
/*
|
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
|
*
|
|
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
|
*
|
|
* Modified by Coffey Labs in 2026 for INBUXA.
|
|
*/
|
|
|
|
use crate::{LONG_1Y_SLUMBER, config::telemetry::WebhookTracer};
|
|
use aws_lc_rs::hmac;
|
|
use base64::{Engine, engine::general_purpose::STANDARD};
|
|
use serde::Serialize;
|
|
use std::{
|
|
sync::{
|
|
Arc,
|
|
atomic::{AtomicBool, Ordering},
|
|
},
|
|
time::Instant,
|
|
};
|
|
use store::write::now;
|
|
use tokio::sync::mpsc;
|
|
use trc::{
|
|
Event, EventDetails, ServerEvent, TelemetryEvent,
|
|
ipc::subscriber::{EventBatch, SubscriberBuilder},
|
|
serializers::json::JsonEventSerializer,
|
|
};
|
|
|
|
pub(crate) fn spawn_webhook_tracer(builder: SubscriberBuilder, settings: WebhookTracer) {
|
|
let (tx, mut rx) = builder.register();
|
|
// inbuxa: failed deliveries come back through a weak sender, so the
|
|
// channel closes when the collector drops this webhook (removed, or
|
|
// replaced after a settings change) and the task ends; upstream held a
|
|
// sender here and the task outlived its subscription
|
|
let tx = tx.downgrade();
|
|
tokio::spawn(async move {
|
|
let settings = Arc::new(settings);
|
|
let mut wakeup_time = LONG_1Y_SLUMBER;
|
|
let discard_after = settings.discard_after.as_secs();
|
|
let mut pending_events = Vec::new();
|
|
let mut next_delivery = Instant::now();
|
|
let in_flight = Arc::new(AtomicBool::new(false));
|
|
|
|
loop {
|
|
// Wait for the next event or timeout
|
|
let event_or_timeout = tokio::time::timeout(wakeup_time, rx.recv()).await;
|
|
let now = now();
|
|
|
|
match event_or_timeout {
|
|
Ok(Some(events)) => {
|
|
let mut discard_count = 0;
|
|
for event in events {
|
|
if now.saturating_sub(event.inner.timestamp) < discard_after {
|
|
pending_events.push(event)
|
|
} else {
|
|
discard_count += 1;
|
|
}
|
|
}
|
|
|
|
if discard_count > 0 {
|
|
trc::event!(
|
|
Telemetry(TelemetryEvent::WebhookError),
|
|
Details = "Discarded stale events",
|
|
Total = discard_count
|
|
);
|
|
}
|
|
}
|
|
Ok(None) => {
|
|
// inbuxa: deliver what is pending rather than drop it
|
|
if !pending_events.is_empty() {
|
|
spawn_webhook_handler(
|
|
settings.clone(),
|
|
in_flight.clone(),
|
|
std::mem::take(&mut pending_events),
|
|
tx.clone(),
|
|
);
|
|
}
|
|
break;
|
|
}
|
|
Err(_) => (),
|
|
}
|
|
|
|
// Process events
|
|
let mut next_retry = None;
|
|
let now = Instant::now();
|
|
if next_delivery <= now {
|
|
if !pending_events.is_empty() {
|
|
next_delivery = now + settings.throttle;
|
|
if !in_flight.load(Ordering::Relaxed) {
|
|
spawn_webhook_handler(
|
|
settings.clone(),
|
|
in_flight.clone(),
|
|
std::mem::take(&mut pending_events),
|
|
tx.clone(),
|
|
);
|
|
}
|
|
}
|
|
} else if !pending_events.is_empty() {
|
|
// Retry later
|
|
let this_retry = next_delivery - now;
|
|
match next_retry {
|
|
Some(next_retry) if this_retry >= next_retry => {}
|
|
_ => {
|
|
next_retry = Some(this_retry);
|
|
}
|
|
}
|
|
}
|
|
wakeup_time = next_retry.unwrap_or(LONG_1Y_SLUMBER);
|
|
}
|
|
});
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
struct EventWrapper {
|
|
events: JsonEventSerializer<Vec<Arc<Event<EventDetails>>>>,
|
|
}
|
|
|
|
fn spawn_webhook_handler(
|
|
settings: Arc<WebhookTracer>,
|
|
in_flight: Arc<AtomicBool>,
|
|
events: EventBatch,
|
|
webhook_tx: mpsc::WeakSender<EventBatch>,
|
|
) {
|
|
tokio::spawn(async move {
|
|
in_flight.store(true, Ordering::Relaxed);
|
|
let wrapper = EventWrapper {
|
|
events: JsonEventSerializer::new(events).with_id().with_spans(),
|
|
};
|
|
|
|
if let Err(err) = post_webhook_events(&settings, &wrapper).await {
|
|
trc::event!(Telemetry(TelemetryEvent::WebhookError), Details = err);
|
|
|
|
let sent = match webhook_tx.upgrade() {
|
|
Some(webhook_tx) => webhook_tx.send(wrapper.events.into_inner()).await.is_ok(),
|
|
None => false,
|
|
};
|
|
if !sent {
|
|
trc::event!(
|
|
Server(ServerEvent::ThreadError),
|
|
Details = "Failed to send failed webhook events back to main thread",
|
|
CausedBy = trc::location!()
|
|
);
|
|
}
|
|
}
|
|
|
|
in_flight.store(false, Ordering::Relaxed);
|
|
});
|
|
}
|
|
|
|
async fn post_webhook_events(
|
|
settings: &WebhookTracer,
|
|
events: &EventWrapper,
|
|
) -> Result<(), String> {
|
|
// Serialize body
|
|
let body = serde_json::to_string(events)
|
|
.map_err(|err| format!("Failed to serialize events: {}", err))?;
|
|
|
|
// Add HMAC-SHA256 signature
|
|
let mut headers = settings.headers.clone();
|
|
if !settings.key.is_empty() {
|
|
let key = hmac::Key::new(hmac::HMAC_SHA256, settings.key.as_bytes());
|
|
let tag = hmac::sign(&key, body.as_bytes());
|
|
|
|
headers.insert(
|
|
"X-Signature",
|
|
STANDARD.encode(tag.as_ref()).parse().unwrap(),
|
|
);
|
|
}
|
|
|
|
// Send request
|
|
let response = settings
|
|
.client
|
|
.post(&settings.url)
|
|
.timeout(settings.timeout)
|
|
.headers(headers)
|
|
.body(body)
|
|
.send()
|
|
.await
|
|
.map_err(|err| format!("Webhook request to {} failed: {err}", settings.url))?;
|
|
|
|
if response.status().is_success() {
|
|
Ok(())
|
|
} else {
|
|
Err(format!(
|
|
"Webhook request to {} failed with code {}: {}",
|
|
settings.url,
|
|
response.status().as_u16(),
|
|
response.status().canonical_reason().unwrap_or("Unknown")
|
|
))
|
|
}
|
|
}
|