Skip to main content

innate_core/kb/
inspection.rs

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        // P3-A: oldest pending chunk timestamp — surfaces long-lived pending debt.
158        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        // Health signal 1: knowledge debt ratio.
165        // Zombie = active chunks with middling confidence (stuck, neither good nor bad)
166        // that are at least 14d old and have been used at least once.
167        // "never-recalled old" chunks are handled by curate 3c (never_used archive).
168        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        // Health signal 5: stale screening count
183        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        // Health signal 4: actual Distill cost within the configured rolling window.
192        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        // Health signal 2: sparks that have been recalled often (soft incubation threshold = 5)
212        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        // agent_coverage misconfiguration (design doc §5.1 / §7): hook/daemon NULL agent
265        // is *by design*, so only an mcp/cli channel with NULL agent is a real signal —
266        // typically a missing `INNATE_AGENT` env in the agent's MCP config. Fire only on
267        // that channel to avoid the perma-misfire a global coverage check would cause.
268        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        // Intuition honesty (PRD §4): does high strength actually predict success, and is
282        // the critic crying wolf? Only nudge once enough appraisals carry an outcome.
283        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        // Storage growth metrics — trace/log bloat is driven by recall/record
310        // activity over time, independent of chunk count, so it is surfaced here
311        // for monitoring before it becomes a problem.
312        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        // ── Observability (P1–P3b) — additive top-level blocks, see design doc §5.6.
319        // Existing keys are never renamed/moved;未就绪的子块整块省略。
320        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    // ------------------------------------------------------------------
400    // Observability blocks (design doc §5.1 / §5.6)
401    // ------------------------------------------------------------------
402
403    /// P1 derived metrics: multi-window rates, per-channel agent coverage,
404    /// zombie chunks, and approximate lifecycle transitions. Pure SQL, no new tables.
405    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        // agent_coverage **by channel** — hook/daemon NULL agent is by design (see §5.1);
413        // only mcp/cli NULL is a real misconfiguration signal.
414        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        // Per-dimension 7d rate breakdown (event_source / agent / top context_key) —
440        // answers "which source/agent/context is degrading" (design doc §5.1).
441        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        // ── recall_pack (§5.1) ──
450        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        // avg retrieved/selected per recall trace (recent usage_trace window).
458        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        // selected-but-never-used top offenders (zombie detail).
478        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        // selected→used (chunk-level): fraction of ever-selected chunks never used.
484        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        // MRR proxy: mean reciprocal selected-rank of chunks that were actually used.
493        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        // hook silence rate (§5.5): hook recalls returning known_none / all hook recalls.
507        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        // selected rank distribution (§5.1): where in the shortlist selected chunks landed.
519        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        // high-rank-unused anomaly: took a top-3 slot but the chunk was never used.
536        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        // low-rank-used anomaly: useful chunk that only surfaced deep (rank>10) — ranking
544        // could promote it. Sample for debugging fused-score weights.
545        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        // ── lifecycle (§5.1) ──
554        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    /// Per-group rate breakdown over `episodic_log` newer than `cutoff`, grouped by
614    /// `group_col` (event_source / agent / context_key). `top_n>0` caps to the busiest
615    /// N groups (for high-cardinality context_key). NULL group → "(null)".
616    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        // selected→used per group: join usage_trace back to episodic_log for the dimension.
635        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        // feedback coverage per group.
655        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    /// Full rate metrics over `episodic_log` rows newer than `cutoff` (§5.1). The
693    /// selected→used and feedback-coverage signals come from usage_trace/feedback_events
694    /// (two cheap supplementary windowed scans), the rest from a single episodic_log scan.
695    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        // selected→used (event counts) over the same window.
714        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    /// P1 daemon health (independent read-only connection) + P3a op aggregation.
746    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        // Omit `ops` entirely until operation_runs has data (table exists from 4.20 but
755        // may be empty before instrumentation lands / on a fresh db).
756        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    /// P3b trend block: current snapshot + delta vs the nearest snapshot ≥7d old.
767    /// Returns None when no snapshot exists yet (block omitted from inspect).
768    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    /// Fused-score weights in effect, as a JSON object (for the recall-eval params
793    /// snapshot so an eval run is self-describing and comparable over time).
794    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    /// Write a state-KPI snapshot row now and return the KPIs (CLI `metrics snapshot`).
806    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    /// State-KPI snapshot payload for `metric_snapshots` (written by curate at cycle
814    /// end). Only numeric, stable KPIs — these can't be reconstructed from event logs.
815    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    // ------------------------------------------------------------------
854    // Intuition honesty metrics (PRD §4 / Spec §7)
855    //
856    // The core KPI is not recall but discrimination quality: "loud when it should
857    // be, silent when it shouldn't." All inputs already exist — appraise persists
858    // {valence, tier, strength} into episodic_log.recall_snapshot, and record fills
859    // in `outcome`. We bucket appraisals by tier and check the actual task_ok rate.
860    // ------------------------------------------------------------------
861
862    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        // Per-tier accumulators: (n_total, n_with_outcome, ok, sum_strength_with_outcome).
870        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        // ECE: evidence-weighted gap between mean strength and actual hit rate per bucket.
932        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        // 方案 B —— verdict_log 仪表盘:可证伪的 ECE / 弃权率(头号体检指标)。
961        // 与上面基于 recall_snapshot 的 tier-bucket 指标互补:verdict_log 直接用
962        // emitted_conf 分桶 + observed 回填算 ECE,且把弃权率作为一等健康信号。
963        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        // ECE 按 **emitted_conf** 分桶:衡量「声称置信度」的真实兑现率(校准映射重算
970        // 则按 strength 分桶,见 curate::recompute_calibration_map —— 两者域不同)。
971        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            // 方案 B verdict_log 仪表盘
998            "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    // ------------------------------------------------------------------
1009    // Public: rebuild_embeddings (evolve --rebuild-embeddings)
1010    // ------------------------------------------------------------------
1011
1012    pub fn rebuild_embeddings(&self) -> Result<usize> {
1013        Ok(self.rebuild_embeddings_capped(None)?.0)
1014    }
1015
1016    /// Bounded re-embed used by latency-sensitive callers (MCP/CLI `evolve
1017    /// --rebuild-embeddings`). When `max` is `Some(n)` at most `n` stale chunks
1018    /// are re-embedded this call, so a large backlog is chipped away across
1019    /// successive evolves instead of blocking one request on the whole queue
1020    /// (each re-embed is a network LLM round-trip). Returns
1021    /// `(rebuilt_this_call, remaining_stale_after)`. `None` rebuilds everything.
1022    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        // Fetch chunks with embed_version=0 (failed writes) or below current meta version.
1029        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        // Bound the batch: keep the first `max`, report the rest as remaining.
1035        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        // Bulk re-embed: drop the warm cache once so the per-row in-place upserts
1044        // stay no-ops (cold) and the loop runs O(N) instead of O(N²). The next
1045        // search reloads the rebuilt vectors from disk.
1046        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                // Restore intended state if encoded in state_reason.
1078                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    // ------------------------------------------------------------------
1108    // Public: inspect_id (inspect <chunk_id> or <trace_id>)
1109    // ------------------------------------------------------------------
1110
1111    pub fn inspect_id(&self, id: &str) -> Result<Value> {
1112        // Try as chunk_id first, then as trace_id.
1113        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        // Try as trace_id.
1134        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    // ------------------------------------------------------------------
1149    // Sanitize
1150    // ------------------------------------------------------------------
1151
1152    pub(super) fn sanitize_content(&self, content: &str) -> (String, SanitizeAction) {
1153        self.sanitizer.sanitize(content)
1154    }
1155}