use std::time::Duration;
use tokio_util::sync::CancellationToken;
use chrono::{DateTime, Utc};
use crate::app::events::BackgroundKind;
use crate::entities::chat::{Chat, DeletedCause};
use crate::entities::message::{Message, MessageRole};
use crate::entities::profile::ToolId;
use crate::entities::sampling::SamplingConfig;
use crate::features::tools::{notes, self_model};
use crate::shared::api::{ApiMessage, ChatRequest};
use super::Orchestrator;
use super::request::last_user_message_at;
use super::tool_loop;
const REFLECT_MAX_TOKENS: usize = 2048;
const REFLECT_MAX_ROUNDS: u32 = 6;
const REFLECT_TIMEOUT: Duration = Duration::from_secs(120);
const REFLECT_TOOL_IDS: &[&str] = &[
self_model::GET_SELF_MODEL_ID,
self_model::UPDATE_SELF_MODEL_ID,
self_model::UPDATE_USER_MODEL_ID,
self_model::ADD_INSIGHT_ID,
notes::NOTE_RECALL_ID,
notes::NOTE_REVISE_ID,
notes::NOTE_SUPERSEDE_ID,
notes::NOTE_MERGE_ID,
notes::NOTE_LINK_ID,
notes::NOTE_NEIGHBORS_ID,
];
fn reflect_system_message(loc: &crate::shared::i18n::Locale) -> String {
loc.tf(
"prompt.reflect.system",
&[("core", self_model::policy_core(loc))],
)
}
pub(super) fn reflect_window(messages: &[Message], reflected_upto: Option<usize>) -> (usize, u32) {
let wm = reflected_upto.unwrap_or(0).min(messages.len());
let count = messages[wm..]
.iter()
.filter(|m| m.role == MessageRole::Assistant && !m.text.trim().is_empty())
.count() as u32;
(wm, count)
}
pub(super) fn behavior_markers(
chat: &Chat,
since: Option<DateTime<Utc>>,
loc: &crate::shared::i18n::Locale,
) -> Option<String> {
let (mut regen, mut del, mut rewrite) = (0u32, 0u32, 0u32);
for d in &chat.deleted {
if let Some(s) = since
&& d.deleted_at <= s
{
continue; }
match d.cause {
Some(DeletedCause::Regenerate) => regen += 1,
Some(DeletedCause::DeleteExchange) => del += 1,
Some(DeletedCause::Rewrite) => rewrite += 1,
None => {}
}
}
if regen == 0 && del == 0 && rewrite == 0 {
return None;
}
let mut about_user: Vec<String> = Vec::new();
if regen > 0 {
about_user.push(loc.tf("reflect.behavior.regen", &[("n", ®en.to_string())]));
}
if del > 0 {
about_user.push(loc.tf("reflect.behavior.deleted", &[("n", &del.to_string())]));
}
let mut out = String::new();
if !about_user.is_empty() {
out.push_str(loc.t("reflect.behavior.header"));
out.push_str(&about_user.join("; "));
out.push('.');
}
if rewrite > 0 {
if !out.is_empty() {
out.push(' ');
}
out.push_str(&loc.tf("reflect.behavior.rewrite", &[("n", &rewrite.to_string())]));
}
Some(out)
}
impl Orchestrator {
pub(super) fn maybe_auto_reflect(&mut self, chat_id: uuid::Uuid) {
let every = self.config.self_model.auto_reflect_every;
if every == 0 {
return;
}
let profile_id;
let lang; let system_message;
let last_user;
let mut digest;
let watermark; let allowed: Vec<ToolId>;
{
let Some(chat) = self.chats.iter().find(|c| c.id == chat_id) else {
return;
};
profile_id = chat.profile_id;
let Some(profile) = self.profiles.iter().find(|p| p.id == profile_id) else {
return;
};
lang = profile.language;
if !profile
.enabled_tools
.iter()
.any(|t| t == self_model::GET_SELF_MODEL_ID)
{
return;
}
let (wm, count) = reflect_window(&chat.messages, chat.reflected_upto);
if !tool_loop::due(count, every) {
return;
}
allowed = REFLECT_TOOL_IDS
.iter()
.filter(|id| profile.enabled_tools.iter().any(|t| t == **id))
.map(ToString::to_string)
.collect();
system_message = chat.system_message.clone();
last_user = last_user_message_at(chat);
let loc = crate::shared::i18n::locale(lang);
let Some(d) =
crate::features::rename_chat::build_conversation_digest(&chat.messages[wm..], loc)
else {
return; };
digest = match behavior_markers(chat, chat.reflected_at, loc) {
Some(markers) => format!("{d}\n\n{markers}"),
None => d,
};
watermark = chat.messages.len();
}
if let Some(overview) = notes::build_self_consolidation_overview(
&self.storage,
profile_id,
crate::shared::i18n::locale(lang),
) {
digest = format!("{digest}\n\n{overview}");
}
if self.bg_running(BackgroundKind::Reflection) {
return;
}
let Ok(backend) = self.engines.backend_if_ready(self.ui_locale()) else {
return;
};
let window = self.chats.iter_mut().find(|c| c.id == chat_id).map(|chat| {
let before = super::background::Window::Reflection {
chat: chat_id,
upto: chat.reflected_upto,
at: chat.reflected_at,
};
chat.reflected_upto = Some(watermark);
chat.reflected_at = Some(Utc::now());
before
});
self.mark_dirty(chat_id);
let cancel = CancellationToken::new();
let sessions = self.session_budget();
let ctx = self.background_tool_ctx(
backend.clone(),
sessions,
profile_id,
chat_id,
system_message,
last_user,
lang,
cancel.clone(),
);
let sampling = SamplingConfig {
max_tokens: Some(REFLECT_MAX_TOKENS),
temperature: Some(0.4),
..Default::default()
};
let request = ChatRequest {
continue_final: false,
system: Some(reflect_system_message(crate::shared::i18n::locale(lang))),
messages: vec![ApiMessage::user(digest)],
sampling,
tools: self
.registry
.schemas_for(&allowed, crate::shared::i18n::locale(lang)),
};
let acted = std::sync::Arc::new(super::background::Acted::default());
tool_loop::spawn_silent_loop(tool_loop::SilentLoop {
backend,
registry: self.registry.clone(),
ctx,
request,
allowed,
cancel: cancel.clone(),
max_rounds: REFLECT_MAX_ROUNDS,
timeout: REFLECT_TIMEOUT,
label: "auto-reflection",
profile_id,
kind: BackgroundKind::Reflection,
done_tx: self.bg_done_tx.clone(),
acted: acted.clone(),
summary_semantics: Some(tool_loop::SummarySemantics {
embedder: self.engines.embedder(),
storage: self.storage.clone(),
profile_id,
loc: crate::shared::i18n::locale(lang),
}),
});
self.begin_bg(
BackgroundKind::Reflection,
cancel,
window.map(|window| super::background::Refund { window, acted }),
);
}
}
#[cfg(test)]
mod tests {
use super::{
Chat, DeletedCause, Message, behavior_markers, reflect_system_message, reflect_window,
};
fn ru() -> &'static crate::shared::i18n::Locale {
crate::shared::i18n::locale(crate::shared::i18n::Lang::Ru)
}
#[test]
fn reflect_system_message_composes_from_policy_core() {
let msg = reflect_system_message(ru());
assert!(msg.contains(crate::features::tools::self_model::policy_core(ru())));
assert!(msg.contains("get_self_model"));
assert!(msg.contains("Поведенческие сигналы"));
assert!(msg.contains("только вызывай инструменты"));
assert!(msg.contains("note_link"));
}
#[test]
fn reflect_system_message_localized_for_all_langs() {
for &lang in crate::shared::i18n::Lang::ALL {
let l = crate::shared::i18n::locale(lang);
let msg = reflect_system_message(l);
assert!(
msg.contains(crate::features::tools::self_model::policy_core(l)),
"{lang:?}: policy_core not embedded"
);
assert!(
!msg.contains("{core}"),
"{lang:?}: placeholder not substituted"
);
assert!(
msg.contains("get_self_model") && msg.contains("note_link"),
"{lang:?}"
);
}
}
#[test]
fn reflect_tools_include_graph() {
use super::{REFLECT_TOOL_IDS, notes};
assert!(REFLECT_TOOL_IDS.contains(¬es::NOTE_LINK_ID));
assert!(REFLECT_TOOL_IDS.contains(¬es::NOTE_NEIGHBORS_ID));
assert!(REFLECT_TOOL_IDS.contains(¬es::NOTE_RECALL_ID));
}
#[test]
fn reflect_message_nudges_cross_organ_linking() {
let msg = super::reflect_system_message(ru());
assert!(msg.contains("note_recall"));
assert!(msg.contains("о собеседнике"));
assert!(msg.contains("Обзор наблюдений для консолидации"));
}
#[test]
fn behavior_markers_counts_by_cause_and_filters_since() {
use crate::entities::chat::DeletedExchange;
use crate::entities::profile::Profile;
use chrono::{Duration, Utc};
let p = Profile::new("P", "s");
let mut chat = Chat::from_profile(&p, "t");
let base = Utc::now();
let mk = |at, cause| DeletedExchange {
deleted_at: at,
messages: vec![Message::user("x")],
draft: String::new(),
cause: Some(cause),
};
chat.deleted = vec![
mk(base - Duration::hours(1), DeletedCause::Regenerate),
mk(base - Duration::hours(2), DeletedCause::Regenerate),
mk(base - Duration::hours(3), DeletedCause::DeleteExchange),
mk(base - Duration::days(5), DeletedCause::Rewrite), DeletedExchange {
deleted_at: base - Duration::hours(1),
messages: vec![Message::user("x")],
draft: String::new(),
cause: None, },
];
let out = behavior_markers(&chat, Some(base - Duration::hours(4)), ru()).unwrap();
assert!(out.contains("перегенерировал твой ответ ×2"));
assert!(out.contains("удалил обмен ×1"));
assert!(!out.contains("переписывал"));
let all = behavior_markers(&chat, None, ru()).unwrap();
assert!(all.contains("Ты сам переписывал свой ответ ×1"));
let empty = Chat::from_profile(&p, "t2");
assert!(behavior_markers(&empty, None, ru()).is_none());
}
#[test]
fn reflect_window_counts_assistant_from_watermark() {
let msgs = vec![
Message::user("u1"),
Message::assistant("a1"),
Message::user("u2"),
Message::assistant("a2"),
Message::assistant(""), Message::user("u3"),
Message::assistant("a3"),
];
assert_eq!(reflect_window(&msgs, None), (0, 3));
assert_eq!(reflect_window(&msgs, Some(4)), (4, 1));
assert_eq!(reflect_window(&msgs, Some(msgs.len())), (7, 0));
}
#[test]
fn reflect_window_clamps_past_watermark_after_truncation() {
let msgs = vec![Message::user("u1"), Message::assistant("a1")];
assert_eq!(reflect_window(&msgs, Some(99)), (2, 0));
}
}