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            "deduped": curate.deduped.len(),
205            "decayed": curate.decayed.len(),
206            "recovered": curate.recovered.len(),
207            "orphans": curate.orphans.len(),
208            "warnings": curate.warnings,
209        })
210    }
211
212    fn distill_batch(&self) -> Result<DistillBatchReport> {
213        let run_id = gen_uuid();
214        let now = utc_now_iso();
215
216        // Atomically claim a batch of 'new' logs → mark 'screening'.
217        self.storage.begin_immediate()?;
218        let logs = match self
219            .storage
220            .claim_distill_batch(&run_id, self.distill_batch_size, &now)
221        {
222            Ok(l) => {
223                self.storage.commit()?;
224                l
225            }
226            Err(e) => {
227                let _ = self.storage.rollback();
228                return Err(e);
229            }
230        };
231
232        let mut chunks_by_log: HashMap<String, Vec<DistilledChunk>> = HashMap::new();
233        let mut failed_logs = HashSet::new();
234        let mut distill_errors = Vec::new();
235        for log in &logs {
236            let log_id = log.get("id").and_then(Value::as_str).unwrap_or("");
237            let context_key = log.get("context_key").and_then(Value::as_str);
238            let related_logs: Vec<Value> = logs
239                .iter()
240                .filter(|other| {
241                    other.get("id").and_then(Value::as_str) == Some(log_id)
242                        || (context_key.is_some()
243                            && other.get("context_key").and_then(Value::as_str) == context_key)
244                })
245                .cloned()
246                .collect();
247            match self.distiller.distill_with_context(log, &related_logs) {
248                Ok(chunks) => {
249                    if chunks.iter().any(|chunk| chunk.source_log_id != log_id) {
250                        let error = "distiller returned a chunk for an unknown source log";
251                        failed_logs.insert(log_id.to_string());
252                        distill_errors.push(format!("{log_id}: {error}"));
253                        self.finish_distill_log(
254                            log_id,
255                            "failed",
256                            Some(&format!("distill_failed:{error}")),
257                            estimate_distill_prompt_tokens(log, &related_logs),
258                            0,
259                        )?;
260                        continue;
261                    }
262                    chunks_by_log.insert(log_id.to_string(), chunks);
263                }
264                Err(error) => {
265                    let note = format!("distill_failed:{error}");
266                    failed_logs.insert(log_id.to_string());
267                    distill_errors.push(format!("{log_id}: {error}"));
268                    self.finish_distill_log(
269                        log_id,
270                        "failed",
271                        Some(&note),
272                        estimate_distill_prompt_tokens(log, &related_logs),
273                        0,
274                    )?;
275                }
276            }
277        }
278
279        let mut count = 0;
280        let provenance = self.distiller.provenance();
281        for log in &logs {
282            let log_id = log.get("id").and_then(Value::as_str).unwrap_or("");
283            if failed_logs.contains(log_id) {
284                continue;
285            }
286            let context_key = log.get("context_key").and_then(Value::as_str);
287            let related_logs: Vec<Value> = logs
288                .iter()
289                .filter(|other| {
290                    other.get("id").and_then(Value::as_str) == Some(log_id)
291                        || (context_key.is_some()
292                            && other.get("context_key").and_then(Value::as_str) == context_key)
293                })
294                .cloned()
295                .collect();
296            let prompt_tokens = estimate_distill_prompt_tokens(log, &related_logs);
297            let chunks = chunks_by_log.remove(log_id).unwrap_or_default();
298            let completion_tokens = chunks
299                .iter()
300                .map(estimate_distilled_chunk_tokens)
301                .sum::<i64>();
302            if chunks.is_empty() {
303                self.finish_distill_log(
304                    log_id,
305                    "discarded",
306                    Some("insufficient_material"),
307                    prompt_tokens,
308                    completion_tokens,
309                )?;
310                continue;
311            }
312
313            // Prepare all chunks + embeddings outside the write transaction so
314            // slow embedding calls do not hold an exclusive SQLite lock.
315            // Supports N >= 1 chunks per log (e.g. multi-concept LLM distillation).
316            // Bad individual chunks are skipped; valid siblings still survive.
317            struct PreparedChunk {
318                row: ChunkRow,
319                cvec_bytes: Vec<u8>,
320                tvec_bytes: Vec<u8>,
321            }
322            let mut prepared: Vec<PreparedChunk> = Vec::with_capacity(chunks.len());
323            let mut embedding_failures = 0_usize;
324            for dc in chunks {
325                let (content, action) = self.sanitize_content(&dc.content);
326                if action == SanitizeAction::Discard {
327                    continue; // skip this chunk, try others
328                }
329                let h = content_hash(&content);
330                if self.storage.is_hash_invalidated(&h)? {
331                    continue; // skip invalidated content, try others
332                }
333                let redacted = action == SanitizeAction::Redact;
334                let conf = if redacted { 0.4 } else { 0.55 };
335                let now2 = utc_now_iso();
336                let chunk_id = gen_uuid();
337                let tokens = estimate_tokens(&content) as i64;
338                // Short label for the web UI's row-skill slot; fall back to the
339                // canonical trigger phrase when the distiller produced none.
340                let skill_name = dc
341                    .skill_name
342                    .clone()
343                    .or_else(|| dc.trigger_desc.clone())
344                    .filter(|s| !s.trim().is_empty());
345                // Per-chunk fallback provenance wins over the batch-level provider,
346                // so heuristic-fallback chunks are tagged accurately.
347                let distill_provider = dc
348                    .provider_override
349                    .clone()
350                    .or_else(|| provenance.provider.clone());
351                let row = ChunkRow {
352                    id: chunk_id,
353                    skill_name,
354                    content: content.clone(),
355                    trigger_desc: dc.trigger_desc.clone(),
356                    anti_trigger_desc: dc.anti_trigger_desc,
357                    content_hash: h,
358                    token_count: Some(tokens),
359                    origin: "distilled".to_string(),
360                    // 蒸馏 chunk 继承源 episodic_log 的 agent(经验来自哪个 agent),
361                    // 而非运行 evolve 的进程(可能是 daemon/cron)。
362                    agent: log
363                        .get("agent")
364                        .and_then(Value::as_str)
365                        .map(str::to_string),
366                    distilled_from: Some(dc.source_log_id),
367                    distill_provider,
368                    distill_model: provenance.model.clone(),
369                    distill_prompt_version: provenance.prompt_version.clone(),
370                    state: "pending".to_string(),
371                    state_reason: Some("init:distilled".to_string()),
372                    confidence: conf,
373                    confidence_reason: Some("init:distilled".to_string()),
374                    version: 1,
375                    embed_version: 1,
376                    created_at: now2.clone(),
377                    updated_at: now2,
378                    ..Default::default()
379                };
380                let (cvec_res, tvec_res) = self.embed_pair(
381                    &content,
382                    row.trigger_desc.as_deref().unwrap_or(&content),
383                    "distill",
384                );
385                let cvec = match cvec_res {
386                    Ok(v) if v.len() == self.embedding.content_dim() => v,
387                    // A failed or wrong-dimension embedding is deferred (not
388                    // written): a dim mismatch would be silently dropped at
389                    // search time, so it must never reach storage.
390                    _ => {
391                        embedding_failures += 1;
392                        continue;
393                    }
394                };
395                let tvec = match tvec_res {
396                    Ok(v) if v.len() == self.embedding.trigger_dim() => v,
397                    _ => {
398                        embedding_failures += 1;
399                        continue;
400                    }
401                };
402                prepared.push(PreparedChunk {
403                    row,
404                    cvec_bytes: pack_embedding(&cvec),
405                    tvec_bytes: pack_embedding(&tvec),
406                });
407            }
408
409            if prepared.is_empty() {
410                let note = if embedding_failures > 0 {
411                    "embedding_failed"
412                } else {
413                    "all_chunks_filtered"
414                };
415                self.finish_distill_log(
416                    log_id,
417                    if embedding_failures > 0 {
418                        "failed"
419                    } else {
420                        "discarded"
421                    },
422                    Some(note),
423                    prompt_tokens,
424                    completion_tokens,
425                )?;
426                if embedding_failures > 0 {
427                    failed_logs.insert(log_id.to_string());
428                }
429                continue;
430            }
431
432            // Write all chunks, vectors, token accounting, and terminal log state atomically.
433            let accounted_at = utc_now_iso();
434            self.storage.begin_immediate()?;
435            let write_result = (|| -> Result<()> {
436                for pc in &prepared {
437                    self.storage.insert_chunk(&pc.row)?;
438                    self.storage
439                        .insert_vec_content(&pc.row.id, &pc.cvec_bytes)?;
440                    self.storage
441                        .insert_vec_trigger(&pc.row.id, &pc.tvec_bytes)?;
442                }
443                let note = (embedding_failures > 0)
444                    .then(|| format!("partial_embedding_failures:{embedding_failures}"));
445                self.storage.finish_distill_log(
446                    log_id,
447                    "distilled",
448                    note.as_deref(),
449                    prompt_tokens,
450                    completion_tokens,
451                    &accounted_at,
452                )?;
453                self.storage.commit()
454            })();
455            if let Err(error) = write_result {
456                let _ = self.storage.rollback();
457                let note = format!("distill_write_failed:{error}");
458                self.finish_distill_log(
459                    log_id,
460                    "failed",
461                    Some(&note),
462                    prompt_tokens,
463                    completion_tokens,
464                )?;
465                failed_logs.insert(log_id.to_string());
466                continue;
467            }
468            count += 1;
469        }
470        if !distill_errors.is_empty() {
471            // Log failures but do not abort: successful chunks are already committed and
472            // failed logs are marked 'failed' for bounded retry. Returning Ok preserves
473            // evolve request state and allows finish_covered_evolve_requests to run.
474            eprintln!(
475                "[innate] distillation partial failure ({} log(s)): {}",
476                distill_errors.len(),
477                distill_errors.join("; ")
478            );
479        }
480        Ok(DistillBatchReport {
481            distilled: count,
482            failed: failed_logs.len(),
483        })
484    }
485
486    fn finish_distill_log(
487        &self,
488        log_id: &str,
489        state: &str,
490        note: Option<&str>,
491        prompt_tokens: i64,
492        completion_tokens: i64,
493    ) -> Result<()> {
494        let accounted_at = utc_now_iso();
495        self.storage.begin_immediate()?;
496        let result = (|| -> Result<()> {
497            self.storage.finish_distill_log(
498                log_id,
499                state,
500                note,
501                prompt_tokens,
502                completion_tokens,
503                &accounted_at,
504            )?;
505            self.storage.commit()
506        })();
507        if result.is_err() {
508            let _ = self.storage.rollback();
509        }
510        result
511    }
512
513    pub(super) fn distill_token_period_start(&self, now: &str) -> Result<String> {
514        let window_hours = self
515            .storage
516            .get_meta("evolve.distill_token_window_hours")?
517            .and_then(|value| value.parse::<i64>().ok())
518            .unwrap_or(24)
519            .max(1);
520        Ok(hours_ago(now, window_hours))
521    }
522}