use std::future::Future;
use std::sync::{Arc, Mutex};
use crate::util::UnwrapPoison;
const ENGINE_TARGET: &str = "turso_core::vdbe::execute";
const CAUSE_PREFIX: &str = "PRAGMA wal_checkpoint failed";
pub(crate) const BLOCKED_REASON: &str = "PRAGMA wal_checkpoint failed: Busy";
tokio::task_local! {
static ARMED: Arc<Mutex<Option<String>>>;
}
pub(crate) fn is_blocked_checkpoint(e: &anyhow::Error) -> bool {
format!("{e:#}")
.split_once(BLOCKED_REASON)
.is_some_and(|(_, rest)| !rest.starts_with(|c: char| c.is_alphanumeric() || c == '_'))
}
pub(crate) fn filter() -> tracing_subscriber::filter::Targets {
tracing_subscriber::filter::Targets::new()
.with_target(ENGINE_TARGET, tracing::level_filters::LevelFilter::DEBUG)
}
#[derive(Clone)]
pub(crate) struct CauseSink(Arc<Mutex<Option<String>>>);
impl CauseSink {
pub(crate) fn new() -> Self {
Self(Arc::new(Mutex::new(None)))
}
pub(crate) fn scoped<T, F>(&self, fut: F) -> impl Future<Output = T> + use<T, F>
where
F: Future<Output = T>,
{
ARMED.scope(Arc::clone(&self.0), fut)
}
pub(crate) fn attach(&self, e: anyhow::Error) -> anyhow::Error {
match self.0.lock().unwrap_poison().take() {
Some(cause) => e.context(format!("engine cause: {cause}")),
None => e,
}
}
}
pub(crate) struct CauseCaptureLayer;
impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for CauseCaptureLayer {
fn on_event(
&self,
event: &tracing::Event<'_>,
_ctx: tracing_subscriber::layer::Context<'_, S>,
) {
let capturing = ARMED
.try_with(|sink| sink.lock().unwrap_poison().is_none())
.unwrap_or(false);
if !capturing {
return;
}
let mut visitor = MessageVisitor(None);
event.record(&mut visitor);
let Some(message) = visitor.0 else {
return;
};
if message.starts_with(CAUSE_PREFIX) {
let _ = ARMED.try_with(|sink| *sink.lock().unwrap_poison() = Some(message));
}
}
}
struct MessageVisitor(Option<String>);
impl tracing::field::Visit for MessageVisitor {
fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
if field.name() == "message" {
self.0 = Some(value.to_string());
}
}
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
self.0 = Some(format!("{value:?}"));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use tracing_subscriber::Layer;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::registry::Registry;
fn capturing_subscriber() -> impl tracing::Subscriber + Send + Sync {
Registry::default().with(CauseCaptureLayer)
}
fn filtered_subscriber() -> impl tracing::Subscriber + Send + Sync {
Registry::default().with(CauseCaptureLayer.with_filter(filter()))
}
#[tokio::test]
async fn only_the_engines_busy_sentence_reads_as_blocked() {
let _guard = tracing::subscriber::set_default(filtered_subscriber());
let sink = CauseSink::new();
sink.scoped(async {
tracing::debug!(target: ENGINE_TARGET, "PRAGMA wal_checkpoint failed: Busy");
})
.await;
let blocked = sink.attach(anyhow::anyhow!(
"Unexpected result from PRAGMA wal_checkpoint"
));
assert!(
blocked
.to_string()
.contains("engine cause: PRAGMA wal_checkpoint failed: Busy"),
"the captured sentence must reach the error: {blocked:#}"
);
assert!(
is_blocked_checkpoint(&blocked),
"the engine's busy reason must read as a blocked checkpoint: {blocked:#}"
);
for genuine in [
"PRAGMA wal_checkpoint failed: BusySnapshot",
"PRAGMA wal_checkpoint failed: Corrupt(\"page 3\")",
"no engine reason was captured",
] {
let e = anyhow::anyhow!("Unexpected result from PRAGMA wal_checkpoint")
.context(format!("engine cause: {genuine}"));
assert!(
!is_blocked_checkpoint(&e),
"{genuine:?} must not read as a blocked checkpoint: {e:#}"
);
}
}
#[tokio::test]
async fn a_sink_keeps_only_its_own_calls_reason() {
let _guard = tracing::subscriber::set_default(capturing_subscriber());
let first = CauseSink::new();
let second = CauseSink::new();
futures_util::join!(
first.scoped(async {
tracing::debug!(target: ENGINE_TARGET, "PRAGMA wal_checkpoint failed: first");
}),
second.scoped(async {
tracing::debug!(target: ENGINE_TARGET, "PRAGMA wal_checkpoint failed: second");
}),
);
let first_text = format!("{:#}", first.attach(anyhow::anyhow!("first call failed")));
assert!(
first_text.contains("PRAGMA wal_checkpoint failed: first"),
"the first sink must keep its own call's reason: {first_text:?}"
);
assert!(
!first_text.contains("second"),
"the first sink must not see the other call's reason: {first_text:?}"
);
tracing::debug!(target: ENGINE_TARGET, "PRAGMA wal_checkpoint failed: unarmed");
let second_text = format!("{:#}", second.attach(anyhow::anyhow!("second call failed")));
assert!(
second_text.contains("PRAGMA wal_checkpoint failed: second"),
"the second sink must keep its own call's reason: {second_text:?}"
);
assert!(
!second_text.contains("unarmed"),
"an engine event outside any scope must reach no sink: {second_text:?}"
);
}
#[tokio::test]
async fn only_the_engines_checkpoint_sentence_is_kept() {
let _guard = tracing::subscriber::set_default(filtered_subscriber());
let sink = CauseSink::new();
sink.scoped(async {
tracing::debug!(target: ENGINE_TARGET, "a different message");
tracing::debug!(
target: "turso_core::storage::pager",
"PRAGMA wal_checkpoint failed: another target"
);
tracing::debug!(target: ENGINE_TARGET, "PRAGMA wal_checkpoint failed: literal");
})
.await;
let rendered = format!("{:#}", sink.attach(anyhow::anyhow!("no reason captured")));
assert!(
rendered.contains("engine cause: PRAGMA wal_checkpoint failed: literal"),
"the engine's own checkpoint sentence at DEBUG must be captured: {rendered:?}"
);
assert!(
!rendered.contains("a different message"),
"a different message on the engine target must be dropped: {rendered:?}"
);
assert!(
!rendered.contains("another target"),
"the engine's sentence on another target must be filtered out: {rendered:?}"
);
let interpolated = CauseSink::new();
interpolated
.scoped(async {
tracing::debug!(
target: ENGINE_TARGET,
"PRAGMA wal_checkpoint failed: {:?}",
"boom"
);
})
.await;
let rendered = format!(
"{:#}",
interpolated.attach(anyhow::anyhow!(
"Unexpected result from PRAGMA wal_checkpoint"
))
);
assert!(
rendered.contains("engine cause: PRAGMA wal_checkpoint failed:"),
"the captured cause must keep the engine's sentence: {rendered:?}"
);
assert!(
rendered.contains("boom"),
"the captured cause must keep the interpolated diagnostics: {rendered:?}"
);
assert!(
rendered.contains("Unexpected result from PRAGMA wal_checkpoint"),
"the product's own text must survive: {rendered:?}"
);
}
#[tokio::test]
async fn the_production_filter_stack_still_keeps_the_engine_reason() {
use tracing_subscriber::EnvFilter;
let _guard = tracing::subscriber::set_default(crate::logs::log_layers(
std::io::sink,
EnvFilter::new(crate::logs::DEFAULT_LOG_FILTER),
));
let sink = CauseSink::new();
sink.scoped(async {
tracing::debug!(
target: ENGINE_TARGET,
"PRAGMA wal_checkpoint failed: the production stack"
);
})
.await;
let rendered = format!("{:#}", sink.attach(anyhow::anyhow!("no reason captured")));
assert!(
rendered.contains("the production stack"),
"the capture layer must keep the engine's reason under the production filter stack: \
{rendered:?}"
);
}
}