feat: observability logs across entire query + compaction pipeline
CI / CI (pull_request) Successful in 11m38s
CI / CI (pull_request) Successful in 11m38s
All components now emit target="observability" structured logs:
chunk_optimizer:
event=chunk_optimize: input, after_threshold_filter, after_dedup,
dedup_removed, selected, budget_bytes
result_compressor:
event=result_compress: input_count, estimated_bytes, compressed_bytes,
budget_bytes, strategy
query_router:
event=query_route: route, candidates, prefiltered, selected, latency_ms
cache_alignment:
event=cache_preload: preloaded, cache_hits, cache_misses, hit_ratio
full_pipeline:
event=full_pipeline_complete: query, candidates, prefiltered, optimized,
dedup_removed, boosts_applied, cache_hit_ratio, budget_bytes, total_ms
compaction:
event=compaction_complete: mode, duration_ms, duplicate_edges_deleted,
stale_facts_deleted, semantic_merged, llm_calls, bytes_freed
781 tests pass.
This commit is contained in:
@@ -224,9 +224,20 @@ impl KvCacheAligner {
|
||||
|
||||
/// Pre-load hot chunks into cache
|
||||
pub fn preload_hot_chunks(&self, hot_chunks: Vec<(&str, &str)>) -> Result<()> {
|
||||
let count = hot_chunks.len();
|
||||
for (chunk_id, text) in hot_chunks {
|
||||
self.cache.put(chunk_id, text);
|
||||
}
|
||||
let metrics = self.cache.metrics();
|
||||
tracing::info!(
|
||||
target: "observability",
|
||||
event = "cache_preload",
|
||||
preloaded = count,
|
||||
cache_hits = metrics.hits,
|
||||
cache_misses = metrics.misses,
|
||||
hit_ratio = format!("{:.2}", metrics.hit_ratio()),
|
||||
"Cache preload complete"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -211,17 +211,32 @@ impl ChunkOptimizer {
|
||||
|
||||
/// End-to-end optimization pipeline
|
||||
pub fn optimize(&self, chunks: Vec<OptimizableChunk>) -> (Vec<OptimizableChunk>, SelectionMetrics) {
|
||||
let input_count = chunks.len();
|
||||
|
||||
// Step 1: Filter by threshold
|
||||
let filtered = self.threshold_filter.filter(chunks.clone());
|
||||
let after_filter = filtered.len();
|
||||
|
||||
// Step 2: Deduplicate
|
||||
let (deduplicated, dedup_removed) = self.deduplicator.deduplicate(filtered);
|
||||
let after_dedup = deduplicated.len();
|
||||
|
||||
// Step 3: Select within budget
|
||||
let (selected, mut metrics) = self.budget_selector.select(deduplicated);
|
||||
|
||||
metrics.dedup_removed = dedup_removed;
|
||||
|
||||
tracing::info!(
|
||||
target: "observability",
|
||||
event = "chunk_optimize",
|
||||
input = input_count,
|
||||
after_threshold_filter = after_filter,
|
||||
after_dedup = after_dedup,
|
||||
dedup_removed = dedup_removed,
|
||||
selected = selected.len(),
|
||||
budget_bytes = metrics.total_bytes,
|
||||
"Chunk optimization complete"
|
||||
);
|
||||
|
||||
(selected, metrics)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -346,7 +346,19 @@ pub async fn compact_memory(
|
||||
}
|
||||
|
||||
total_stats.duration_ms = start.elapsed().as_millis() as u64;
|
||||
info!("Compaction complete in {}ms: {:?}", total_stats.duration_ms, total_stats);
|
||||
info!(
|
||||
target: "observability",
|
||||
event = "compaction_complete",
|
||||
mode = ?mode,
|
||||
duration_ms = total_stats.duration_ms,
|
||||
duplicate_edges_deleted = total_stats.duplicate_edges_deleted,
|
||||
stale_facts_deleted = total_stats.stale_facts_deleted,
|
||||
semantic_merged = total_stats.semantic_merged,
|
||||
llm_calls = total_stats.llm_calls,
|
||||
bytes_freed = total_stats.bytes_freed,
|
||||
human_reviews_queued = total_stats.human_reviews_queued,
|
||||
"Compaction complete"
|
||||
);
|
||||
|
||||
Ok(total_stats)
|
||||
}
|
||||
|
||||
@@ -344,6 +344,22 @@ impl FullPipeline {
|
||||
|
||||
metrics.total_latency_ms = start.elapsed().as_millis() as u64;
|
||||
|
||||
tracing::info!(
|
||||
target: "observability",
|
||||
event = "full_pipeline_complete",
|
||||
query = query,
|
||||
candidates = metrics.wiki_scope_docs,
|
||||
prefiltered = metrics.prefilter_candidates,
|
||||
optimized = metrics.post_optimization_count,
|
||||
dedup_removed = metrics.dedup_removed,
|
||||
boosts_applied = metrics.metadata_boosts_applied,
|
||||
cache_hit_ratio = format!("{:.2}", metrics.cache_hit_ratio),
|
||||
budget_bytes = metrics.budget_used_bytes,
|
||||
total_ms = metrics.total_latency_ms,
|
||||
"Full query pipeline complete"
|
||||
);
|
||||
|
||||
|
||||
Ok(PipelineResult {
|
||||
query: query.to_string(),
|
||||
query_intent,
|
||||
@@ -467,6 +483,22 @@ impl FullPipeline {
|
||||
|
||||
metrics.total_latency_ms = start.elapsed().as_millis() as u64;
|
||||
|
||||
tracing::info!(
|
||||
target: "observability",
|
||||
event = "full_pipeline_complete",
|
||||
query = query,
|
||||
candidates = metrics.wiki_scope_docs,
|
||||
prefiltered = metrics.prefilter_candidates,
|
||||
optimized = metrics.post_optimization_count,
|
||||
dedup_removed = metrics.dedup_removed,
|
||||
boosts_applied = metrics.metadata_boosts_applied,
|
||||
cache_hit_ratio = format!("{:.2}", metrics.cache_hit_ratio),
|
||||
budget_bytes = metrics.budget_used_bytes,
|
||||
total_ms = metrics.total_latency_ms,
|
||||
"Full query pipeline complete"
|
||||
);
|
||||
|
||||
|
||||
Ok(PipelineResult {
|
||||
query: query.to_string(),
|
||||
query_intent,
|
||||
|
||||
@@ -238,6 +238,17 @@ impl QueryRouter {
|
||||
|
||||
let latency_ms = start.elapsed().as_millis() as u64;
|
||||
|
||||
tracing::info!(
|
||||
target: "observability",
|
||||
event = "query_route",
|
||||
route = "direct",
|
||||
candidates = all_candidates.len(),
|
||||
prefiltered = prefilter_size,
|
||||
selected = selected_chunks.len(),
|
||||
latency_ms = latency_ms,
|
||||
"Query routing complete"
|
||||
);
|
||||
|
||||
Ok(RoutedResult {
|
||||
selected_chunks,
|
||||
route,
|
||||
|
||||
@@ -235,6 +235,18 @@ impl BudgetCompressor {
|
||||
let strategy = self.select_strategy(estimated);
|
||||
let compressed = self.compressor.compress_batch(results, strategy);
|
||||
|
||||
let compressed_size: usize = compressed.iter().map(|c| c.text.as_ref().map_or(0, |t| t.len())).sum();
|
||||
tracing::info!(
|
||||
target: "observability",
|
||||
event = "result_compress",
|
||||
input_count = compressed.len(),
|
||||
estimated_bytes = estimated,
|
||||
compressed_bytes = compressed_size,
|
||||
budget_bytes = self.max_budget_bytes,
|
||||
strategy = ?strategy,
|
||||
"Result compression complete"
|
||||
);
|
||||
|
||||
(compressed, strategy)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user