1use super::*;
2
3impl KnowledgeBase {
4 pub fn inspect(&self) -> Result<Value> {
5 let total: i64 = count_query(
6 &self.storage,
7 "SELECT COUNT(*) FROM chunks WHERE origin!='spark'",
8 )?;
9 let active: i64 = count_query(
10 &self.storage,
11 "SELECT COUNT(*) FROM chunks WHERE state='active' AND origin!='spark'",
12 )?;
13 let pending: i64 = count_query(
14 &self.storage,
15 "SELECT COUNT(*) FROM chunks WHERE state='pending' AND origin!='spark'",
16 )?;
17 let archived: i64 = count_query(
18 &self.storage,
19 "SELECT COUNT(*) FROM chunks WHERE state='archived' AND origin!='spark'",
20 )?;
21 let sparks: i64 = count_query(
22 &self.storage,
23 "SELECT COUNT(*) FROM chunks WHERE origin='spark' AND state!='archived'",
24 )?;
25 let open_logs: i64 = count_query(
26 &self.storage,
27 "SELECT COUNT(*) FROM episodic_log WHERE distill_state='open'",
28 )?;
29 let new_logs: i64 = count_query(
30 &self.storage,
31 "SELECT COUNT(*) FROM episodic_log WHERE distill_state='new'",
32 )?;
33 let embed_rebuild: i64 = count_query(&self.storage,
34 "SELECT COUNT(*) FROM chunks WHERE embed_version=0 OR embed_version < (SELECT COALESCE(CAST(value AS INTEGER),1) FROM meta WHERE key='embed_version')")?;
35 let schema_version = self.storage.get_meta_or("schema_version", "?");
36 let lib_id = self.storage.get_meta_or("lib_id", "?");
37 let last_agg = self.storage.get_meta_or("last_agg_ts", "never");
38
39 let metric_window_start = days_ago(&utc_now_iso(), 30);
40 let trace_metrics = self.storage.query_chunks_params(
41 "SELECT COUNT(*) AS total,
42 SUM(CASE WHEN task_state='completed' THEN 1 ELSE 0 END) AS completed,
43 SUM(CASE WHEN task_state='timed_out' THEN 1 ELSE 0 END) AS timed_out,
44 SUM(CASE WHEN task_state='completed' AND usage_state!='unknown'
45 THEN 1 ELSE 0 END) AS usage_known,
46 SUM(CASE WHEN task_state='completed' AND usage_state='known_some'
47 THEN 1 ELSE 0 END) AS usage_some,
48 SUM(CASE WHEN task_state='completed'
49 AND outcome IN ('ok','fail')
50 THEN 1 ELSE 0 END) AS outcome_known,
51 SUM(CASE WHEN outcome='ok' THEN 1 ELSE 0 END) AS succeeded
52 FROM episodic_log WHERE ts >= ?",
53 rusqlite::params![metric_window_start],
54 )?;
55 let trace_row = trace_metrics.first();
56 let trace_total = trace_row
57 .and_then(|row| row.get("total"))
58 .and_then(Value::as_i64)
59 .unwrap_or(0);
60 let trace_completed = trace_row
61 .and_then(|row| row.get("completed"))
62 .and_then(Value::as_i64)
63 .unwrap_or(0);
64 let trace_timed_out = trace_row
65 .and_then(|row| row.get("timed_out"))
66 .and_then(Value::as_i64)
67 .unwrap_or(0);
68 let usage_known = trace_row
69 .and_then(|row| row.get("usage_known"))
70 .and_then(Value::as_i64)
71 .unwrap_or(0);
72 let usage_some = trace_row
73 .and_then(|row| row.get("usage_some"))
74 .and_then(Value::as_i64)
75 .unwrap_or(0);
76 let succeeded = trace_row
77 .and_then(|row| row.get("succeeded"))
78 .and_then(Value::as_i64)
79 .unwrap_or(0);
80 let outcome_known = trace_row
81 .and_then(|row| row.get("outcome_known"))
82 .and_then(Value::as_i64)
83 .unwrap_or(0);
84 let usage_rows = self.storage.query_chunks_params(
85 "SELECT recall_snapshot, used_ids FROM episodic_log
86 WHERE task_state='completed'
87 AND usage_state!='unknown' AND used_complete=1
88 AND recall_snapshot IS NOT NULL AND used_ids IS NOT NULL
89 AND ts >= ?",
90 rusqlite::params![metric_window_start],
91 )?;
92 let mut selected_total = 0_i64;
93 let mut selected_used = 0_i64;
94 for row in usage_rows {
95 let selected: HashSet<String> = row
96 .get("recall_snapshot")
97 .and_then(Value::as_str)
98 .and_then(|raw| serde_json::from_str::<Value>(raw).ok())
99 .and_then(|snapshot| snapshot.get("selected").cloned())
100 .and_then(|value| serde_json::from_value::<Vec<String>>(value).ok())
101 .unwrap_or_default()
102 .into_iter()
103 .collect();
104 let used: HashSet<String> = row
105 .get("used_ids")
106 .and_then(Value::as_str)
107 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
108 .unwrap_or_default()
109 .into_iter()
110 .collect();
111 selected_total += selected.len() as i64;
112 selected_used += selected.intersection(&used).count() as i64;
113 }
114 let feedback_count = count_query_params(
115 &self.storage,
116 "SELECT COUNT(*) FROM feedback_events WHERE ts >= ?",
117 rusqlite::params![metric_window_start],
118 )?;
119 let feedback_traces = count_query_params(
120 &self.storage,
121 "SELECT COUNT(DISTINCT f.trace_id)
122 FROM feedback_events f
123 JOIN episodic_log e ON e.trace_id=f.trace_id
124 WHERE f.ts >= ? AND e.ts >= ? AND e.task_state='completed'",
125 rusqlite::params![metric_window_start, metric_window_start],
126 )?;
127 let pending_evolve = count_query(
128 &self.storage,
129 "SELECT COUNT(*) FROM evolve_requests WHERE state IN ('pending','running')",
130 )?;
131 let governance_pending = count_query(
132 &self.storage,
133 "SELECT COUNT(*) FROM governance_proposals WHERE state='pending'",
134 )?;
135 let failed_evolve = count_query_params(
136 &self.storage,
137 "SELECT COUNT(*) FROM evolve_requests
138 WHERE last_failed_at >= ?",
139 rusqlite::params![metric_window_start],
140 )?;
141 let failed_distill = count_query_params(
142 &self.storage,
143 "SELECT COUNT(*) FROM episodic_log
144 WHERE distill_last_failed_at >= ?",
145 rusqlite::params![metric_window_start],
146 )?;
147 let confidence_buckets = self.storage.query_chunks(&format!(
148 "SELECT
149 SUM(CASE WHEN confidence < 0.25 THEN 1 ELSE 0 END) AS low,
150 SUM(CASE WHEN confidence >= 0.25 AND confidence < {0} THEN 1 ELSE 0 END) AS medium,
151 SUM(CASE WHEN confidence >= {0} THEN 1 ELSE 0 END) AS high
152 FROM chunks WHERE origin!='spark' AND state!='archived'",
153 self.promote_confidence_min
154 ))?;
155 let confidence_row = confidence_buckets.first();
156
157 let pending_oldest_ts = self.storage.query_chunks(
159 "SELECT MIN(created_at) AS oldest FROM chunks WHERE state='pending' AND origin!='spark'",
160 )?.into_iter().next()
161 .and_then(|r| r.get("oldest").cloned())
162 .filter(|v| !v.is_null());
163
164 let zombie_cutoff = days_ago(&utc_now_iso(), 14);
169 let zombie: i64 = count_query_params(
170 &self.storage,
171 "SELECT COUNT(*) FROM chunks
172 WHERE origin!='spark' AND state='active'
173 AND confidence >= 0.4 AND confidence <= 0.6
174 AND last_used_at IS NOT NULL
175 AND created_at < ?",
176 rusqlite::params![zombie_cutoff],
177 )?;
178 let debt_numerator = pending + zombie;
179 let debt_denominator = active.max(1);
180 let debt_ratio = debt_numerator as f64 / debt_denominator as f64;
181
182 let screening_cutoff = minutes_ago(&utc_now_iso(), self.screening_timeout_minutes);
184 let stale_screening: i64 = count_query_params(
185 &self.storage,
186 "SELECT COUNT(*) FROM episodic_log
187 WHERE distill_state='screening' AND distill_locked_at < ?",
188 rusqlite::params![screening_cutoff],
189 )?;
190
191 let distill_period_start = self.distill_token_period_start(&utc_now_iso())?;
193 let distill_cost = self.storage.query_chunks_params(
194 "SELECT COALESCE(SUM(prompt_tokens),0) AS pt,
195 COALESCE(SUM(completion_tokens),0) AS ct
196 FROM distill_token_usage
197 WHERE accounted_at >= ?",
198 rusqlite::params![distill_period_start],
199 )?;
200 let prompt_tokens = distill_cost
201 .first()
202 .and_then(|r| r.get("pt"))
203 .and_then(Value::as_i64)
204 .unwrap_or(0);
205 let completion_tokens = distill_cost
206 .first()
207 .and_then(|r| r.get("ct"))
208 .and_then(Value::as_i64)
209 .unwrap_or(0);
210
211 let spark_threshold: i64 = self
213 .storage
214 .get_meta("curate.soft_mature_threshold")
215 .ok()
216 .flatten()
217 .and_then(|v| v.parse::<i64>().ok())
218 .unwrap_or(5);
219 let recurring_sparks = self.storage.query_chunks_params(
220 "SELECT ut.chunk_id, COUNT(*) AS cnt,
221 c.content, c.trigger_desc, c.maturity
222 FROM usage_trace ut
223 JOIN chunks c ON c.id = ut.chunk_id
224 WHERE ut.event='retrieved'
225 AND c.origin='spark'
226 GROUP BY ut.chunk_id HAVING cnt >= ?",
227 rusqlite::params![spark_threshold],
228 )?;
229 let recurring_spark_ids: Vec<Value> = recurring_sparks
230 .iter()
231 .map(|r| {
232 json!({
233 "id": r.get("chunk_id").and_then(Value::as_str).unwrap_or(""),
234 "retrieved_count": r.get("cnt").and_then(Value::as_i64).unwrap_or(0),
235 "maturity": r.get("maturity").and_then(Value::as_str).unwrap_or(""),
236 "content_preview": r.get("content").and_then(Value::as_str).unwrap_or("")
237 .chars().take(80).collect::<String>(),
238 })
239 })
240 .collect();
241
242 let mut suggestions: Vec<Value> = Vec::new();
243 if embed_rebuild > 0 {
244 suggestions.push(json!({"action": "innate evolve --rebuild-embeddings", "reason": format!("{embed_rebuild} chunk(s) missing embeddings")}));
245 }
246 if new_logs > 0 {
247 suggestions.push(json!({"action": "innate evolve --trigger manual", "reason": format!("{new_logs} episodic log(s) ready to distill")}));
248 }
249 if pending > 0 {
250 suggestions.push(json!({"action": "innate approve <id> # or innate archive <id>", "reason": format!("{pending} pending chunk(s) awaiting review")}));
251 }
252 if !recurring_spark_ids.is_empty() {
253 suggestions.push(json!({"action": "innate promote-spark <id> --to note", "reason": format!("{} spark(s) recalled ≥{spark_threshold}× — consider promoting", recurring_spark_ids.len())}));
254 }
255 if stale_screening > 0 {
256 suggestions.push(json!({"action": "innate evolve --trigger manual", "reason": format!("{stale_screening} episodic log(s) stuck in screening")}));
257 }
258 if governance_pending > 0 {
259 suggestions.push(json!({
260 "action": "review governance_proposals",
261 "reason": format!("{governance_pending} chunk(s) have repeated negative feedback")
262 }));
263 }
264 let mcpcli_null: i64 = count_query_params(
269 &self.storage,
270 "SELECT COUNT(*) FROM episodic_log
271 WHERE event_source IN ('mcp','cli') AND agent IS NULL AND ts >= ?1",
272 rusqlite::params![days_ago(&utc_now_iso(), 7)],
273 )?;
274 if mcpcli_null > 0 {
275 suggestions.push(json!({
276 "action": "set INNATE_AGENT in the agent's MCP env (innate install)",
277 "reason": format!("{mcpcli_null} mcp/cli call(s) in 7d have NULL agent — agent attribution misconfigured")
278 }));
279 }
280
281 let intuition = self.intuition_calibration(&metric_window_start)?;
284 let appraisals = intuition
285 .get("appraisals")
286 .and_then(Value::as_i64)
287 .unwrap_or(0);
288 let mono_gap = intuition
289 .get("monotonicity_gap")
290 .and_then(Value::as_f64)
291 .unwrap_or(0.0);
292 let false_alarm = intuition
293 .get("false_alarm_rate")
294 .and_then(Value::as_f64)
295 .unwrap_or(0.0);
296 if appraisals >= 20 && mono_gap <= 0.0 {
297 suggestions.push(json!({
298 "action": "tune recall.w_* / situation.coarse_keys",
299 "reason": "appraise strength may be noise — strong tier does not beat weak on task_ok"
300 }));
301 }
302 if appraisals >= 20 && false_alarm >= 0.5 {
303 suggestions.push(json!({
304 "action": "review caution chunks / raise appraise.tier_strong",
305 "reason": format!("intuition false-alarm rate {false_alarm} — strong cautions often end ok")
306 }));
307 }
308
309 let usage_trace_total = count_query(&self.storage, "SELECT COUNT(*) FROM usage_trace")?;
313 let episodic_log_total = count_query(&self.storage, "SELECT COUNT(*) FROM episodic_log")?;
314 let page_count = count_query(&self.storage, "PRAGMA page_count")?;
315 let page_size = count_query(&self.storage, "PRAGMA page_size")?;
316 let db_size_bytes = page_count * page_size;
317
318 let now_obs = utc_now_iso();
321 let observability = self.observability_block(&now_obs)?;
322 let operational = self.operational_block(&now_obs)?;
323 let trends = self.trends_block(&now_obs)?;
324
325 let mut out = json!({
326 "schema_version": schema_version,
327 "lib_id": lib_id,
328 "last_agg_ts": last_agg,
329 "chunks": {
330 "total": total, "active": active, "pending": pending, "archived": archived,
331 "pending_oldest_ts": pending_oldest_ts,
332 },
333 "storage": {
334 "usage_trace_rows": usage_trace_total,
335 "episodic_log_rows": episodic_log_total,
336 "db_size_bytes": db_size_bytes,
337 "db_size_mb": (db_size_bytes as f64 / 1_048_576.0 * 100.0).round() / 100.0,
338 },
339 "sparks": sparks,
340 "episodic_log": {"open": open_logs, "new": new_logs},
341 "embed_rebuild_queue": embed_rebuild,
342 "knowledge_debt_ratio": (debt_ratio * 100.0).round() / 100.0,
343 "stale_screening_count": stale_screening,
344 "feedback_loop": {
345 "trace_completion_rate": ratio(trace_completed, trace_total),
346 "usage_annotation_rate": ratio(usage_known, trace_completed),
347 "trace_use_rate": ratio(usage_some, usage_known),
348 "selected_to_used_rate": ratio(selected_used, selected_total),
349 "task_success_rate": ratio(succeeded, outcome_known),
350 "feedback_coverage": ratio(feedback_traces, trace_completed),
351 "feedback_events": feedback_count,
352 "timed_out_traces": trace_timed_out,
353 "pending_evolve_requests": pending_evolve,
354 "failed_evolve_requests_30d": failed_evolve,
355 "failed_distill_logs_30d": failed_distill,
356 "pending_governance_proposals": governance_pending,
357 "window_days": 30,
358 "confidence_distribution": {
359 "low": confidence_row.and_then(|row| row.get("low")).and_then(Value::as_i64).unwrap_or(0),
360 "medium": confidence_row.and_then(|row| row.get("medium")).and_then(Value::as_i64).unwrap_or(0),
361 "high": confidence_row.and_then(|row| row.get("high")).and_then(Value::as_i64).unwrap_or(0),
362 }
363 },
364 "intuition_calibration": intuition,
365 "distill_cost_estimate": {"prompt_tokens": prompt_tokens, "completion_tokens": completion_tokens},
366 "recurring_sparks": recurring_sparks.len(),
367 "recurring_spark_ids": recurring_spark_ids,
368 "params": {
369 "recall.w_content": self.w_content,
370 "recall.w_trigger": self.w_trigger,
371 "recall.w_context": self.w_context,
372 "recall.w_activation": self.w_activation,
373 "recall.w_spread": self.w_spread,
374 "recall.top_k_candidates": self.top_k_candidates,
375 "curate.low_conf_threshold": self.low_conf_threshold,
376 "curate.low_conf_idle_days": self.low_conf_idle_days,
377 "curate.repeat_select_min": self.repeat_select_min,
378 "curate.never_used_age_days": self.never_used_age_days,
379 "curate.promote_used_success_min": self.promote_used_success_min,
380 "curate.promote_confidence_min": self.promote_confidence_min,
381 "curate.screening_timeout_minutes": self.screening_timeout_minutes,
382 "curate.open_ttl_days": self.open_ttl_days,
383 "curate.log_compact_days": self.log_compact_days,
384 "evolve.schedule_interval_hours": self.evolve_schedule_interval_hours,
385 },
386 "suggestions": suggestions
387 });
388
389 if let Some(obj) = out.as_object_mut() {
390 obj.insert("observability".to_string(), observability);
391 obj.insert("operational".to_string(), operational);
392 if let Some(trends) = trends {
393 obj.insert("trends".to_string(), trends);
394 }
395 }
396 Ok(out)
397 }
398
399 fn observability_block(&self, now: &str) -> Result<Value> {
406 let windows = json!({
407 "1d": self.window_rates(&days_ago(now, 1))?,
408 "7d": self.window_rates(&days_ago(now, 7))?,
409 "30d": self.window_rates(&days_ago(now, 30))?,
410 });
411
412 let mut by_source = serde_json::Map::new();
415 let rows = self.storage.query_chunks_params(
416 "SELECT COALESCE(event_source,'unknown') AS src, COUNT(*) AS total,
417 SUM(CASE WHEN agent IS NOT NULL THEN 1 ELSE 0 END) AS with_agent
418 FROM episodic_log WHERE ts >= ?1 GROUP BY event_source",
419 rusqlite::params![days_ago(now, 30)],
420 )?;
421 for r in &rows {
422 let src = r
423 .get("src")
424 .and_then(Value::as_str)
425 .unwrap_or("unknown")
426 .to_string();
427 let total = r.get("total").and_then(Value::as_i64).unwrap_or(0);
428 let with_agent = r.get("with_agent").and_then(Value::as_i64).unwrap_or(0);
429 by_source.insert(
430 src,
431 json!({
432 "total": total,
433 "with_agent": with_agent,
434 "agent_coverage": ratio(with_agent, total),
435 }),
436 );
437 }
438
439 let cut7 = days_ago(now, 7);
442 let by_dimension = json!({
443 "event_source": { "agent_coverage": by_source,
444 "rates": self.dimension_rates("event_source", &cut7, 0)? },
445 "agent": self.dimension_rates("agent", &cut7, 0)?,
446 "context_key": self.dimension_rates("context_key", &cut7, 10)?,
447 });
448
449 let zombie: i64 = count_query_params(
451 &self.storage,
452 "SELECT COUNT(*) FROM chunks
453 WHERE state!='archived' AND origin!='spark'
454 AND selected_count >= ?1 AND used_count = 0",
455 rusqlite::params![self.repeat_select_min],
456 )?;
457 let avg_pack = self.storage.query_chunks_params(
459 "SELECT
460 AVG(r) AS avg_retrieved,
461 AVG(s) AS avg_selected
462 FROM (
463 SELECT trace_id,
464 SUM(CASE WHEN event='retrieved' THEN 1 ELSE 0 END) AS r,
465 SUM(CASE WHEN event='selected' THEN 1 ELSE 0 END) AS s
466 FROM usage_trace WHERE ts >= ?1 GROUP BY trace_id)",
467 rusqlite::params![cut7],
468 )?;
469 let avg_get = |k: &str| {
470 avg_pack
471 .first()
472 .and_then(|r| r.get(k))
473 .and_then(Value::as_f64)
474 .map(|x| (x * 100.0).round() / 100.0)
475 .unwrap_or(0.0)
476 };
477 let selected_unused_top = self.storage.query_chunks(
479 "SELECT id, selected_count FROM chunks
480 WHERE state!='archived' AND origin!='spark' AND used_count=0 AND selected_count>0
481 ORDER BY selected_count DESC LIMIT 5",
482 )?;
483 let sel_any: i64 = count_query(
485 &self.storage,
486 "SELECT COUNT(*) FROM chunks WHERE origin!='spark' AND selected_count>0",
487 )?;
488 let sel_unused: i64 = count_query(
489 &self.storage,
490 "SELECT COUNT(*) FROM chunks WHERE origin!='spark' AND selected_count>0 AND used_count=0",
491 )?;
492 let mrr_rows = self.storage.query_chunks_params(
494 "SELECT AVG(1.0/ut.rank) AS mrr FROM usage_trace ut
495 JOIN chunks c ON c.id = ut.chunk_id
496 WHERE ut.event='selected' AND ut.rank IS NOT NULL AND ut.rank>0
497 AND c.used_count>0 AND ut.ts >= ?1",
498 rusqlite::params![cut7],
499 )?;
500 let used_rank_mrr = mrr_rows
501 .first()
502 .and_then(|r| r.get("mrr"))
503 .and_then(Value::as_f64)
504 .map(|x| (x * 1000.0).round() / 1000.0)
505 .unwrap_or(0.0);
506 let hook_total: i64 = count_query_params(
508 &self.storage,
509 "SELECT COUNT(*) FROM episodic_log WHERE event_source='hook' AND ts >= ?1",
510 rusqlite::params![cut7],
511 )?;
512 let hook_silent: i64 = count_query_params(
513 &self.storage,
514 "SELECT COUNT(*) FROM episodic_log
515 WHERE event_source='hook' AND usage_state='known_none' AND ts >= ?1",
516 rusqlite::params![cut7],
517 )?;
518 let rank_hist = self.storage.query_chunks_params(
520 "SELECT
521 SUM(CASE WHEN rank=1 THEN 1 ELSE 0 END) AS r1,
522 SUM(CASE WHEN rank BETWEEN 2 AND 3 THEN 1 ELSE 0 END) AS r2_3,
523 SUM(CASE WHEN rank BETWEEN 4 AND 10 THEN 1 ELSE 0 END) AS r4_10,
524 SUM(CASE WHEN rank > 10 THEN 1 ELSE 0 END) AS r11plus
525 FROM usage_trace WHERE event='selected' AND rank IS NOT NULL AND ts >= ?1",
526 rusqlite::params![cut7],
527 )?;
528 let rh = |k: &str| {
529 rank_hist
530 .first()
531 .and_then(|r| r.get(k))
532 .and_then(Value::as_i64)
533 .unwrap_or(0)
534 };
535 let high_rank_unused = self.storage.query_chunks_params(
537 "SELECT ut.chunk_id AS id, MIN(ut.rank) AS best_rank
538 FROM usage_trace ut JOIN chunks c ON c.id = ut.chunk_id
539 WHERE ut.event='selected' AND ut.rank<=3 AND c.used_count=0 AND ut.ts >= ?1
540 GROUP BY ut.chunk_id ORDER BY best_rank LIMIT 5",
541 rusqlite::params![cut7],
542 )?;
543 let low_rank_used = self.storage.query_chunks_params(
546 "SELECT ut.chunk_id AS id, MIN(ut.rank) AS best_rank
547 FROM usage_trace ut JOIN chunks c ON c.id = ut.chunk_id
548 WHERE ut.event='selected' AND ut.rank>10 AND c.used_count>0 AND ut.ts >= ?1
549 GROUP BY ut.chunk_id ORDER BY best_rank DESC LIMIT 5",
550 rusqlite::params![cut7],
551 )?;
552
553 let promotions: i64 = count_query_params(
555 &self.storage,
556 "SELECT COUNT(*) FROM chunks
557 WHERE state='active' AND origin!='spark' AND state_updated_at >= ?1
558 AND state_reason IN ('repeated_success','approved','restore')",
559 rusqlite::params![cut7],
560 )?;
561 let evictions: i64 = count_query_params(
562 &self.storage,
563 "SELECT COUNT(*) FROM chunks
564 WHERE state='archived' AND origin!='spark' AND state_updated_at >= ?1",
565 rusqlite::params![cut7],
566 )?;
567 let pending_oldest = self.storage.query_chunks(
568 "SELECT MIN(created_at) AS oldest FROM chunks WHERE state='pending' AND origin!='spark'",
569 )?;
570 let pending_oldest_ts = pending_oldest
571 .first()
572 .and_then(|r| r.get("oldest"))
573 .and_then(Value::as_str)
574 .map(|s| s.to_string());
575 let gov_backlog = self.storage.query_chunks(
576 "SELECT MIN(created_at) AS oldest FROM governance_proposals WHERE state='pending'",
577 )?;
578 let gov_backlog_oldest_ts = gov_backlog
579 .first()
580 .and_then(|r| r.get("oldest"))
581 .and_then(Value::as_str)
582 .map(|s| s.to_string());
583
584 Ok(json!({
585 "windows": windows,
586 "by_dimension": by_dimension,
587 "recall_pack": {
588 "zombie_chunks": zombie,
589 "avg_retrieved": avg_get("avg_retrieved"),
590 "avg_selected": avg_get("avg_selected"),
591 "selected_unused_rate": ratio(sel_unused, sel_any),
592 "selected_unused_top": selected_unused_top,
593 "used_rank_mrr": used_rank_mrr,
594 "hook_silence_rate": ratio(hook_silent, hook_total),
595 "selected_rank_distribution": {
596 "1": rh("r1"), "2-3": rh("r2_3"), "4-10": rh("r4_10"), "11+": rh("r11plus")
597 },
598 "high_rank_unused": high_rank_unused,
599 "low_rank_used": low_rank_used,
600 },
601 "lifecycle": {
602 "pending_oldest_ts": pending_oldest_ts,
603 "governance_backlog_oldest_ts": gov_backlog_oldest_ts,
604 "state_transition_approx": {
605 "promotions_7d": promotions,
606 "evictions_7d": evictions,
607 "note": "approx via state_updated_at/state_reason; not a strict rate"
608 }
609 }
610 }))
611 }
612
613 fn dimension_rates(&self, group_col: &str, cutoff: &str, top_n: i64) -> Result<Value> {
617 let limit = if top_n > 0 {
618 format!(" ORDER BY recalls DESC LIMIT {top_n}")
619 } else {
620 String::new()
621 };
622 let sql = format!(
623 "SELECT COALESCE({group_col},'(null)') AS g,
624 COUNT(*) AS recalls,
625 SUM(CASE WHEN usage_state='known_none' THEN 1 ELSE 0 END) AS empty,
626 SUM(CASE WHEN task_state='completed' THEN 1 ELSE 0 END) AS completed,
627 SUM(CASE WHEN outcome='ok' THEN 1 ELSE 0 END) AS ok,
628 SUM(CASE WHEN outcome IN ('ok','fail') THEN 1 ELSE 0 END) AS outcome_known,
629 SUM(CASE WHEN task_state='completed' AND usage_state!='unknown' THEN 1 ELSE 0 END) AS annotated
630 FROM episodic_log WHERE ts >= ?1 GROUP BY {group_col}{limit}"
631 );
632 let rows = self.storage.query_chunks_params(&sql, rusqlite::params![cutoff])?;
633
634 let su_sql = format!(
636 "SELECT COALESCE(el.{group_col},'(null)') AS g,
637 SUM(CASE WHEN ut.event='selected' THEN 1 ELSE 0 END) AS sel,
638 SUM(CASE WHEN ut.event='used' THEN 1 ELSE 0 END) AS used
639 FROM usage_trace ut JOIN episodic_log el ON el.trace_id = ut.trace_id
640 WHERE ut.ts >= ?1 GROUP BY el.{group_col}"
641 );
642 let su_rows = self.storage.query_chunks_params(&su_sql, rusqlite::params![cutoff])?;
643 let mut su_map: std::collections::HashMap<String, (i64, i64)> = std::collections::HashMap::new();
644 for r in &su_rows {
645 let g = r.get("g").and_then(Value::as_str).unwrap_or("(null)").to_string();
646 su_map.insert(
647 g,
648 (
649 r.get("sel").and_then(Value::as_i64).unwrap_or(0),
650 r.get("used").and_then(Value::as_i64).unwrap_or(0),
651 ),
652 );
653 }
654 let fb_sql = format!(
656 "SELECT COALESCE(el.{group_col},'(null)') AS g,
657 COUNT(DISTINCT fe.trace_id) AS fb
658 FROM feedback_events fe JOIN episodic_log el ON el.trace_id = fe.trace_id
659 WHERE fe.ts >= ?1 GROUP BY el.{group_col}"
660 );
661 let fb_rows = self.storage.query_chunks_params(&fb_sql, rusqlite::params![cutoff])?;
662 let mut fb_map: std::collections::HashMap<String, i64> = std::collections::HashMap::new();
663 for r in &fb_rows {
664 let g = r.get("g").and_then(Value::as_str).unwrap_or("(null)").to_string();
665 fb_map.insert(g, r.get("fb").and_then(Value::as_i64).unwrap_or(0));
666 }
667
668 let mut out = serde_json::Map::new();
669 for r in &rows {
670 let g = r.get("g").and_then(Value::as_str).unwrap_or("(null)").to_string();
671 let i = |k: &str| r.get(k).and_then(Value::as_i64).unwrap_or(0);
672 let recalls = i("recalls");
673 let completed = i("completed");
674 let (sel, used) = su_map.get(&g).copied().unwrap_or((0, 0));
675 let fb = fb_map.get(&g).copied().unwrap_or(0);
676 out.insert(
677 g,
678 json!({
679 "recalls": recalls,
680 "empty_recall_rate": ratio(i("empty"), recalls),
681 "completed_rate": ratio(completed, recalls),
682 "task_success_rate": ratio(i("ok"), i("outcome_known")),
683 "usage_annotation_rate": ratio(i("annotated"), completed),
684 "selected_to_used_rate": ratio(used, sel),
685 "feedback_coverage": ratio(fb, completed),
686 }),
687 );
688 }
689 Ok(Value::Object(out))
690 }
691
692 fn window_rates(&self, cutoff: &str) -> Result<Value> {
696 let rows = self.storage.query_chunks_params(
697 "SELECT COUNT(*) AS recalls,
698 SUM(CASE WHEN usage_state='known_none' THEN 1 ELSE 0 END) AS empty,
699 SUM(CASE WHEN task_state='completed' THEN 1 ELSE 0 END) AS completed,
700 SUM(CASE WHEN task_state='timed_out' THEN 1 ELSE 0 END) AS timed_out,
701 SUM(CASE WHEN outcome='ok' THEN 1 ELSE 0 END) AS ok,
702 SUM(CASE WHEN outcome IN ('ok','fail') THEN 1 ELSE 0 END) AS outcome_known,
703 SUM(CASE WHEN task_state='completed' AND usage_state!='unknown' THEN 1 ELSE 0 END) AS annotated
704 FROM episodic_log WHERE ts >= ?1",
705 rusqlite::params![cutoff],
706 )?;
707 let g = |k: &str| {
708 rows.first()
709 .and_then(|r| r.get(k))
710 .and_then(Value::as_i64)
711 .unwrap_or(0)
712 };
713 let su = self.storage.query_chunks_params(
715 "SELECT SUM(CASE WHEN event='selected' THEN 1 ELSE 0 END) AS sel,
716 SUM(CASE WHEN event='used' THEN 1 ELSE 0 END) AS used
717 FROM usage_trace WHERE ts >= ?1",
718 rusqlite::params![cutoff],
719 )?;
720 let sug = |k: &str| {
721 su.first()
722 .and_then(|r| r.get(k))
723 .and_then(Value::as_i64)
724 .unwrap_or(0)
725 };
726 let fb_traces: i64 = count_query_params(
727 &self.storage,
728 "SELECT COUNT(DISTINCT trace_id) FROM feedback_events WHERE ts >= ?1",
729 rusqlite::params![cutoff],
730 )?;
731 let recalls = g("recalls");
732 let completed = g("completed");
733 Ok(json!({
734 "recalls": recalls,
735 "empty_recall_rate": ratio(g("empty"), recalls),
736 "completed_rate": ratio(completed, recalls),
737 "timeout_rate": ratio(g("timed_out"), recalls),
738 "task_success_rate": ratio(g("ok"), g("outcome_known")),
739 "usage_annotation_rate": ratio(g("annotated"), completed),
740 "selected_to_used_rate": ratio(sug("used"), sug("sel")),
741 "feedback_coverage": ratio(fb_traces, completed),
742 }))
743 }
744
745 fn operational_block(&self, now: &str) -> Result<Value> {
747 let daemon = crate::daemon::health(
748 &crate::paths::daemon_state_path(),
749 &crate::paths::daemon_pid_path(),
750 now,
751 );
752 let mut block = serde_json::Map::new();
753 block.insert("daemon".to_string(), daemon);
754 if self.storage.count_operation_runs().unwrap_or(0) > 0 {
757 let rows = self.storage.operation_runs_since(&days_ago(now, 7))?;
758 block.insert(
759 "ops".to_string(),
760 crate::storage::metrics::aggregate_ops(&rows),
761 );
762 }
763 Ok(Value::Object(block))
764 }
765
766 fn trends_block(&self, now: &str) -> Result<Option<Value>> {
769 let Some((cur_ts, cur_json)) = self.storage.latest_snapshot()? else {
770 return Ok(None);
771 };
772 let cur: Value = serde_json::from_str(&cur_json).unwrap_or_else(|_| json!({}));
773 let mut out = json!({ "current_ts": cur_ts, "current": cur.clone() });
774 if let Some((base_ts, base_json)) =
775 self.storage.snapshot_at_or_before(&days_ago(now, 7))?
776 {
777 let base: Value = serde_json::from_str(&base_json).unwrap_or_else(|_| json!({}));
778 let mut delta = serde_json::Map::new();
779 if let (Some(c), Some(b)) = (cur.as_object(), base.as_object()) {
780 for (k, cv) in c {
781 if let (Some(cf), Some(bf)) = (cv.as_f64(), b.get(k).and_then(Value::as_f64)) {
782 delta.insert(k.clone(), json!(((cf - bf) * 1000.0).round() / 1000.0));
783 }
784 }
785 }
786 out["baseline_ts"] = json!(base_ts);
787 out["delta_vs_7d"] = Value::Object(delta);
788 }
789 Ok(Some(out))
790 }
791
792 pub fn recall_weights(&self) -> Value {
795 json!({
796 "w_content": self.w_content,
797 "w_trigger": self.w_trigger,
798 "w_lexical": self.w_lexical,
799 "w_context": self.w_context,
800 "w_activation": self.w_activation,
801 "w_spread": self.w_spread,
802 })
803 }
804
805 pub fn write_metric_snapshot(&self) -> Result<Value> {
807 let now = utc_now_iso();
808 let kpis = self.collect_kpis(&now)?;
809 self.storage.insert_metric_snapshot(&now, &kpis.to_string())?;
810 Ok(kpis)
811 }
812
813 pub(crate) fn collect_kpis(&self, now: &str) -> Result<Value> {
816 let active: i64 = count_query(
817 &self.storage,
818 "SELECT COUNT(*) FROM chunks WHERE state='active' AND origin!='spark'",
819 )?;
820 let pending: i64 = count_query(
821 &self.storage,
822 "SELECT COUNT(*) FROM chunks WHERE state='pending' AND origin!='spark'",
823 )?;
824 let zombie: i64 = count_query_params(
825 &self.storage,
826 "SELECT COUNT(*) FROM chunks
827 WHERE state!='archived' AND origin!='spark'
828 AND selected_count >= ?1 AND used_count = 0",
829 rusqlite::params![self.repeat_select_min],
830 )?;
831 let debt_ratio = if active > 0 {
832 (pending as f64 / active as f64 * 100.0).round() / 100.0
833 } else {
834 pending as f64
835 };
836 let w7 = self.window_rates(&days_ago(now, 7))?;
837 let daemon = crate::daemon::health(
838 &crate::paths::daemon_state_path(),
839 &crate::paths::daemon_pid_path(),
840 now,
841 );
842 Ok(json!({
843 "active": active,
844 "pending": pending,
845 "knowledge_debt_ratio": debt_ratio,
846 "zombie_chunks": zombie,
847 "task_success_rate_7d": w7.get("task_success_rate").and_then(Value::as_f64).unwrap_or(0.0),
848 "empty_recall_rate_7d": w7.get("empty_recall_rate").and_then(Value::as_f64).unwrap_or(0.0),
849 "daemon_errors_24h": daemon.get("errors_24h").and_then(Value::as_i64).unwrap_or(0),
850 }))
851 }
852
853 fn intuition_calibration(&self, window_start: &str) -> Result<Value> {
863 let rows = self.storage.query_chunks_params(
864 "SELECT recall_snapshot, outcome FROM episodic_log
865 WHERE ts >= ? AND recall_snapshot LIKE '%\"appraise\"%'",
866 rusqlite::params![window_start],
867 )?;
868
869 let mut buckets: std::collections::BTreeMap<String, [f64; 4]> =
871 std::collections::BTreeMap::new();
872 for tier in ["weak", "medium", "strong"] {
873 buckets.insert(tier.to_string(), [0.0; 4]);
874 }
875 let mut total = 0_i64;
876 let mut silent = 0_i64;
877 let mut caution_strong = 0_i64;
878 let mut caution_strong_false = 0_i64;
879
880 for row in &rows {
881 let snapshot = row
882 .get("recall_snapshot")
883 .and_then(Value::as_str)
884 .and_then(|raw| serde_json::from_str::<Value>(raw).ok());
885 let Some(appraise) = snapshot.as_ref().and_then(|s| s.get("appraise")) else {
886 continue;
887 };
888 let tier = appraise
889 .get("tier")
890 .and_then(Value::as_str)
891 .unwrap_or("weak");
892 let valence = appraise
893 .get("valence")
894 .and_then(Value::as_str)
895 .unwrap_or("neutral");
896 let strength = appraise
897 .get("strength")
898 .and_then(Value::as_f64)
899 .unwrap_or(0.0);
900 let outcome = row.get("outcome").and_then(Value::as_str);
901
902 total += 1;
903 if tier == "weak" || valence == "neutral" {
904 silent += 1;
905 }
906 let has_outcome = matches!(outcome, Some("ok") | Some("fail"));
907 let is_ok = outcome == Some("ok");
908 if let Some(b) = buckets.get_mut(tier) {
909 b[0] += 1.0;
910 if has_outcome {
911 b[1] += 1.0;
912 b[3] += strength;
913 if is_ok {
914 b[2] += 1.0;
915 }
916 }
917 }
918 if valence == "caution" && tier == "strong" && has_outcome {
919 caution_strong += 1;
920 if is_ok {
921 caution_strong_false += 1;
922 }
923 }
924 }
925
926 let hit_rate = |b: &[f64; 4]| if b[1] > 0.0 { b[2] / b[1] } else { 0.0 };
927 let weak = buckets.get("weak").copied().unwrap_or([0.0; 4]);
928 let strong = buckets.get("strong").copied().unwrap_or([0.0; 4]);
929 let monotonicity_gap = hit_rate(&strong) - hit_rate(&weak);
930
931 let outcome_total: f64 = buckets.values().map(|b| b[1]).sum();
933 let ece = if outcome_total > 0.0 {
934 buckets
935 .values()
936 .filter(|b| b[1] > 0.0)
937 .map(|b| {
938 let avg_strength = b[3] / b[1];
939 (b[1] / outcome_total) * (avg_strength - hit_rate(b)).abs()
940 })
941 .sum::<f64>()
942 } else {
943 0.0
944 };
945
946 let bucket_detail: Vec<Value> = ["weak", "medium", "strong"]
947 .iter()
948 .map(|tier| {
949 let b = buckets.get(*tier).copied().unwrap_or([0.0; 4]);
950 json!({
951 "tier": tier,
952 "n": b[0] as i64,
953 "n_with_outcome": b[1] as i64,
954 "avg_strength": if b[1] > 0.0 { (b[3] / b[1] * 1000.0).round() / 1000.0 } else { 0.0 },
955 "actual_hit_rate": (hit_rate(&b) * 1000.0).round() / 1000.0,
956 })
957 })
958 .collect();
959
960 let (vl_total, vl_abstained, vl_observed) =
964 self.storage.verdict_log_overview().unwrap_or((0, 0, 0));
965 let samples = self.storage.verdict_calibration_samples().unwrap_or_default();
966 let bins = self.calibration_bins.max(2);
967 let mut vhit = vec![0.0_f64; bins as usize];
968 let mut vtot = vec![0.0_f64; bins as usize];
969 for (_strength, conf, h) in &samples {
972 let b = ((conf * bins as f64).floor() as i64).clamp(0, bins - 1) as usize;
973 vtot[b] += 1.0;
974 vhit[b] += *h;
975 }
976 let n_obs: f64 = vtot.iter().sum();
977 let verdict_ece = if n_obs > 0.0 {
978 (0..bins as usize)
979 .filter(|&b| vtot[b] > 0.0)
980 .map(|b| {
981 let claimed = (b as f64 + 0.5) / bins as f64;
982 let actual = vhit[b] / vtot[b];
983 (vtot[b] / n_obs) * (claimed - actual).abs()
984 })
985 .sum::<f64>()
986 } else {
987 0.0
988 };
989
990 Ok(json!({
991 "appraisals": total,
992 "monotonicity_gap": (monotonicity_gap * 1000.0).round() / 1000.0,
993 "ece": (ece * 1000.0).round() / 1000.0,
994 "false_alarm_rate": ratio(caution_strong_false, caution_strong),
995 "silence_rate": ratio(silent, total),
996 "buckets": bucket_detail,
997 "verdict_log": {
999 "total": vl_total,
1000 "abstained": vl_abstained,
1001 "abstain_rate": ratio(vl_abstained, vl_total),
1002 "observed": vl_observed,
1003 "ece": (verdict_ece * 1000.0).round() / 1000.0,
1004 },
1005 }))
1006 }
1007
1008 pub fn rebuild_embeddings(&self) -> Result<usize> {
1013 Ok(self.rebuild_embeddings_capped(None)?.0)
1014 }
1015
1016 pub fn rebuild_embeddings_capped(&self, max: Option<usize>) -> Result<(usize, usize)> {
1023 let meta_version = self
1024 .storage
1025 .get_meta("embed_version")?
1026 .and_then(|v| v.parse::<i64>().ok())
1027 .unwrap_or(1);
1028 let mut stale = self.storage.query_chunks_params(
1030 "SELECT id, content, trigger_desc, state_reason FROM chunks
1031 WHERE embed_version = 0 OR embed_version < ?",
1032 rusqlite::params![meta_version],
1033 )?;
1034 let total_stale = stale.len();
1036 let remaining = match max {
1037 Some(n) if n < total_stale => {
1038 stale.truncate(n);
1039 total_stale - n
1040 }
1041 _ => 0,
1042 };
1043 self.storage.invalidate_vector_caches();
1047 let mut count = 0;
1048 for row in &stale {
1049 let id = match row.get("id").and_then(Value::as_str) {
1050 Some(v) => v,
1051 None => continue,
1052 };
1053 let content = row.get("content").and_then(Value::as_str).unwrap_or("");
1054 let trigger = row
1055 .get("trigger_desc")
1056 .and_then(Value::as_str)
1057 .unwrap_or(content);
1058 let state_reason = row
1059 .get("state_reason")
1060 .and_then(Value::as_str)
1061 .unwrap_or("");
1062
1063 let (cvec_res, tvec_res) = self.embed_pair(content, trigger, "rebuild");
1064 let cvec = match cvec_res {
1065 Ok(v) => v,
1066 Err(_) => continue,
1067 };
1068 let tvec = match tvec_res {
1069 Ok(v) => v,
1070 Err(_) => continue,
1071 };
1072
1073 self.storage.begin_immediate()?;
1074 let r = (|| -> Result<()> {
1075 self.store_vec_content(id, &cvec)?;
1076 self.store_vec_trigger(id, &tvec)?;
1077 let new_reason = if state_reason.starts_with("embedding_pending:target=") {
1079 let target_state = state_reason.trim_start_matches("embedding_pending:target=");
1080 let now = utc_now_iso();
1081 self.storage.update_chunk_state(
1082 id,
1083 target_state,
1084 Some("embedding_rebuilt"),
1085 &now,
1086 )?;
1087 "embedding_rebuilt".to_string()
1088 } else {
1089 "embedding_rebuilt".to_string()
1090 };
1091 let now = utc_now_iso();
1092 self.storage.conn_execute(
1093 "UPDATE chunks SET embed_version=?, state_reason=?, updated_at=? WHERE id=?",
1094 rusqlite::params![meta_version, new_reason, now, id],
1095 )?;
1096 self.storage.commit()
1097 })();
1098 if r.is_err() {
1099 let _ = self.storage.rollback();
1100 } else {
1101 count += 1;
1102 }
1103 }
1104 Ok((count, remaining))
1105 }
1106
1107 pub fn inspect_id(&self, id: &str) -> Result<Value> {
1112 if let Some(chunk) = self.storage.get_chunk(id)? {
1114 let traces = self.storage.query_chunks_params(
1115 "SELECT * FROM usage_trace WHERE chunk_id=? ORDER BY ts DESC LIMIT 20",
1116 rusqlite::params![id],
1117 )?;
1118 let derived = self.storage.query_chunks_params(
1119 "SELECT id, state, confidence FROM chunks WHERE distilled_from IN (
1120 SELECT id FROM episodic_log WHERE trace_id IN (
1121 SELECT trace_id FROM usage_trace WHERE chunk_id=?
1122 )
1123 ) LIMIT 10",
1124 rusqlite::params![id],
1125 )?;
1126 return Ok(json!({
1127 "kind": "chunk",
1128 "chunk": chunk,
1129 "recent_traces": traces,
1130 "derived_chunks": derived,
1131 }));
1132 }
1133 if let Some(log) = self.storage.get_episodic_log(id)? {
1135 let traces = self.storage.query_chunks_params(
1136 "SELECT * FROM usage_trace WHERE trace_id=? ORDER BY ts ASC",
1137 rusqlite::params![id],
1138 )?;
1139 return Ok(json!({
1140 "kind": "trace",
1141 "episodic_log": log,
1142 "usage_traces": traces,
1143 }));
1144 }
1145 Err(InnateError::ChunkNotFound(id.to_string()))
1146 }
1147
1148 pub(super) fn sanitize_content(&self, content: &str) -> (String, SanitizeAction) {
1153 self.sanitizer.sanitize(content)
1154 }
1155}