use std::fs::{File, OpenOptions};
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc::{self, Receiver, RecvTimeoutError, SyncSender};
use std::sync::{Arc, OnceLock};
use std::thread;
use std::time::{Duration, Instant};
use crate::pool::{ConnectionPool, TEST_HARNESS_ENV};
const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(5 * 60);
const DRAIN_INTERVAL: Duration = Duration::from_secs(1);
const QUEUE_CAPACITY: usize = 1024;
const MAX_ERROR_BYTES: usize = 512;
const DRAIN_BATCH_CAP: usize = 256;
const WRITE_DELAY_MS_OVERRIDE_ENV: &str = "KHIVE_WRITER_TIMEOUT_SINK_WRITE_DELAY_MS";
const STARTUP_BARRIER_DIR_ENV: &str = "KHIVE_WRITER_TIMEOUT_SINK_STARTUP_BARRIER_DIR";
const SLOW_WRITE_THRESHOLD: Duration = Duration::from_secs(1);
const SLOW_WRITE_THRESHOLD_MS_OVERRIDE_ENV: &str = "KHIVE_SLOW_WRITE_THRESHOLD_MS";
pub(crate) fn slow_write_threshold() -> Option<Duration> {
match std::env::var(SLOW_WRITE_THRESHOLD_MS_OVERRIDE_ENV) {
Ok(v) => match v.parse::<u64>() {
Ok(0) => None,
Ok(ms) => Some(Duration::from_millis(ms)),
Err(_) => Some(SLOW_WRITE_THRESHOLD),
},
Err(_) => Some(SLOW_WRITE_THRESHOLD),
}
}
fn ndjson_file_name() -> String {
format!("writer_timeouts.{}.ndjson", std::process::id())
}
const ROTATE_COLLISION_RETRIES: u32 = 16;
fn rotate_if_regular_file(path: &Path, dir: &Path) -> Result<(), ()> {
let is_regular_file = std::fs::symlink_metadata(path)
.map(|meta| meta.is_file())
.unwrap_or(false);
if !is_regular_file {
return Ok(());
}
let mut millis = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis())
.unwrap_or(0);
for _ in 0..ROTATE_COLLISION_RETRIES {
let candidate = dir.join(format!(
"writer_timeouts.{}.r{millis}.ndjson",
std::process::id()
));
if candidate.exists() {
millis += 1;
continue;
}
return std::fs::rename(path, &candidate).map_err(|_| ());
}
Err(())
}
const FALLBACK_LOG_SUBDIR: &str = ".khive-logs";
const SINK_DIR_OVERRIDE_ENV: &str = "KHIVE_WRITER_TIMEOUT_SINK_DIR";
const HEARTBEAT_MS_OVERRIDE_ENV: &str = "KHIVE_WRITER_TIMEOUT_SINK_HEARTBEAT_MS";
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum Site {
PoolAdmission,
StandaloneGraph,
StandaloneEvent,
StandaloneText,
StandaloneSqlBridge,
DirectRouteSqlBridgeWriter,
DirectRouteAtomicUnit,
DirectRouteVecDeleteSubjects,
DirectRouteOrphanSweep,
DirectRouteFtsRenameNamespace,
DirectRouteFtsGeneralWrite,
DirectRouteVecGeneralWrite,
DirectRouteEntity,
DirectRouteNote,
DirectRouteGraphGeneralWrite,
DirectRouteEventGeneralWrite,
DirectRouteSparseGeneralWrite,
DirectRouteAgentGeneralWrite,
DirectRouteRuntimeMergeEntity,
DirectRouteRuntimeMergeNote,
DirectRouteRuntimeUpdateSymmetricEdge,
}
impl Site {
fn as_str(self) -> &'static str {
match self {
Site::PoolAdmission => "pool_admission",
Site::StandaloneGraph => "standalone:graph",
Site::StandaloneEvent => "standalone:event",
Site::StandaloneText => "standalone:text",
Site::StandaloneSqlBridge => "standalone:sql_bridge",
Site::DirectRouteSqlBridgeWriter => "direct_route:sql_bridge_writer",
Site::DirectRouteAtomicUnit => "direct_route:atomic_unit",
Site::DirectRouteVecDeleteSubjects => "direct_route:vec_delete_subjects",
Site::DirectRouteOrphanSweep => "direct_route:orphan_sweep",
Site::DirectRouteFtsRenameNamespace => "direct_route:fts_rename_namespace",
Site::DirectRouteFtsGeneralWrite => "direct_route:fts_general_write",
Site::DirectRouteVecGeneralWrite => "direct_route:vec_general_write",
Site::DirectRouteEntity => "direct_route:entity",
Site::DirectRouteNote => "direct_route:note",
Site::DirectRouteGraphGeneralWrite => "direct_route:graph_general_write",
Site::DirectRouteEventGeneralWrite => "direct_route:event_general_write",
Site::DirectRouteSparseGeneralWrite => "direct_route:sparse_general_write",
Site::DirectRouteAgentGeneralWrite => "direct_route:agent_general_write",
Site::DirectRouteRuntimeMergeEntity => "direct_route:runtime_merge_entity",
Site::DirectRouteRuntimeMergeNote => "direct_route:runtime_merge_note",
Site::DirectRouteRuntimeUpdateSymmetricEdge => {
"direct_route:runtime_update_symmetric_edge"
}
}
}
}
struct QueuedEvent {
ts_utc: String,
kind: &'static str,
db: String,
site: Option<&'static str>,
error: Option<String>,
timeout_ms: Option<u64>,
elapsed_ms: Option<u64>,
queue_depth: Option<u64>,
writer_stages: Option<WriterStageFields>,
}
#[derive(Clone, serde::Serialize)]
struct WriterStageFields {
queue_wait_micros: u64,
transaction_acquire_micros: u64,
body_micros: u64,
commit_micros: u64,
total_micros: u64,
}
struct SinkHandle {
sender: SyncSender<QueuedEvent>,
dropped: Arc<AtomicU64>,
}
static SINK: OnceLock<SinkHandle> = OnceLock::new();
#[derive(serde::Serialize)]
struct EventRecord<'a> {
ts_utc: String,
kind: &'a str,
db: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
site: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
error: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
timeout_ms: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
elapsed_ms: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
queue_depth: Option<u64>,
#[serde(flatten, skip_serializing_if = "Option::is_none")]
writer_stages: Option<WriterStageFields>,
#[serde(skip_serializing_if = "Option::is_none")]
pid: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
version: Option<&'a str>,
}
#[allow(clippy::too_many_arguments)]
fn build_line(
ts_utc: String,
kind: &str,
db: &str,
site: Option<&str>,
error: Option<&str>,
timeout_ms: Option<u64>,
elapsed_ms: Option<u64>,
queue_depth: Option<u64>,
writer_stages: Option<WriterStageFields>,
pid: Option<u32>,
version: Option<&str>,
) -> String {
let record = EventRecord {
ts_utc,
kind,
db,
site,
error,
timeout_ms,
elapsed_ms,
queue_depth,
writer_stages,
pid,
version,
};
let mut line = serde_json::to_string(&record).unwrap_or_else(|_| "{}".to_string());
line.push('\n');
line
}
#[allow(clippy::too_many_arguments)]
fn build_line_now(
kind: &str,
db: &str,
site: Option<&str>,
error: Option<&str>,
timeout_ms: Option<u64>,
pid: Option<u32>,
version: Option<&str>,
) -> String {
build_line(
now_rfc3339(),
kind,
db,
site,
error,
timeout_ms,
None,
None,
None,
pid,
version,
)
}
fn now_rfc3339() -> String {
chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Millis, true)
}
fn truncate_error(s: &str, max_bytes: usize) -> String {
if s.len() <= max_bytes {
return s.to_string();
}
let mut end = max_bytes;
while end > 0 && !s.is_char_boundary(end) {
end -= 1;
}
let mut truncated = s[..end].to_string();
truncated.push_str("...(truncated)");
truncated
}
fn enqueue(sender: &SyncSender<QueuedEvent>, dropped: &AtomicU64, event: QueuedEvent) {
if sender.try_send(event).is_err() {
dropped.fetch_add(1, Ordering::Relaxed);
}
}
struct AppendSink<W> {
target: Option<W>,
poisoned: bool,
write_delay: Duration,
}
impl<W: Write> AppendSink<W> {
fn new() -> Self {
Self {
target: None,
poisoned: false,
write_delay: Duration::ZERO,
}
}
fn set_write_delay(&mut self, delay: Duration) {
self.write_delay = delay;
}
fn is_open(&self) -> bool {
self.target.is_some()
}
fn is_poisoned(&self) -> bool {
self.poisoned
}
fn set_target(&mut self, target: W) {
if !self.poisoned {
self.target = Some(target);
}
}
fn write_line(&mut self, line: &str) -> bool {
if self.poisoned {
return false;
}
let Some(target) = self.target.as_mut() else {
return false;
};
if !self.write_delay.is_zero() {
thread::sleep(self.write_delay);
}
match target.write(line.as_bytes()) {
Ok(n) if n == line.len() => true,
Ok(_) => {
self.poisoned = true;
self.target = None;
false
}
Err(_) => {
self.target = None;
false
}
}
}
#[cfg(test)]
fn into_target(self) -> Option<W> {
self.target
}
}
impl AppendSink<File> {
fn ensure_open(&mut self, dir: &Path) {
if self.poisoned || self.target.is_some() {
return;
}
let path = dir.join(ndjson_file_name());
if rotate_if_regular_file(&path, dir).is_err() {
return;
}
let opened = std::fs::create_dir_all(dir)
.and_then(|_| OpenOptions::new().create(true).append(true).open(&path))
.ok();
if let Some(file) = opened {
self.set_target(file);
}
}
fn sync(&self) {
if let Some(file) = self.target.as_ref() {
let _ = file.sync_data();
}
}
}
const RETENTION_MAX_AGE: Duration = Duration::from_secs(14 * 24 * 60 * 60);
fn prune_expired_sink_files(dir: &Path) {
let Ok(entries) = std::fs::read_dir(dir) else {
return;
};
let now = std::time::SystemTime::now();
for entry in entries.flatten() {
let name = entry.file_name();
let name = name.to_string_lossy();
if !name.starts_with("writer_timeouts.") || !name.ends_with(".ndjson") {
continue;
}
let Ok(age) = entry
.metadata()
.and_then(|meta| meta.modified())
.and_then(|modified| {
now.duration_since(modified)
.map_err(|e| std::io::Error::other(e.to_string()))
})
else {
continue;
};
if age > RETENTION_MAX_AGE {
let _ = std::fs::remove_file(entry.path());
}
}
}
fn retry_open_if_healthy(sink: &mut AppendSink<File>, dir: &Path) {
if !sink.is_poisoned() {
sink.ensure_open(dir);
}
}
fn write_event(sink: &mut AppendSink<File>, dropped: &AtomicU64, event: QueuedEvent) {
if !sink.is_open() {
dropped.fetch_add(1, Ordering::Relaxed);
return;
}
let line = build_line(
event.ts_utc,
event.kind,
&event.db,
event.site,
event.error.as_deref(),
event.timeout_ms,
event.elapsed_ms,
event.queue_depth,
event.writer_stages,
None,
None,
);
if sink.write_line(&line) {
sink.sync();
} else {
dropped.fetch_add(1, Ordering::Relaxed);
}
}
fn maybe_emit_heartbeat(
sink: &mut AppendSink<File>,
dir: &Path,
db_identity: &str,
dropped: &AtomicU64,
last_heartbeat: &mut Instant,
last_summary_dropped: &mut u64,
heartbeat_interval: Duration,
) {
if last_heartbeat.elapsed() < heartbeat_interval {
return;
}
*last_heartbeat = Instant::now();
sink.ensure_open(dir);
if !sink.is_open() {
return;
}
let hb = build_line_now(
"heartbeat",
db_identity,
None,
None,
None,
Some(std::process::id()),
Some(env!("CARGO_PKG_VERSION")),
);
sink.write_line(&hb);
let current_dropped = dropped.load(Ordering::Relaxed);
if current_dropped > *last_summary_dropped {
let message = format!(
"sink drops since last summary: {}",
current_dropped - *last_summary_dropped
);
let summary = build_line_now(
"sink_error_summary",
"-",
None,
Some(&message),
None,
None,
None,
);
if sink.write_line(&summary) {
*last_summary_dropped = current_dropped;
}
}
}
fn writer_thread_loop(
go: Receiver<()>,
receiver: Receiver<QueuedEvent>,
dir: PathBuf,
db_identity: String,
dropped: Arc<AtomicU64>,
heartbeat_interval: Duration,
) {
if go.recv().is_err() {
return;
}
pause_at_startup_barrier_for_test();
prune_expired_sink_files(&dir);
let mut sink = AppendSink::<File>::new();
sink.set_write_delay(write_delay_from_env());
let mut last_heartbeat = Instant::now();
let mut last_summary_dropped = 0u64;
sink.ensure_open(&dir);
if sink.is_open() {
let startup = build_line_now(
"startup",
&db_identity,
None,
None,
None,
Some(std::process::id()),
Some(env!("CARGO_PKG_VERSION")),
);
sink.write_line(&startup);
}
loop {
match receiver.recv_timeout(DRAIN_INTERVAL) {
Ok(first) => {
retry_open_if_healthy(&mut sink, &dir);
write_event(&mut sink, &dropped, first);
maybe_emit_heartbeat(
&mut sink,
&dir,
&db_identity,
&dropped,
&mut last_heartbeat,
&mut last_summary_dropped,
heartbeat_interval,
);
let mut drained_this_wakeup = 1usize;
while drained_this_wakeup < DRAIN_BATCH_CAP {
match receiver.try_recv() {
Ok(event) => {
write_event(&mut sink, &dropped, event);
drained_this_wakeup += 1;
maybe_emit_heartbeat(
&mut sink,
&dir,
&db_identity,
&dropped,
&mut last_heartbeat,
&mut last_summary_dropped,
heartbeat_interval,
);
}
Err(_) => break,
}
}
}
Err(RecvTimeoutError::Timeout) => {
retry_open_if_healthy(&mut sink, &dir);
}
Err(RecvTimeoutError::Disconnected) => break,
}
maybe_emit_heartbeat(
&mut sink,
&dir,
&db_identity,
&dropped,
&mut last_heartbeat,
&mut last_summary_dropped,
heartbeat_interval,
);
}
}
fn resolve_log_dir(db_parent: Option<&Path>) -> Option<PathBuf> {
if let Some(dir) = std::env::var_os(SINK_DIR_OVERRIDE_ENV) {
return Some(PathBuf::from(dir));
}
let under_test_harness = std::env::var(TEST_HARNESS_ENV).as_deref() == Ok("1");
if !under_test_harness {
if let Some(home) = std::env::var_os("HOME") {
return Some(Path::new(&home).join(".khive").join("logs"));
}
}
db_parent.map(|p| p.join(FALLBACK_LOG_SUBDIR))
}
fn heartbeat_interval_from_env() -> Duration {
std::env::var(HEARTBEAT_MS_OVERRIDE_ENV)
.ok()
.and_then(|v| v.parse::<u64>().ok())
.filter(|&ms| ms > 0)
.map(Duration::from_millis)
.unwrap_or(HEARTBEAT_INTERVAL)
}
fn write_delay_from_env() -> Duration {
std::env::var(WRITE_DELAY_MS_OVERRIDE_ENV)
.ok()
.and_then(|v| v.parse::<u64>().ok())
.map(Duration::from_millis)
.unwrap_or(Duration::ZERO)
}
fn pause_at_startup_barrier_for_test() {
if std::env::var(TEST_HARNESS_ENV).as_deref() != Ok("1") {
return;
}
let Some(dir) = std::env::var_os(STARTUP_BARRIER_DIR_ENV).map(PathBuf::from) else {
return;
};
if std::fs::write(dir.join("reached"), b"reached").is_err() {
return;
}
while !dir.join("release").exists() {
thread::sleep(Duration::from_millis(1));
}
let _ = std::fs::write(dir.join("resumed"), b"resumed");
}
pub(crate) fn init(db_parent: Option<&Path>, db_identity: &str) {
let Some(db_parent) = db_parent else {
return;
};
if SINK.get().is_some() {
return;
}
let Some(dir) = resolve_log_dir(Some(db_parent)) else {
return;
};
let heartbeat_interval = heartbeat_interval_from_env();
let (sender, receiver) = mpsc::sync_channel::<QueuedEvent>(QUEUE_CAPACITY);
let (go_sender, go_receiver) = mpsc::sync_channel::<()>(1);
let dropped = Arc::new(AtomicU64::new(0));
let thread_dropped = Arc::clone(&dropped);
let db_identity = db_identity.to_string();
let spawn_result = thread::Builder::new()
.name("khive-writer-timeout-sink".to_string())
.spawn(move || {
writer_thread_loop(
go_receiver,
receiver,
dir,
db_identity,
thread_dropped,
heartbeat_interval,
)
});
if spawn_result.is_ok() {
let handle = SinkHandle { sender, dropped };
if SINK.set(handle).is_ok() {
let _ = go_sender.send(());
}
}
}
pub(crate) fn db_label(pool: &ConnectionPool) -> String {
pool.canonical_path()
.map(|p| p.display().to_string())
.unwrap_or_else(|| "memory".to_string())
}
pub(crate) fn emit_timeout(db: &str, site: Site, error: &str, timeout_ms: Option<u64>) {
let Some(handle) = SINK.get() else {
return;
};
let event = QueuedEvent {
ts_utc: now_rfc3339(),
kind: "timeout",
db: db.to_string(),
site: Some(site.as_str()),
error: Some(truncate_error(error, MAX_ERROR_BYTES)),
timeout_ms,
elapsed_ms: None,
queue_depth: None,
writer_stages: None,
};
enqueue(&handle.sender, &handle.dropped, event);
}
pub(crate) fn emit_queue_saturation(db: &str, timeout_ms: u64) {
let Some(handle) = SINK.get() else {
return;
};
let event = QueuedEvent {
ts_utc: now_rfc3339(),
kind: "queue_saturation",
db: db.to_string(),
site: None,
error: None,
timeout_ms: Some(timeout_ms),
elapsed_ms: None,
queue_depth: None,
writer_stages: None,
};
enqueue(&handle.sender, &handle.dropped, event);
}
pub(crate) fn emit_writer_task_retirement(db: &str, reason: &str) {
let Some(handle) = SINK.get() else {
return;
};
let event = QueuedEvent {
ts_utc: now_rfc3339(),
kind: "writer_task_retirement",
db: db.to_string(),
site: None,
error: Some(truncate_error(reason, MAX_ERROR_BYTES)),
timeout_ms: None,
elapsed_ms: None,
queue_depth: None,
writer_stages: None,
};
enqueue(&handle.sender, &handle.dropped, event);
}
#[cfg(test)]
thread_local! {
static DIRECT_ROUTE_CAPTURE: std::cell::RefCell<Option<Vec<(String, Site)>>> = const {
std::cell::RefCell::new(None)
};
}
#[cfg(test)]
pub(crate) fn capture_direct_routes<R>(f: impl FnOnce() -> R) -> (R, Vec<(String, Site)>) {
struct CaptureGuard;
impl Drop for CaptureGuard {
fn drop(&mut self) {
DIRECT_ROUTE_CAPTURE.with(|capture| *capture.borrow_mut() = None);
}
}
DIRECT_ROUTE_CAPTURE.with(|capture| {
assert!(capture.borrow().is_none(), "nested direct-route capture");
*capture.borrow_mut() = Some(Vec::new());
});
let _guard = CaptureGuard;
let result = f();
let events = DIRECT_ROUTE_CAPTURE.with(|capture| capture.borrow_mut().take().unwrap());
(result, events)
}
pub(crate) fn emit_direct_route_violation(db: &str, site: Site) {
#[cfg(test)]
DIRECT_ROUTE_CAPTURE.with(|capture| {
if let Some(events) = capture.borrow_mut().as_mut() {
events.push((db.to_owned(), site));
}
});
let Some(handle) = SINK.get() else {
return;
};
let event = QueuedEvent {
ts_utc: now_rfc3339(),
kind: "direct_route_violation",
db: db.to_string(),
site: Some(site.as_str()),
error: None,
timeout_ms: None,
elapsed_ms: None,
queue_depth: None,
writer_stages: None,
};
enqueue(&handle.sender, &handle.dropped, event);
}
pub(crate) fn emit_slow_write(db: &str, stages: &crate::writer_task::WriterStageObservation) {
let Some(handle) = SINK.get() else {
return;
};
let event = QueuedEvent {
ts_utc: now_rfc3339(),
kind: "slow_write",
db: db.to_string(),
site: None,
error: None,
timeout_ms: None,
elapsed_ms: Some(stages.total_micros / 1_000),
queue_depth: Some(stages.queue_depth_at_entry),
writer_stages: Some(WriterStageFields {
queue_wait_micros: stages.queue_wait_micros,
transaction_acquire_micros: stages.transaction_acquire_micros,
body_micros: stages.body_micros,
commit_micros: stages.commit_micros,
total_micros: stages.total_micros,
}),
};
enqueue(&handle.sender, &handle.dropped, event);
}
pub(crate) fn is_busy_or_locked(err: &rusqlite::Error) -> bool {
matches!(
err,
rusqlite::Error::SqliteFailure(inner, _)
if matches!(
inner.code,
rusqlite::ErrorCode::DatabaseBusy | rusqlite::ErrorCode::DatabaseLocked
)
)
}
pub(crate) fn maybe_emit_busy(db: &str, site: Site, err: &rusqlite::Error) {
if is_busy_or_locked(err) {
emit_timeout(db, site, &err.to_string(), None);
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::thread;
#[test]
fn store_direct_route_sites_have_stable_names() {
assert_eq!(Site::DirectRouteEntity.as_str(), "direct_route:entity");
assert_eq!(Site::DirectRouteNote.as_str(), "direct_route:note");
assert_eq!(
Site::DirectRouteGraphGeneralWrite.as_str(),
"direct_route:graph_general_write"
);
assert_eq!(
Site::DirectRouteEventGeneralWrite.as_str(),
"direct_route:event_general_write"
);
assert_eq!(
Site::DirectRouteSparseGeneralWrite.as_str(),
"direct_route:sparse_general_write"
);
assert_eq!(
Site::DirectRouteAgentGeneralWrite.as_str(),
"direct_route:agent_general_write"
);
assert_eq!(
Site::DirectRouteRuntimeMergeEntity.as_str(),
"direct_route:runtime_merge_entity"
);
assert_eq!(
Site::DirectRouteRuntimeMergeNote.as_str(),
"direct_route:runtime_merge_note"
);
assert_eq!(
Site::DirectRouteRuntimeUpdateSymmetricEdge.as_str(),
"direct_route:runtime_update_symmetric_edge"
);
}
#[test]
fn resolve_log_dir_prefers_explicit_override() {
let dir = tempfile::tempdir().unwrap();
std::env::set_var(SINK_DIR_OVERRIDE_ENV, dir.path());
let resolved = resolve_log_dir(None);
std::env::remove_var(SINK_DIR_OVERRIDE_ENV);
assert_eq!(resolved, Some(dir.path().to_path_buf()));
}
#[test]
fn resolve_log_dir_falls_back_to_db_parent_under_test_harness() {
std::env::remove_var(SINK_DIR_OVERRIDE_ENV);
assert_eq!(std::env::var(TEST_HARNESS_ENV).as_deref(), Ok("1"));
let parent = Path::new("/tmp/does-not-need-to-exist");
let resolved = resolve_log_dir(Some(parent));
assert_eq!(resolved, Some(parent.join(FALLBACK_LOG_SUBDIR)));
}
#[test]
fn resolve_log_dir_none_without_db_parent_under_test_harness() {
std::env::remove_var(SINK_DIR_OVERRIDE_ENV);
assert_eq!(resolve_log_dir(None), None);
}
#[test]
fn slow_write_threshold_env_override_and_zero_disable() {
std::env::remove_var(SLOW_WRITE_THRESHOLD_MS_OVERRIDE_ENV);
assert_eq!(slow_write_threshold(), Some(SLOW_WRITE_THRESHOLD));
std::env::set_var(SLOW_WRITE_THRESHOLD_MS_OVERRIDE_ENV, "250");
assert_eq!(slow_write_threshold(), Some(Duration::from_millis(250)));
std::env::set_var(SLOW_WRITE_THRESHOLD_MS_OVERRIDE_ENV, "0");
assert_eq!(slow_write_threshold(), None);
std::env::set_var(SLOW_WRITE_THRESHOLD_MS_OVERRIDE_ENV, "not-a-number");
assert_eq!(slow_write_threshold(), Some(SLOW_WRITE_THRESHOLD));
std::env::remove_var(SLOW_WRITE_THRESHOLD_MS_OVERRIDE_ENV);
}
#[test]
fn slow_write_line_carries_elapsed_and_depth_and_omits_inapplicable_fields() {
let stages = WriterStageFields {
queue_wait_micros: 11,
transaction_acquire_micros: 22,
body_micros: 1_111_000,
commit_micros: 44,
total_micros: 1_234_000,
};
let line = build_line(
now_rfc3339(),
"slow_write",
"test-db",
None,
None,
None,
Some(1234),
Some(7),
Some(stages),
None,
None,
);
let parsed: serde_json::Value = serde_json::from_str(line.trim()).unwrap();
assert_eq!(parsed["kind"], "slow_write");
assert_eq!(parsed["elapsed_ms"], 1234);
assert_eq!(parsed["queue_depth"], 7);
assert_eq!(parsed["queue_wait_micros"], 11);
assert_eq!(parsed["transaction_acquire_micros"], 22);
assert_eq!(parsed["body_micros"], 1_111_000);
assert_eq!(parsed["commit_micros"], 44);
assert_eq!(parsed["total_micros"], 1_234_000);
assert!(parsed.get("site").is_none());
assert!(parsed.get("error").is_none());
assert!(parsed.get("timeout_ms").is_none());
}
#[test]
fn truncate_error_leaves_short_strings_untouched() {
let short = "database is locked";
assert_eq!(truncate_error(short, MAX_ERROR_BYTES), short);
}
#[test]
fn truncate_error_bounds_long_strings() {
let long = "x".repeat(MAX_ERROR_BYTES * 4);
let truncated = truncate_error(&long, MAX_ERROR_BYTES);
assert!(
truncated.len() <= MAX_ERROR_BYTES + "...(truncated)".len(),
"truncated length {} exceeds bound",
truncated.len()
);
assert!(truncated.ends_with("...(truncated)"));
}
#[test]
fn truncate_error_respects_utf8_char_boundaries() {
let s = "a".repeat(MAX_ERROR_BYTES - 1) + "€€€€";
let truncated = truncate_error(&s, MAX_ERROR_BYTES);
assert!(truncated.is_char_boundary(truncated.len() - "...(truncated)".len()));
}
#[test]
fn enqueue_drops_and_counts_on_full_channel() {
let (sender, _receiver) = mpsc::sync_channel::<QueuedEvent>(1);
let dropped = AtomicU64::new(0);
let make_event = || QueuedEvent {
ts_utc: now_rfc3339(),
kind: "timeout",
db: "test-db".to_string(),
site: Some(Site::PoolAdmission.as_str()),
error: Some("boom".to_string()),
timeout_ms: Some(5),
elapsed_ms: None,
queue_depth: None,
writer_stages: None,
};
enqueue(&sender, &dropped, make_event());
assert_eq!(dropped.load(Ordering::Relaxed), 0);
enqueue(&sender, &dropped, make_event());
assert_eq!(dropped.load(Ordering::Relaxed), 1);
}
struct ShortWriteOnceThenOk {
calls: usize,
written: Vec<u8>,
}
impl Write for ShortWriteOnceThenOk {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.calls += 1;
if self.calls == 1 && buf.len() > 1 {
let n = buf.len() - 1;
self.written.extend_from_slice(&buf[..n]);
Ok(n)
} else {
self.written.extend_from_slice(buf);
Ok(buf.len())
}
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[test]
fn append_sink_short_write_poisons_permanently() {
let mut sink = AppendSink::<ShortWriteOnceThenOk>::new();
sink.set_target(ShortWriteOnceThenOk {
calls: 0,
written: Vec::new(),
});
assert!(
!sink.write_line("first line\n"),
"a short write must be reported as failed"
);
assert!(sink.is_poisoned());
assert!(!sink.is_open(), "the target must be dropped on poisoning");
sink.set_target(ShortWriteOnceThenOk {
calls: 0,
written: Vec::new(),
});
assert!(
!sink.is_open(),
"set_target must refuse to arm a poisoned sink"
);
assert!(
!sink.write_line("second line\n"),
"a poisoned sink must never write again"
);
}
struct ErrorOnceThenOk {
already_failed: Arc<std::sync::atomic::AtomicBool>,
written: Vec<u8>,
}
impl Write for ErrorOnceThenOk {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
if !self.already_failed.swap(true, Ordering::Relaxed) {
return Err(std::io::Error::other("injected failure"));
}
self.written.extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[test]
fn append_sink_hard_error_is_recoverable_not_poisoned() {
let already_failed = Arc::new(std::sync::atomic::AtomicBool::new(false));
let mut sink = AppendSink::<ErrorOnceThenOk>::new();
sink.set_target(ErrorOnceThenOk {
already_failed: Arc::clone(&already_failed),
written: Vec::new(),
});
assert!(!sink.write_line("first line\n"));
assert!(
!sink.is_poisoned(),
"an ordinary write error must not be terminal"
);
assert!(!sink.is_open(), "the failed target is dropped");
sink.set_target(ErrorOnceThenOk {
already_failed: Arc::clone(&already_failed),
written: Vec::new(),
});
assert!(sink.write_line("second line\n"));
}
#[test]
fn append_sink_recovers_after_transient_directory_failure() {
let dir = tempfile::tempdir().unwrap();
let blocker = dir.path().join("blocker");
std::fs::write(&blocker, b"not a directory").unwrap();
let bogus_dir = blocker.join("logs");
let mut sink = AppendSink::<File>::new();
sink.ensure_open(&bogus_dir);
assert!(!sink.is_open(), "open against a blocked path must fail");
assert!(!sink.is_poisoned(), "a failed open must not be terminal");
std::fs::remove_file(&blocker).unwrap();
sink.ensure_open(&bogus_dir);
assert!(
sink.is_open(),
"a later ensure_open against a now-writable directory must succeed"
);
assert!(sink.write_line("recovered\n"));
let contents = std::fs::read_to_string(bogus_dir.join(ndjson_file_name())).unwrap();
assert!(contents.contains("recovered"));
}
#[test]
fn concurrent_enqueue_then_single_threaded_drain_produces_intact_lines() {
let (sender, receiver) = mpsc::sync_channel::<QueuedEvent>(QUEUE_CAPACITY);
let dropped = Arc::new(AtomicU64::new(0));
const THREADS: usize = 64;
let handles: Vec<_> = (0..THREADS)
.map(|i| {
let sender = sender.clone();
let dropped = Arc::clone(&dropped);
thread::spawn(move || {
let event = QueuedEvent {
ts_utc: now_rfc3339(),
kind: "timeout",
db: "concurrent-test".to_string(),
site: Some(Site::StandaloneGraph.as_str()),
error: Some(format!("contention-{i}")),
timeout_ms: Some(i as u64),
elapsed_ms: None,
queue_depth: None,
writer_stages: None,
};
enqueue(&sender, &dropped, event);
})
})
.collect();
for h in handles {
h.join().unwrap();
}
drop(sender);
let mut sink = AppendSink::<Vec<u8>>::new();
sink.set_target(Vec::new());
let mut written_lines = 0usize;
while let Ok(event) = receiver.try_recv() {
let line = build_line(
event.ts_utc,
event.kind,
&event.db,
event.site,
event.error.as_deref(),
event.timeout_ms,
event.elapsed_ms,
event.queue_depth,
event.writer_stages,
None,
None,
);
if sink.write_line(&line) {
written_lines += 1;
}
}
assert_eq!(
written_lines + dropped.load(Ordering::Relaxed) as usize,
THREADS,
"every enqueue must be either written or counted as dropped"
);
assert_eq!(
dropped.load(Ordering::Relaxed),
0,
"queue capacity exceeds THREADS"
);
let bytes = sink.into_target().expect("target must still be armed");
let contents = String::from_utf8(bytes).unwrap();
let lines: Vec<&str> = contents.lines().collect();
assert_eq!(lines.len(), THREADS);
for line in &lines {
let parsed: serde_json::Value =
serde_json::from_str(line).unwrap_or_else(|e| panic!("corrupt line {line:?}: {e}"));
assert_eq!(parsed["kind"], "timeout");
assert_eq!(parsed["site"], "standalone:graph");
}
}
#[test]
fn is_busy_or_locked_classifies_sqlite_error_codes() {
let busy = rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
Some("database is locked".to_string()),
);
assert!(is_busy_or_locked(&busy));
let locked = rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_LOCKED),
Some("table is locked".to_string()),
);
assert!(is_busy_or_locked(&locked));
let other = rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
Some("constraint failed".to_string()),
);
assert!(!is_busy_or_locked(&other));
assert!(!is_busy_or_locked(&rusqlite::Error::QueryReturnedNoRows));
}
}