Skip to main content

innate_core/kb/
evolve.rs

1use super::*;
2
3impl KnowledgeBase {
4    pub fn evolve(&self, trigger: &str) -> Result<Value> {
5        self.measure("evolve", None, None, || self.evolve_inner(trigger))
6    }
7
8    fn evolve_inner(&self, trigger: &str) -> Result<Value> {
9        if !matches!(trigger, "manual" | "scheduled" | "threshold") {
10            return Err(InnateError::InvalidState(format!(
11                "invalid evolve trigger: {trigger}"
12            )));
13        }
14        let evolve_started_at = utc_now_iso();
15        let retry_cutoff = minutes_ago(&evolve_started_at, 5);
16        let recovered_failed = self.storage.conn_execute_count(
17            "UPDATE episodic_log
18             SET distill_state='new', distill_note='retry_failed',
19                 distill_locked_at=NULL, distill_run_id=NULL
20             WHERE distill_state='failed'
21               AND distill_attempts < 3
22               AND COALESCE(distill_last_failed_at, distill_accounted_at, ts) < ?",
23            rusqlite::params![retry_cutoff],
24        )?;
25        if recovered_failed > 0 {
26            self.storage
27                .request_evolve(&gen_uuid(), "distill_retry", &evolve_started_at)?;
28        }
29        if trigger == "scheduled" {
30            let age_cutoff = hours_ago(&evolve_started_at, self.evolve_schedule_interval_hours);
31            let aged_new = count_query_params(
32                &self.storage,
33                "SELECT COUNT(*) FROM episodic_log
34                 WHERE distill_state='new' AND ts <= ?",
35                rusqlite::params![age_cutoff],
36            )?;
37            if aged_new > 0 {
38                self.storage
39                    .request_evolve(&gen_uuid(), "scheduled", &evolve_started_at)?;
40            }
41        }
42        let request = self.storage.claim_evolve_request_with_reason(
43            &evolve_started_at,
44            &minutes_ago(&evolve_started_at, self.screening_timeout_minutes),
45        )?;
46        let request_id = request.as_ref().map(|claim| claim.id.as_str());
47        let request_reason = request.as_ref().map(|claim| claim.reason.as_str());
48        // Issue 7: scheduled without pending request → still run curate (time-based maintenance).
49        if trigger == "scheduled" && request_id.is_none() {
50            let curator = Arc::clone(&self.curator);
51            let curate = curator.run(self, &CurateScope::default())?;
52            return Ok(json!({
53                "distilled": 0,
54                "curate": self.format_curate_report(&curate),
55                "skipped": "no_evolve_request"
56            }));
57        }
58
59        // Issue 8: threshold / token-limit gates should not suppress curate.
60        if trigger == "threshold" {
61            let rows = self.storage.query_chunks(
62                "SELECT COUNT(*) AS cnt FROM episodic_log WHERE distill_state='new'",
63            )?;
64            let cnt = rows
65                .first()
66                .and_then(|r| r.get("cnt"))
67                .and_then(Value::as_i64)
68                .unwrap_or(0);
69            if cnt < self.evolve_threshold {
70                let curator = Arc::clone(&self.curator);
71                let curate = curator.run(self, &CurateScope::default())?;
72                if let Some(id) = request_id {
73                    if matches!(request_reason, Some("governance" | "governance_ready")) {
74                        self.storage.finish_evolve_request(
75                            id,
76                            "completed",
77                            Some("curate_only"),
78                            &utc_now_iso(),
79                        )?;
80                    } else {
81                        self.storage.defer_evolve_request(
82                            id,
83                            "below_threshold",
84                            &hours_after(
85                                &evolve_started_at,
86                                self.evolve_schedule_interval_hours.max(1),
87                            ),
88                        )?;
89                    }
90                }
91                return Ok(json!({
92                    "distilled": 0,
93                    "curate": self.format_curate_report(&curate),
94                    "skipped": "below_threshold"
95                }));
96            }
97        }
98
99        if trigger != "manual" {
100            if let Some(limit) = self
101                .storage
102                .get_meta("max_distill_tokens_per_period")?
103                .and_then(|value| value.parse::<i64>().ok())
104                .filter(|value| *value > 0)
105            {
106                let period_start = self.distill_token_period_start(&evolve_started_at)?;
107                let rows = self.storage.query_chunks_params(
108                    "SELECT COALESCE(SUM(prompt_tokens + completion_tokens),0) AS used
109                     FROM distill_token_usage
110                     WHERE accounted_at >= ?",
111                    rusqlite::params![period_start],
112                )?;
113                let used_tokens = rows
114                    .first()
115                    .and_then(|row| row.get("used"))
116                    .and_then(Value::as_i64)
117                    .unwrap_or(0);
118                if used_tokens >= limit {
119                    let curator = Arc::clone(&self.curator);
120                    let curate = curator.run(self, &CurateScope::default())?;
121                    if let Some(id) = request_id {
122                        if matches!(request_reason, Some("governance" | "governance_ready")) {
123                            self.storage.finish_evolve_request(
124                                id,
125                                "completed",
126                                Some("curate_only"),
127                                &utc_now_iso(),
128                            )?;
129                        } else {
130                            self.storage.defer_evolve_request(
131                                id,
132                                "distill_token_limit",
133                                &hours_after(&evolve_started_at, 1),
134                            )?;
135                        }
136                    }
137                    return Ok(json!({
138                        "distilled": 0,
139                        "curate": self.format_curate_report(&curate),
140                        "skipped": "distill_token_limit",
141                        "distill_tokens_used": used_tokens,
142                        "distill_token_limit": limit,
143                        "period_start": period_start,
144                    }));
145                }
146            }
147        }
148
149        let result = (|| -> Result<Value> {
150            let distill = self.measure("distill", None, None, || self.distill_batch())?;
151            let curator = Arc::clone(&self.curator);
152            let curate = curator.run(self, &CurateScope::default())?;
153            Ok(json!({
154                "distilled": distill.distilled,
155                "distill_failed": distill.failed,
156                "curate": self.format_curate_report(&curate),
157            }))
158        })();
159        if let Some(id) = request_id {
160            let (state, note) = match &result {
161                Ok(_) => ("completed", None),
162                Err(error) => ("failed", Some(error.to_string())),
163            };
164            self.storage
165                .finish_evolve_request(id, state, note.as_deref(), &utc_now_iso())?;
166        }
167        if result.is_ok() {
168            self.storage
169                .finish_covered_evolve_requests(&evolve_started_at, &utc_now_iso())?;
170        }
171
172        let failed_remaining = count_query(
173            &self.storage,
174            "SELECT COUNT(*) FROM episodic_log
175             WHERE distill_state='failed' AND distill_attempts < 3",
176        )?;
177        if failed_remaining > 0 {
178            self.storage.request_evolve_at(
179                &gen_uuid(),
180                "distill_retry",
181                &utc_now_iso(),
182                Some(&minutes_after(&utc_now_iso(), 5)),
183            )?;
184        }
185
186        // Issue 9: if more 'new' logs remain after this batch, self-queue so the next evolve drains them.
187        if result.is_ok() {
188            let remaining = count_query(
189                &self.storage,
190                "SELECT COUNT(*) FROM episodic_log WHERE distill_state='new'",
191            )?;
192            if remaining > 0 {
193                let _ = self
194                    .storage
195                    .request_evolve(&gen_uuid(), "batch_continue", &utc_now_iso());
196            }
197        }
198        result
199    }
200
201    fn format_curate_report(&self, curate: &CurateReport) -> Value {
202        json!({
203            "archived": curate.archived.len(),
204            "promoted": curate.promoted.len(),
205            "deduped": curate.deduped.len(),
206            "decayed": curate.decayed.len(),
207            "recovered": curate.recovered.len(),
208            "orphans": curate.orphans.len(),
209            "warnings": curate.warnings,
210        })
211    }
212
213    fn distill_batch(&self) -> Result<DistillBatchReport> {
214        let run_id = gen_uuid();
215        let now = utc_now_iso();
216
217        // Atomically claim a batch of 'new' logs → mark 'screening'.
218        self.storage.begin_immediate()?;
219        let logs = match self
220            .storage
221            .claim_distill_batch(&run_id, self.distill_batch_size, &now)
222        {
223            Ok(l) => {
224                self.storage.commit()?;
225                l
226            }
227            Err(e) => {
228                let _ = self.storage.rollback();
229                return Err(e);
230            }
231        };
232
233        let mut chunks_by_log: HashMap<String, Vec<DistilledChunk>> = HashMap::new();
234        let mut failed_logs = HashSet::new();
235        let mut distill_errors = Vec::new();
236        for log in &logs {
237            let log_id = log.get("id").and_then(Value::as_str).unwrap_or("");
238            let context_key = log.get("context_key").and_then(Value::as_str);
239            let related_logs: Vec<Value> = logs
240                .iter()
241                .filter(|other| {
242                    other.get("id").and_then(Value::as_str) == Some(log_id)
243                        || (context_key.is_some()
244                            && other.get("context_key").and_then(Value::as_str) == context_key)
245                })
246                .cloned()
247                .collect();
248            match self.distiller.distill_with_context(log, &related_logs) {
249                Ok(chunks) => {
250                    if chunks.iter().any(|chunk| chunk.source_log_id != log_id) {
251                        let error = "distiller returned a chunk for an unknown source log";
252                        failed_logs.insert(log_id.to_string());
253                        distill_errors.push(format!("{log_id}: {error}"));
254                        self.finish_distill_log(
255                            log_id,
256                            "failed",
257                            Some(&format!("distill_failed:{error}")),
258                            estimate_distill_prompt_tokens(log, &related_logs),
259                            0,
260                        )?;
261                        continue;
262                    }
263                    chunks_by_log.insert(log_id.to_string(), chunks);
264                }
265                Err(error) => {
266                    let note = format!("distill_failed:{error}");
267                    failed_logs.insert(log_id.to_string());
268                    distill_errors.push(format!("{log_id}: {error}"));
269                    self.finish_distill_log(
270                        log_id,
271                        "failed",
272                        Some(&note),
273                        estimate_distill_prompt_tokens(log, &related_logs),
274                        0,
275                    )?;
276                }
277            }
278        }
279
280        let mut count = 0;
281        let provenance = self.distiller.provenance();
282        for log in &logs {
283            let log_id = log.get("id").and_then(Value::as_str).unwrap_or("");
284            if failed_logs.contains(log_id) {
285                continue;
286            }
287            let context_key = log.get("context_key").and_then(Value::as_str);
288            let related_logs: Vec<Value> = logs
289                .iter()
290                .filter(|other| {
291                    other.get("id").and_then(Value::as_str) == Some(log_id)
292                        || (context_key.is_some()
293                            && other.get("context_key").and_then(Value::as_str) == context_key)
294                })
295                .cloned()
296                .collect();
297            let prompt_tokens = estimate_distill_prompt_tokens(log, &related_logs);
298            let chunks = chunks_by_log.remove(log_id).unwrap_or_default();
299            let completion_tokens = chunks
300                .iter()
301                .map(estimate_distilled_chunk_tokens)
302                .sum::<i64>();
303            if chunks.is_empty() {
304                self.finish_distill_log(
305                    log_id,
306                    "discarded",
307                    Some("insufficient_material"),
308                    prompt_tokens,
309                    completion_tokens,
310                )?;
311                continue;
312            }
313
314            // Prepare all chunks + embeddings outside the write transaction so
315            // slow embedding calls do not hold an exclusive SQLite lock.
316            // Supports N >= 1 chunks per log (e.g. multi-concept LLM distillation).
317            // Bad individual chunks are skipped; valid siblings still survive.
318            struct PreparedChunk {
319                row: ChunkRow,
320                cvec_bytes: Vec<u8>,
321                tvec_bytes: Vec<u8>,
322            }
323            let mut prepared: Vec<PreparedChunk> = Vec::with_capacity(chunks.len());
324            let mut embedding_failures = 0_usize;
325            for dc in chunks {
326                let (content, action) = self.sanitize_content(&dc.content);
327                if action == SanitizeAction::Discard {
328                    continue; // skip this chunk, try others
329                }
330                let h = content_hash(&content);
331                if self.storage.is_hash_invalidated(&h)? {
332                    continue; // skip invalidated content, try others
333                }
334                let redacted = action == SanitizeAction::Redact;
335                let conf = if redacted { 0.4 } else { 0.55 };
336                let now2 = utc_now_iso();
337                let chunk_id = gen_uuid();
338                let tokens = estimate_tokens(&content) as i64;
339                // Short label for the web UI's row-skill slot; fall back to the
340                // canonical trigger phrase when the distiller produced none.
341                let skill_name = dc
342                    .skill_name
343                    .clone()
344                    .or_else(|| dc.trigger_desc.clone())
345                    .filter(|s| !s.trim().is_empty());
346                // Per-chunk fallback provenance wins over the batch-level provider,
347                // so heuristic-fallback chunks are tagged accurately.
348                let distill_provider = dc
349                    .provider_override
350                    .clone()
351                    .or_else(|| provenance.provider.clone());
352                let row = ChunkRow {
353                    id: chunk_id,
354                    skill_name,
355                    content: content.clone(),
356                    trigger_desc: dc.trigger_desc.clone(),
357                    anti_trigger_desc: dc.anti_trigger_desc,
358                    content_hash: h,
359                    token_count: Some(tokens),
360                    origin: "distilled".to_string(),
361                    // 蒸馏 chunk 继承源 episodic_log 的 agent(经验来自哪个 agent),
362                    // 而非运行 evolve 的进程(可能是 daemon/cron)。
363                    agent: log
364                        .get("agent")
365                        .and_then(Value::as_str)
366                        .map(str::to_string),
367                    distilled_from: Some(dc.source_log_id),
368                    distill_provider,
369                    distill_model: provenance.model.clone(),
370                    distill_prompt_version: provenance.prompt_version.clone(),
371                    state: "pending".to_string(),
372                    state_reason: Some("init:distilled".to_string()),
373                    confidence: conf,
374                    confidence_reason: Some("init:distilled".to_string()),
375                    version: 1,
376                    embed_version: 1,
377                    created_at: now2.clone(),
378                    updated_at: now2,
379                    ..Default::default()
380                };
381                let (cvec_res, tvec_res) = self.embed_pair(
382                    &content,
383                    row.trigger_desc.as_deref().unwrap_or(&content),
384                    "distill",
385                );
386                let cvec = match cvec_res {
387                    Ok(v) if v.len() == self.embedding.content_dim() => v,
388                    // A failed or wrong-dimension embedding is deferred (not
389                    // written): a dim mismatch would be silently dropped at
390                    // search time, so it must never reach storage.
391                    _ => {
392                        embedding_failures += 1;
393                        continue;
394                    }
395                };
396                let tvec = match tvec_res {
397                    Ok(v) if v.len() == self.embedding.trigger_dim() => v,
398                    _ => {
399                        embedding_failures += 1;
400                        continue;
401                    }
402                };
403                prepared.push(PreparedChunk {
404                    row,
405                    cvec_bytes: pack_embedding(&cvec),
406                    tvec_bytes: pack_embedding(&tvec),
407                });
408            }
409
410            if prepared.is_empty() {
411                let note = if embedding_failures > 0 {
412                    "embedding_failed"
413                } else {
414                    "all_chunks_filtered"
415                };
416                self.finish_distill_log(
417                    log_id,
418                    if embedding_failures > 0 {
419                        "failed"
420                    } else {
421                        "discarded"
422                    },
423                    Some(note),
424                    prompt_tokens,
425                    completion_tokens,
426                )?;
427                if embedding_failures > 0 {
428                    failed_logs.insert(log_id.to_string());
429                }
430                continue;
431            }
432
433            // Write all chunks, vectors, token accounting, and terminal log state atomically.
434            let accounted_at = utc_now_iso();
435            self.storage.begin_immediate()?;
436            let write_result = (|| -> Result<()> {
437                for pc in &prepared {
438                    self.storage.insert_chunk(&pc.row)?;
439                    self.storage
440                        .insert_vec_content(&pc.row.id, &pc.cvec_bytes)?;
441                    self.storage
442                        .insert_vec_trigger(&pc.row.id, &pc.tvec_bytes)?;
443                }
444                let note = (embedding_failures > 0)
445                    .then(|| format!("partial_embedding_failures:{embedding_failures}"));
446                self.storage.finish_distill_log(
447                    log_id,
448                    "distilled",
449                    note.as_deref(),
450                    prompt_tokens,
451                    completion_tokens,
452                    &accounted_at,
453                )?;
454                self.storage.commit()
455            })();
456            if let Err(error) = write_result {
457                let _ = self.storage.rollback();
458                let note = format!("distill_write_failed:{error}");
459                self.finish_distill_log(
460                    log_id,
461                    "failed",
462                    Some(&note),
463                    prompt_tokens,
464                    completion_tokens,
465                )?;
466                failed_logs.insert(log_id.to_string());
467                continue;
468            }
469            count += 1;
470        }
471        if !distill_errors.is_empty() {
472            // Log failures but do not abort: successful chunks are already committed and
473            // failed logs are marked 'failed' for bounded retry. Returning Ok preserves
474            // evolve request state and allows finish_covered_evolve_requests to run.
475            eprintln!(
476                "[innate] distillation partial failure ({} log(s)): {}",
477                distill_errors.len(),
478                distill_errors.join("; ")
479            );
480        }
481        Ok(DistillBatchReport {
482            distilled: count,
483            failed: failed_logs.len(),
484        })
485    }
486
487    fn finish_distill_log(
488        &self,
489        log_id: &str,
490        state: &str,
491        note: Option<&str>,
492        prompt_tokens: i64,
493        completion_tokens: i64,
494    ) -> Result<()> {
495        let accounted_at = utc_now_iso();
496        self.storage.begin_immediate()?;
497        let result = (|| -> Result<()> {
498            self.storage.finish_distill_log(
499                log_id,
500                state,
501                note,
502                prompt_tokens,
503                completion_tokens,
504                &accounted_at,
505            )?;
506            self.storage.commit()
507        })();
508        if result.is_err() {
509            let _ = self.storage.rollback();
510        }
511        result
512    }
513
514    pub(super) fn distill_token_period_start(&self, now: &str) -> Result<String> {
515        let window_hours = self
516            .storage
517            .get_meta("evolve.distill_token_window_hours")?
518            .and_then(|value| value.parse::<i64>().ok())
519            .unwrap_or(24)
520            .max(1);
521        Ok(hours_ago(now, window_hours))
522    }
523}