#[cfg(feature = "decode")]
use keyhog_core::ChunkMetadata;
use serde::{Deserialize, Serialize};
use std::borrow::Cow;
use std::cell::RefCell;
use std::collections::{BTreeMap, HashSet};
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum DogfoodEvent {
ExampleSuppressed {
detector: String,
path: Option<String>,
credential_redacted: String,
reason: Cow<'static, str>,
},
ShapeSuppressed {
path: Option<String>,
credential_redacted: String,
reason: Cow<'static, str>,
},
StaticRecoveryRejected {
path: Option<String>,
expression_offset: usize,
decoder: Cow<'static, str>,
reason: Cow<'static, str>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum StaticRecoveryRejection {
LiteralByteArrayElement,
JsonBase64,
JsonUtf8,
JsonByteArray,
XorPlaintextUtf8,
StringJoinJson,
BufferBase64,
BufferHex,
AesKeyLength,
AesIvLength,
AesCiphertextBlockLength,
AesPadding,
AesPlaintextUtf8,
}
impl StaticRecoveryRejection {
const ALL: [Self; 13] = [
Self::LiteralByteArrayElement,
Self::JsonBase64,
Self::JsonUtf8,
Self::JsonByteArray,
Self::XorPlaintextUtf8,
Self::StringJoinJson,
Self::BufferBase64,
Self::BufferHex,
Self::AesKeyLength,
Self::AesIvLength,
Self::AesCiphertextBlockLength,
Self::AesPadding,
Self::AesPlaintextUtf8,
];
const fn index(self) -> usize {
match self {
Self::LiteralByteArrayElement => 0,
Self::JsonBase64 => 1,
Self::JsonUtf8 => 2,
Self::JsonByteArray => 3,
Self::XorPlaintextUtf8 => 4,
Self::StringJoinJson => 5,
Self::BufferBase64 => 6,
Self::BufferHex => 7,
Self::AesKeyLength => 8,
Self::AesIvLength => 9,
Self::AesCiphertextBlockLength => 10,
Self::AesPadding => 11,
Self::AesPlaintextUtf8 => 12,
}
}
pub(crate) const fn as_str(self) -> &'static str {
match self {
Self::LiteralByteArrayElement => "literal_byte_array_element",
Self::JsonBase64 => "json_base64",
Self::JsonUtf8 => "json_utf8",
Self::JsonByteArray => "json_byte_array",
Self::XorPlaintextUtf8 => "xor_plaintext_utf8",
Self::StringJoinJson => "string_join_json",
Self::BufferBase64 => "buffer_base64",
Self::BufferHex => "buffer_hex",
Self::AesKeyLength => "aes_key_length",
Self::AesIvLength => "aes_iv_length",
Self::AesCiphertextBlockLength => "aes_ciphertext_block_length",
Self::AesPadding => "aes_padding",
Self::AesPlaintextUtf8 => "aes_plaintext_utf8",
}
}
}
pub const DOGFOOD_DETAIL_EVENT_LIMIT: usize = 1024;
fn record_dropped_detail(counter: &AtomicUsize) {
let mut current = counter.load(Ordering::Relaxed);
while current != usize::MAX {
match counter.compare_exchange_weak(
current,
current + 1,
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => return,
Err(observed) => current = observed,
}
}
}
fn push_dogfood_detail(
events: &Mutex<Vec<DogfoodEvent>>,
detail_events_dropped: &AtomicUsize,
event: DogfoodEvent,
) -> bool {
match events.lock() {
Ok(mut events) if events.len() < DOGFOOD_DETAIL_EVENT_LIMIT => {
events.push(event);
true
}
Ok(_) | Err(_) => {
record_dropped_detail(detail_events_dropped);
false
}
}
}
fn recover_telemetry_lock<'a, T>(mutex: &'a Mutex<T>) -> std::sync::MutexGuard<'a, T> {
match mutex.lock() {
Ok(guard) => guard,
Err(poisoned) => {
let guard = poisoned.into_inner();
mutex.clear_poison();
guard
}
}
}
#[derive(Default)]
struct StaticRecoveryTelemetry {
counts: [AtomicU64; StaticRecoveryRejection::ALL.len()],
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
enum EmittedDogfoodKey {
Suppression(String),
#[cfg(feature = "decode")]
StaticRecovery {
source_type: Arc<str>,
path: Option<Arc<str>>,
commit: Option<Arc<str>>,
expression_offset: usize,
reason: &'static str,
},
}
impl StaticRecoveryTelemetry {
fn record(&self, reason: StaticRecoveryRejection) {
self.add(reason, 1);
}
fn add(&self, reason: StaticRecoveryRejection, amount: u64) {
let counter = &self.counts[reason.index()];
let mut current = counter.load(Ordering::Relaxed);
while current != u64::MAX {
let next = current.saturating_add(amount);
match counter.compare_exchange_weak(current, next, Ordering::Relaxed, Ordering::Relaxed)
{
Ok(_) => return,
Err(observed) => current = observed,
}
}
}
fn snapshot(&self) -> BTreeMap<String, u64> {
StaticRecoveryRejection::ALL
.iter()
.filter_map(|reason| {
let count = self.counts[reason.index()].load(Ordering::Relaxed);
(count != 0).then(|| (reason.as_str().to_owned(), count))
})
.collect()
}
fn reset(&self) {
for count in &self.counts {
count.store(0, Ordering::Relaxed);
}
}
}
#[derive(Default)]
struct Telemetry {
dogfood_enabled: AtomicBool,
example_suppressions: AtomicUsize,
events: Mutex<Vec<DogfoodEvent>>,
emitted_suppression_events: Mutex<HashSet<EmittedDogfoodKey>>,
detail_events_dropped: AtomicUsize,
static_recovery: StaticRecoveryTelemetry,
}
#[derive(Default)]
pub struct ScanTelemetry {
dogfood_enabled: AtomicBool,
example_suppressions: AtomicUsize,
events: Mutex<Vec<DogfoodEvent>>,
emitted_suppression_events: Mutex<HashSet<EmittedDogfoodKey>>,
detail_events_dropped: AtomicUsize,
static_recovery: StaticRecoveryTelemetry,
}
impl ScanTelemetry {
pub fn new() -> Self {
Self::default()
}
pub fn enable_dogfood(&self) {
self.dogfood_enabled.store(true, Ordering::Relaxed);
}
fn is_dogfood_enabled(&self) -> bool {
self.dogfood_enabled.load(Ordering::Relaxed)
}
fn example_suppression_count(&self) -> usize {
self.example_suppressions.load(Ordering::Relaxed)
}
fn drain_events(&self) -> Vec<DogfoodEvent> {
drain_event_buffers(&self.events, &self.emitted_suppression_events)
}
pub fn drain(&self) -> ScanTelemetrySnapshot {
ScanTelemetrySnapshot {
example_suppressions: self.example_suppression_count() as u64,
dogfood_events: self.drain_events(),
dogfood_detail_events_dropped: self.detail_events_dropped.load(Ordering::Relaxed)
as u64,
static_recovery_rejections: self.static_recovery.snapshot(),
}
}
}
pub struct ScanTelemetrySnapshot {
pub example_suppressions: u64,
pub dogfood_events: Vec<DogfoodEvent>,
pub dogfood_detail_events_dropped: u64,
pub static_recovery_rejections: BTreeMap<String, u64>,
}
thread_local! {
static CURRENT_SCAN_TELEMETRY: RefCell<Option<Arc<ScanTelemetry>>> = RefCell::new(None);
}
struct ScanTelemetryRestore {
previous: Option<Arc<ScanTelemetry>>,
}
impl Drop for ScanTelemetryRestore {
fn drop(&mut self) {
let previous = self.previous.take();
CURRENT_SCAN_TELEMETRY.with(|slot| {
*slot.borrow_mut() = previous;
});
}
}
pub fn with_scan_telemetry<R>(telemetry: &Arc<ScanTelemetry>, f: impl FnOnce() -> R) -> R {
let previous = CURRENT_SCAN_TELEMETRY.with(|slot| {
let mut slot = slot.borrow_mut();
slot.replace(Arc::clone(telemetry))
});
let _restore = ScanTelemetryRestore { previous };
f()
}
fn current_scan_telemetry() -> Option<Arc<ScanTelemetry>> {
CURRENT_SCAN_TELEMETRY.with(|slot| slot.borrow().clone())
}
pub(crate) fn capture_scan_telemetry() -> Option<Arc<ScanTelemetry>> {
current_scan_telemetry()
}
pub(crate) fn with_captured_scan_telemetry<R>(
telemetry: Option<&Arc<ScanTelemetry>>,
f: impl FnOnce() -> R,
) -> R {
match telemetry {
Some(telemetry) => with_scan_telemetry(telemetry, f),
None => f(),
}
}
fn current_scan_dogfood_enabled() -> Option<bool> {
CURRENT_SCAN_TELEMETRY.with(|slot| {
slot.borrow()
.as_ref()
.map(|telemetry| telemetry.is_dogfood_enabled())
})
}
static FILES_SCANNED: AtomicUsize = AtomicUsize::new(0);
static BYTES_SCANNED: AtomicUsize = AtomicUsize::new(0);
static SKIPPED_FILES: AtomicUsize = AtomicUsize::new(0);
static TOTAL_MATCHES: AtomicUsize = AtomicUsize::new(0);
static GPU_DISPATCHES: AtomicUsize = AtomicUsize::new(0);
static STRUCTURED_PARSE_FAILURES: AtomicUsize = AtomicUsize::new(0);
static STRUCTURED_OVERSIZE_SKIPS: AtomicUsize = AtomicUsize::new(0);
static DECODE_TRUNCATIONS: AtomicUsize = AtomicUsize::new(0);
#[cfg(test)]
thread_local! {
static THREAD_DECODE_TRUNCATIONS: std::cell::Cell<usize> =
const { std::cell::Cell::new(0) };
}
static INVALID_PATTERN_INDEX_SKIPS: AtomicUsize = AtomicUsize::new(0);
static BOUNDARY_RESULT_CARDINALITY_MISMATCHES: AtomicUsize = AtomicUsize::new(0);
static LINE_OFFSET_MAPPING_MISMATCHES: AtomicUsize = AtomicUsize::new(0);
static CHUNK_DEADLINE_ABORTS: AtomicUsize = AtomicUsize::new(0);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ScannerCoverageGapEvent {
StructuredParseFailure,
StructuredOversizeSkip,
DecodeTruncation,
InvalidPatternIndexSkip,
BoundaryResultCardinalityMismatch,
LineOffsetMappingMismatch,
ChunkDeadlineAbort,
}
impl ScannerCoverageGapEvent {
pub(crate) const ALL: [Self; 7] = [
Self::StructuredParseFailure,
Self::StructuredOversizeSkip,
Self::DecodeTruncation,
Self::InvalidPatternIndexSkip,
Self::BoundaryResultCardinalityMismatch,
Self::LineOffsetMappingMismatch,
Self::ChunkDeadlineAbort,
];
pub(crate) fn counter(self) -> &'static AtomicUsize {
match self {
Self::StructuredParseFailure => &STRUCTURED_PARSE_FAILURES,
Self::StructuredOversizeSkip => &STRUCTURED_OVERSIZE_SKIPS,
Self::DecodeTruncation => &DECODE_TRUNCATIONS,
Self::InvalidPatternIndexSkip => &INVALID_PATTERN_INDEX_SKIPS,
Self::BoundaryResultCardinalityMismatch => &BOUNDARY_RESULT_CARDINALITY_MISMATCHES,
Self::LineOffsetMappingMismatch => &LINE_OFFSET_MAPPING_MISMATCHES,
Self::ChunkDeadlineAbort => &CHUNK_DEADLINE_ABORTS,
}
}
const fn label(self) -> &'static str {
match self {
Self::StructuredParseFailure => "structured_parse_failures",
Self::StructuredOversizeSkip => "structured_oversize_skips",
Self::DecodeTruncation => "decode_truncations",
Self::InvalidPatternIndexSkip => "invalid_pattern_index_skips",
Self::BoundaryResultCardinalityMismatch => "boundary_result_cardinality_mismatches",
Self::LineOffsetMappingMismatch => "line_offset_mapping_mismatches",
Self::ChunkDeadlineAbort => "chunk_deadline_aborts",
}
}
}
#[derive(Clone, Copy, Default, Eq, PartialEq)]
pub struct ScannerCoverageSnapshot {
counts: [usize; ScannerCoverageGapEvent::ALL.len()],
}
impl ScannerCoverageSnapshot {
#[must_use]
pub fn capture() -> Self {
Self {
counts: std::array::from_fn(|index| {
ScannerCoverageGapEvent::ALL[index]
.counter()
.load(Ordering::Relaxed)
}),
}
}
#[must_use]
pub fn saturating_delta(self, earlier: Self) -> Self {
Self {
counts: std::array::from_fn(|index| {
self.counts[index].saturating_sub(earlier.counts[index])
}),
}
}
#[must_use]
pub fn is_empty(self) -> bool {
self.counts.iter().all(|count| *count == 0)
}
}
impl std::fmt::Debug for ScannerCoverageSnapshot {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let mut gaps = formatter.debug_map();
for (event, count) in ScannerCoverageGapEvent::ALL.into_iter().zip(self.counts) {
if count > 0 {
gaps.entry(&event.label(), &count);
}
}
gaps.finish()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[must_use = "scanner coverage gaps must be recorded through the typed recorder so partial coverage remains surfaced"]
pub(crate) struct RecordedScannerCoverageGap {
event: ScannerCoverageGapEvent,
previous: usize,
delta: usize,
}
pub(crate) fn record_scanner_coverage_gap(
event: ScannerCoverageGapEvent,
) -> RecordedScannerCoverageGap {
let previous = event.counter().fetch_add(1, Ordering::Relaxed);
RecordedScannerCoverageGap {
event,
previous,
delta: 1,
}
}
static DOGFOOD_ENABLED: AtomicBool = AtomicBool::new(false);
fn cell() -> &'static Telemetry {
static CELL: OnceLock<Telemetry> = OnceLock::new();
CELL.get_or_init(Telemetry::default)
}
pub fn enable_dogfood() {
DOGFOOD_ENABLED.store(true, Ordering::Relaxed);
cell().dogfood_enabled.store(true, Ordering::Relaxed);
}
pub fn is_dogfood_enabled() -> bool {
if let Some(enabled) = current_scan_dogfood_enabled() {
return enabled;
}
DOGFOOD_ENABLED.load(Ordering::Relaxed)
}
pub fn record_example_suppression(
detector: &str,
path: Option<&str>,
credential: &str,
reason: &'static str,
) {
if let Some(t) = current_scan_telemetry() {
record_example_suppression_in(
&t.example_suppressions,
&t.events,
&t.emitted_suppression_events,
&t.detail_events_dropped,
detector,
path,
credential,
reason,
);
return;
}
let t = cell();
record_example_suppression_in(
&t.example_suppressions,
&t.events,
&t.emitted_suppression_events,
&t.detail_events_dropped,
detector,
path,
credential,
reason,
);
}
fn record_example_suppression_in(
example_suppressions: &AtomicUsize,
events: &Mutex<Vec<DogfoodEvent>>,
emitted_suppression_events: &Mutex<HashSet<EmittedDogfoodKey>>,
detail_events_dropped: &AtomicUsize,
detector: &str,
path: Option<&str>,
credential: &str,
reason: &'static str,
) {
example_suppressions.fetch_add(1, Ordering::Relaxed);
if !is_dogfood_enabled() {
return;
}
let credential_hash = keyhog_core::hex_encode(&keyhog_core::sha256_hash(credential));
if !mark_suppression_event_emitted(
emitted_suppression_events,
detail_events_dropped,
&credential_hash,
) {
return;
}
let redacted = keyhog_core::redact(credential).into_owned();
push_dogfood_detail(
events,
detail_events_dropped,
DogfoodEvent::ExampleSuppressed {
detector: detector.to_string(),
path: path.map(str::to_string),
credential_redacted: redacted,
reason: Cow::Borrowed(reason),
},
);
}
fn mark_suppression_event_emitted(
emitted_suppression_events: &Mutex<HashSet<EmittedDogfoodKey>>,
detail_events_dropped: &AtomicUsize,
credential_hash: &str,
) -> bool {
match emitted_suppression_events.lock() {
Ok(mut emitted) => {
let key = EmittedDogfoodKey::Suppression(credential_hash.to_owned());
if emitted.contains(&key) {
return false;
}
if emitted.len() >= DOGFOOD_DETAIL_EVENT_LIMIT {
record_dropped_detail(detail_events_dropped);
return false;
}
emitted.insert(key)
}
Err(_) => {
record_dropped_detail(detail_events_dropped);
false }
}
}
pub(crate) fn record_shape_suppression(path: Option<&str>, credential: &str, reason: &'static str) {
if !is_dogfood_enabled() {
return;
}
if let Some(t) = current_scan_telemetry() {
record_shape_suppression_in(
&t.events,
&t.emitted_suppression_events,
&t.detail_events_dropped,
path,
credential,
reason,
);
return;
}
let t = cell();
record_shape_suppression_in(
&t.events,
&t.emitted_suppression_events,
&t.detail_events_dropped,
path,
credential,
reason,
);
}
#[cfg(feature = "decode")]
pub(crate) fn record_static_recovery_rejection(
metadata: &ChunkMetadata,
expression_offset: usize,
reason: StaticRecoveryRejection,
) {
if !is_dogfood_enabled() {
return;
}
if let Some(t) = current_scan_telemetry() {
t.static_recovery.record(reason);
if !mark_static_recovery_event_emitted(
&t.emitted_suppression_events,
&t.detail_events_dropped,
metadata,
expression_offset,
reason,
) {
return;
}
push_dogfood_detail(
&t.events,
&t.detail_events_dropped,
static_recovery_event(metadata, expression_offset, reason),
);
return;
}
let t = cell();
t.static_recovery.record(reason);
if !mark_static_recovery_event_emitted(
&t.emitted_suppression_events,
&t.detail_events_dropped,
metadata,
expression_offset,
reason,
) {
return;
}
push_dogfood_detail(
&t.events,
&t.detail_events_dropped,
static_recovery_event(metadata, expression_offset, reason),
);
}
#[cfg(feature = "decode")]
fn static_recovery_event(
metadata: &ChunkMetadata,
expression_offset: usize,
reason: StaticRecoveryRejection,
) -> DogfoodEvent {
DogfoodEvent::StaticRecoveryRejected {
path: metadata.path.as_deref().map(str::to_owned),
expression_offset,
decoder: Cow::Borrowed("javascript-static"),
reason: Cow::Borrowed(reason.as_str()),
}
}
#[cfg(feature = "decode")]
fn mark_static_recovery_event_emitted(
emitted_events: &Mutex<HashSet<EmittedDogfoodKey>>,
detail_events_dropped: &AtomicUsize,
metadata: &ChunkMetadata,
expression_offset: usize,
reason: StaticRecoveryRejection,
) -> bool {
let key = EmittedDogfoodKey::StaticRecovery {
source_type: Arc::clone(&metadata.source_type),
path: metadata.path.clone(),
commit: metadata.commit.clone(),
expression_offset,
reason: reason.as_str(),
};
match emitted_events.lock() {
Ok(mut emitted) => {
if emitted.contains(&key) {
return false;
}
if emitted.len() >= DOGFOOD_DETAIL_EVENT_LIMIT {
record_dropped_detail(detail_events_dropped);
return false;
}
emitted.insert(key)
}
Err(_) => {
record_dropped_detail(detail_events_dropped);
false }
}
}
fn record_shape_suppression_in(
events: &Mutex<Vec<DogfoodEvent>>,
emitted_suppression_events: &Mutex<HashSet<EmittedDogfoodKey>>,
detail_events_dropped: &AtomicUsize,
path: Option<&str>,
credential: &str,
reason: &'static str,
) {
let credential_hash = keyhog_core::hex_encode(&keyhog_core::sha256_hash(credential));
if !mark_suppression_event_emitted(
emitted_suppression_events,
detail_events_dropped,
&credential_hash,
) {
return;
}
let redacted = keyhog_core::redact(credential).into_owned();
push_dogfood_detail(
events,
detail_events_dropped,
DogfoodEvent::ShapeSuppressed {
path: path.map(str::to_string),
credential_redacted: redacted,
reason: Cow::Borrowed(reason),
},
);
}
pub fn example_suppression_count() -> usize {
cell().example_suppressions.load(Ordering::Relaxed)
}
#[cfg(test)]
pub(crate) fn reset_example_suppression_count() {
cell().example_suppressions.store(0, Ordering::Relaxed);
}
pub fn add_example_suppressions(n: usize) {
cell().example_suppressions.fetch_add(n, Ordering::Relaxed);
}
pub(crate) fn record_structured_parse_failure() {
let _receipt = record_scanner_coverage_gap(ScannerCoverageGapEvent::StructuredParseFailure);
}
pub fn structured_parse_failure_count() -> usize {
STRUCTURED_PARSE_FAILURES.load(Ordering::Relaxed)
}
pub(crate) fn record_structured_oversize_skip() {
let _receipt = record_scanner_coverage_gap(ScannerCoverageGapEvent::StructuredOversizeSkip);
}
pub fn structured_oversize_skip_count() -> usize {
STRUCTURED_OVERSIZE_SKIPS.load(Ordering::Relaxed)
}
pub(crate) fn record_decode_truncation() {
let _receipt = record_scanner_coverage_gap(ScannerCoverageGapEvent::DecodeTruncation);
#[cfg(test)]
THREAD_DECODE_TRUNCATIONS.with(|count| count.set(count.get() + 1));
}
#[cfg(not(test))]
pub fn decode_truncation_count() -> usize {
DECODE_TRUNCATIONS.load(Ordering::Relaxed)
}
#[cfg(test)]
pub fn decode_truncation_count() -> usize {
THREAD_DECODE_TRUNCATIONS.with(|count| count.get())
}
pub(crate) fn record_invalid_pattern_index_skip() {
let _receipt = record_scanner_coverage_gap(ScannerCoverageGapEvent::InvalidPatternIndexSkip);
}
pub fn invalid_pattern_index_skip_count() -> usize {
INVALID_PATTERN_INDEX_SKIPS.load(Ordering::Relaxed)
}
pub(crate) fn record_boundary_result_cardinality_mismatch() {
let _receipt =
record_scanner_coverage_gap(ScannerCoverageGapEvent::BoundaryResultCardinalityMismatch);
}
pub fn boundary_result_cardinality_mismatch_count() -> usize {
BOUNDARY_RESULT_CARDINALITY_MISMATCHES.load(Ordering::Relaxed)
}
#[cfg(feature = "multiline")]
pub(crate) fn record_line_offset_mapping_mismatch() {
let _receipt = record_scanner_coverage_gap(ScannerCoverageGapEvent::LineOffsetMappingMismatch);
}
pub(crate) fn record_chunk_deadline_abort() {
let _receipt = record_scanner_coverage_gap(ScannerCoverageGapEvent::ChunkDeadlineAbort);
}
pub fn chunk_deadline_abort_count() -> usize {
CHUNK_DEADLINE_ABORTS.load(Ordering::Relaxed)
}
pub fn line_offset_mapping_mismatch_count() -> usize {
LINE_OFFSET_MAPPING_MISMATCHES.load(Ordering::Relaxed)
}
pub fn append_events<I: IntoIterator<Item = DogfoodEvent>>(events: I) {
append_event_details(events, true);
}
pub fn append_daemon_events<I: IntoIterator<Item = DogfoodEvent>>(events: I) {
append_event_details(events, false);
}
fn append_event_details<I: IntoIterator<Item = DogfoodEvent>>(
events: I,
infer_static_recovery_counts: bool,
) {
let t = cell();
for event in events {
if infer_static_recovery_counts {
let DogfoodEvent::StaticRecoveryRejected { reason, .. } = &event else {
push_dogfood_detail(&t.events, &t.detail_events_dropped, event);
continue;
};
if let Some(reason) = StaticRecoveryRejection::ALL
.iter()
.find(|candidate| candidate.as_str() == reason.as_ref())
{
t.static_recovery.record(*reason);
}
}
push_dogfood_detail(&t.events, &t.detail_events_dropped, event);
}
}
pub fn merge_daemon_aggregates(
static_recovery_rejections: &BTreeMap<String, u64>,
detail_events_dropped: u64,
) -> Result<(), String> {
let mut resolved = Vec::with_capacity(static_recovery_rejections.len());
for (name, count) in static_recovery_rejections {
let Some(reason) = StaticRecoveryRejection::ALL
.iter()
.copied()
.find(|candidate| candidate.as_str() == name)
else {
return Err(format!(
"daemon returned unknown static-recovery rejection reason {name:?}; restart it with this KeyHog build"
));
};
resolved.push((reason, *count));
}
let telemetry = cell();
for (reason, count) in resolved {
telemetry.static_recovery.add(reason, count);
}
let dropped = usize::try_from(detail_events_dropped).unwrap_or(usize::MAX); let counter = &telemetry.detail_events_dropped;
let mut current = counter.load(Ordering::Relaxed);
while current != usize::MAX {
match counter.compare_exchange_weak(
current,
current.saturating_add(dropped),
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => break,
Err(observed) => current = observed,
}
}
Ok(())
}
pub fn static_recovery_rejection_counts() -> BTreeMap<String, u64> {
cell().static_recovery.snapshot()
}
pub fn dogfood_detail_events_dropped() -> usize {
cell().detail_events_dropped.load(Ordering::Relaxed)
}
pub fn drain_events() -> Vec<DogfoodEvent> {
let t = cell();
drain_event_buffers(&t.events, &t.emitted_suppression_events)
}
fn drain_event_buffers(
events: &Mutex<Vec<DogfoodEvent>>,
emitted_suppression_events: &Mutex<HashSet<EmittedDogfoodKey>>,
) -> Vec<DogfoodEvent> {
recover_telemetry_lock(emitted_suppression_events).clear();
std::mem::take(&mut *recover_telemetry_lock(events))
}
pub(crate) fn record_file_scanned(bytes: usize) {
FILES_SCANNED.fetch_add(1, Ordering::Relaxed);
BYTES_SCANNED.fetch_add(bytes, Ordering::Relaxed);
}
pub(crate) fn global_scan_counts() -> (usize, usize) {
(
FILES_SCANNED.load(Ordering::Relaxed),
BYTES_SCANNED.load(Ordering::Relaxed),
)
}
pub(crate) fn record_file_skipped() {
SKIPPED_FILES.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_match_found() {
TOTAL_MATCHES.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_gpu_dispatch() {
GPU_DISPATCHES.fetch_add(1, Ordering::Relaxed);
}
pub fn reset_for_scan() {
let t = cell();
DOGFOOD_ENABLED.store(false, Ordering::Relaxed);
t.dogfood_enabled.store(false, Ordering::Relaxed);
t.example_suppressions.store(0, Ordering::Relaxed);
t.detail_events_dropped.store(0, Ordering::Relaxed);
t.static_recovery.reset();
FILES_SCANNED.store(0, Ordering::Relaxed);
BYTES_SCANNED.store(0, Ordering::Relaxed);
SKIPPED_FILES.store(0, Ordering::Relaxed);
TOTAL_MATCHES.store(0, Ordering::Relaxed);
GPU_DISPATCHES.store(0, Ordering::Relaxed);
for gap in ScannerCoverageGapEvent::ALL {
gap.counter().store(0, Ordering::Relaxed);
}
#[cfg(test)]
THREAD_DECODE_TRUNCATIONS.with(|count| count.set(0));
recover_telemetry_lock(&t.events).clear();
recover_telemetry_lock(&t.emitted_suppression_events).clear();
CURRENT_SCAN_TELEMETRY.with(|slot| {
*slot.borrow_mut() = None;
});
}
#[cfg(test)]
#[doc(hidden)]
pub mod testing {
use std::sync::Arc;
pub fn reset() {
super::reset_for_scan();
}
pub(crate) fn poison_events(telemetry: &Arc<super::ScanTelemetry>) {
let telemetry = Arc::clone(telemetry);
let _ = std::thread::spawn(move || {
let Ok(_events) = telemetry.events.lock() else {
panic!("fresh telemetry event buffer was already poisoned");
};
panic!("poison scoped telemetry event buffer");
})
.join();
}
}