use std::sync::Arc;
use std::time::Duration;
use futures_util::StreamExt;
use tokio::sync::mpsc::UnboundedSender;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use crate::app::events::AppEvent;
use crate::entities::chat::Chat;
use crate::entities::message::{Message, MessageRole};
use crate::entities::sampling::{ReasoningEffort, SamplingConfig};
use crate::features::tools::self_model;
use crate::shared::api::{ApiMessage, ChatChunk, ChatRequest, EngineBackend, FinishReason};
use super::Orchestrator;
const IMPERSONATION_TIMEOUT: Duration = Duration::from_secs(600);
pub(super) struct ImpDone {
pub(super) id: Uuid,
pub(super) reason: FinishReason,
pub(super) prefill: Option<crate::shared::api::contract::Prefill>,
}
impl Orchestrator {
pub(super) fn impersonation_system(
&self,
profile: Option<&crate::entities::profile::Profile>,
loc: &'static crate::shared::i18n::Locale,
) -> String {
profile
.and_then(|p| p.impersonation_profile_id)
.and_then(|id| {
self.config
.impersonation_profiles
.iter()
.find(|ip| ip.id == id)
})
.map(|ip| ip.system_message.trim().to_string())
.filter(|s| !s.is_empty())
.unwrap_or_else(|| loc.t("prompt.impersonation.default").to_string())
}
pub(super) fn handle_impersonate(&mut self, seed: String) {
if !self.gen_state.is_idle() || self.imp_gen.is_some() {
return;
}
let Some(active_id) = self.active_id else {
let _ = self.evt_tx.send(AppEvent::Error(
self.ui_locale().t("ui.err.no_active_chat").into(),
));
return;
};
let backend = match self
.engines
.impersonation_backend_if_ready(self.config.impersonation_engine.mode, self.ui_locale())
{
Ok(backend) => backend,
Err(msg) => {
let _ = self.evt_tx.send(AppEvent::Error(msg));
return;
}
};
let Some(chat) = self.chats.iter().find(|c| c.id == active_id) else {
return;
};
let loc = self.profile_locale(chat.profile_id);
let profile = self.profiles.iter().find(|p| p.id == chat.profile_id);
let imp_system = self.impersonation_system(profile, loc);
let user_hint = profile
.filter(|p| {
p.enabled_tools
.iter()
.any(|t| t == self_model::GET_SELF_MODEL_ID)
})
.and_then(|_| {
self.storage
.db()
.self_model_get(chat.profile_id)
.ok()
.flatten()
})
.and_then(|m| {
let cap = crate::entities::self_model::SelfModelParams::from_settings(
&self.config.self_model,
)
.prompt_cap;
m.user_model.render_for_impersonation(cap, loc)
});
let request = build_impersonation_request(
chat,
chat.compaction_view(self.config.compaction.enabled),
imp_system,
&seed,
self.config.impersonation_sampling.clone(),
user_hint.as_deref(),
loc,
);
let id = Uuid::new_v4();
let cancel = CancellationToken::new();
self.imp_gen = Some(id);
self.imp_cancel = Some(cancel.clone());
let _ = self
.evt_tx
.send(AppEvent::ImpersonationStarted { generation_id: id });
let sessions = (self.config.impersonation_engine.mode
== crate::shared::config::ImpersonationMode::Shared)
.then(|| self.session_budget());
spawn_impersonation(
backend,
request,
id,
cancel,
self.ui_locale(),
self.evt_tx.clone(),
self.imp_done_tx.clone(),
sessions,
);
}
pub(super) fn handle_cancel_impersonation(&mut self) {
if let Some(token) = &self.imp_cancel {
token.cancel();
}
}
pub(super) fn handle_imp_done(&mut self, done: ImpDone) {
let ImpDone {
id,
reason,
prefill,
} = done;
if self.imp_gen == Some(id) {
self.imp_gen = None;
self.imp_cancel = None;
let _ = self.evt_tx.send(AppEvent::ImpersonationFinished {
generation_id: id,
reason,
});
}
self.note_slow_prefill(prefill);
}
}
pub(super) fn build_impersonation_request(
chat: &Chat,
compaction: Option<(&str, usize)>,
mut system: String,
seed: &str,
mut sampling: SamplingConfig,
user_hint: Option<&str>,
loc: &crate::shared::i18n::Locale,
) -> ChatRequest {
sampling.thinking = Some(false);
sampling.reasoning_effort = Some(ReasoningEffort::None);
sampling.reasoning_budget = Some(0);
let (summary, upto) = match compaction {
Some((s, i)) => (Some(s), i),
None => (None, 0),
};
system =
super::request::inject_compaction(Some(system), summary, false, loc).unwrap_or_default();
if let Some(hint) = user_hint.map(str::trim).filter(|s| !s.is_empty()) {
system.push_str("\n\n");
system.push_str(hint);
}
let messages = chat.messages[upto..]
.iter()
.filter_map(swap_role_message)
.collect();
let messages = alternate_for_template(messages, &mut system, loc);
let seed = seed.trim();
if !seed.is_empty() {
system.push_str("\n\n");
system.push_str(&loc.tf("prompt.impersonation.continue", &[("seed", seed)]));
}
ChatRequest {
continue_final: false,
system: Some(system),
messages,
sampling,
tools: Vec::new(),
}
}
pub(super) fn alternate_for_template(
messages: Vec<ApiMessage>,
system: &mut String,
loc: &crate::shared::i18n::Locale,
) -> Vec<ApiMessage> {
let mut merged: Vec<ApiMessage> = Vec::with_capacity(messages.len());
for message in messages {
match merged.last_mut() {
Some(last) if last.role == message.role => {
last.content.push_str("\n\n");
last.content.push_str(&message.content);
}
_ => merged.push(message),
}
}
if merged.len() >= 2 && merged[0].role == crate::shared::api::contract::ApiRole::Assistant {
let opening = merged.remove(0);
system.push_str("\n\n");
system.push_str(&loc.tf(
"prompt.impersonation.opening",
&[("text", &opening.content)],
));
}
merged
}
pub(super) fn swap_role_message(message: &Message) -> Option<ApiMessage> {
if message.text.trim().is_empty() {
return None;
}
match message.role {
MessageRole::User => Some(ApiMessage::assistant(&message.text)),
MessageRole::Assistant => Some(ApiMessage::user(&message.text)),
MessageRole::System | MessageRole::Tool => None,
}
}
#[allow(clippy::too_many_arguments)]
pub(super) fn spawn_impersonation(
backend: Arc<dyn EngineBackend>,
request: ChatRequest,
id: Uuid,
cancel: CancellationToken,
loc: &'static crate::shared::i18n::Locale,
evt_tx: UnboundedSender<AppEvent>,
done_tx: UnboundedSender<ImpDone>,
sessions: Option<Arc<crate::shared::session_budget::SessionBudget>>,
) {
tokio::spawn(async move {
let run = async {
let mut reason = FinishReason::Stop;
let mut prefill = None;
let estimate = super::generation::estimate_prompt_tokens(&request);
let _lane = match sessions.as_deref() {
Some(budget) => {
let need = budget.price(
crate::shared::session_budget::Shape::Impersonation,
estimate,
0,
request.sampling.max_tokens.map(|m| m as u64),
);
match budget
.acquire_silent(need, &cancel, "impersonation", false)
.await
{
Some(reservation) => Some(reservation),
None => return Ok((FinishReason::Cancelled, None)),
}
}
None => None,
};
let mut stream = backend.chat_stream(request, cancel.clone()).await?;
while let Some(chunk) = stream.next().await {
match chunk {
ChatChunk::Text(t) => {
let _ = evt_tx.send(AppEvent::ImpersonationChunk {
generation_id: id,
text: t,
});
}
ChatChunk::Retry {
attempt,
max,
delay,
} => {
tracing::info!(attempt, max, ?delay, "retrying a an impersonation turn");
}
ChatChunk::Error { message, .. } => {
tracing::warn!(error = %message, "engine error while impersonating");
}
ChatChunk::Usage(u) => {
if let Some(budget) = sessions.as_deref() {
budget.record_usage(
crate::shared::session_budget::Shape::Impersonation,
estimate,
u.prompt_tokens as u64,
);
prefill = u.prefill;
}
}
ChatChunk::Thoughts(_)
| ChatChunk::ThoughtsSignature(_)
| ChatChunk::ToolCall(_) => {}
ChatChunk::Finished(r) => {
reason = r;
break;
}
}
}
Ok::<(FinishReason, Option<crate::shared::api::contract::Prefill>), anyhow::Error>((
reason, prefill,
))
};
let (reason, prefill) = match tokio::time::timeout(IMPERSONATION_TIMEOUT, run).await {
Ok(Ok(landed)) => landed,
Ok(Err(err)) => {
let _ = evt_tx.send(AppEvent::Error(
loc.tf("ui.err.impersonation_failed", &[("err", &err.to_string())]),
));
(FinishReason::Error, None)
}
Err(_) => {
cancel.cancel();
(FinishReason::Length, None)
}
};
let _ = done_tx.send(ImpDone {
id,
reason,
prefill,
});
});
}