fix(runtime): eliminate unowned persistence tasks

This commit is contained in:
xiaoxixi 2026-07-14 11:50:55 +08:00
parent 59ecb27c06
commit 24d3e26b43
4 changed files with 71 additions and 72 deletions

View File

@ -536,22 +536,19 @@ impl ContextCompressor {
summary summary
); );
let key = format!("ctx_compressed_{}", uuid::Uuid::new_v4()); let key = format!("ctx_compressed_{}", uuid::Uuid::new_v4());
let mm = self.memory.clone(); if let Err(e) = self
let sid = self.session_id.clone(); .memory
tokio::spawn(async move { .store(
if let Err(e) = mm &key,
.store( &timeline_content,
&key, crate::memory::MemoryCategory::Timeline,
&timeline_content, self.session_id.as_deref(),
crate::memory::MemoryCategory::Timeline, Some(0.3),
sid.as_deref(), )
Some(0.3), .await
) {
.await tracing::warn!(error = %e, "Failed to store compressed context as timeline");
{ }
tracing::warn!(error = %e, "Failed to store compressed context as timeline");
}
});
// Add summary as a special user message // Add summary as a special user message
new_messages.push(ChatMessage::user(format!( new_messages.push(ChatMessage::user(format!(

View File

@ -1252,29 +1252,24 @@ impl FeishuChannel {
forwarded_metadata.insert("feishu.parent_id".to_string(), pid.clone()); forwarded_metadata.insert("feishu.parent_id".to_string(), pid.clone());
} }
// Publish to bus asynchronously #[cfg(debug_assertions)]
let channel = self.clone(); tracing::debug!(open_id = %parsed.open_id, chat_id = %parsed.chat_id, content_len = %parsed.content.len(), media_count = %parsed.media.len(), "Publishing message to bus");
let bus = bus.clone(); let msg = crate::bus::InboundMessage {
tokio::spawn(async move { channel: "feishu".to_string(),
sender_id: parsed.open_id.clone(),
chat_id: parsed.chat_id.clone(),
content: parsed.content.clone(),
timestamp: crate::bus::message::current_timestamp(),
media: parsed.media.clone(),
metadata: std::collections::HashMap::new(),
forwarded_metadata,
};
if let Err(e) = self.handle_and_publish(&bus, &msg).await {
tracing::error!(error = %e, open_id = %parsed.open_id, chat_id = %parsed.chat_id, "Failed to publish Feishu message to bus");
} else {
#[cfg(debug_assertions)] #[cfg(debug_assertions)]
tracing::debug!(open_id = %parsed.open_id, chat_id = %parsed.chat_id, content_len = %parsed.content.len(), media_count = %parsed.media.len(), "Publishing message to bus"); tracing::debug!(open_id = %parsed.open_id, chat_id = %parsed.chat_id, "Message published to bus successfully");
let msg = crate::bus::InboundMessage { }
channel: "feishu".to_string(),
sender_id: parsed.open_id.clone(),
chat_id: parsed.chat_id.clone(),
content: parsed.content.clone(),
timestamp: crate::bus::message::current_timestamp(),
media: parsed.media.clone(),
metadata: std::collections::HashMap::new(),
forwarded_metadata,
};
if let Err(e) = channel.handle_and_publish(&bus, &msg).await {
tracing::error!(error = %e, open_id = %parsed.open_id, chat_id = %parsed.chat_id, "Failed to publish Feishu message to bus");
} else {
#[cfg(debug_assertions)]
tracing::debug!(open_id = %parsed.open_id, chat_id = %parsed.chat_id, "Message published to bus successfully");
}
});
} }
Ok(None) => {} Ok(None) => {}
Err(e) => { Err(e) => {

View File

@ -330,24 +330,29 @@ impl LLMProvider for AnthropicProvider {
return Err(format!("API error ({}): {}", status.as_u16(), error_msg).into()); return Err(format!("API error ({}): {}", status.as_u16(), error_msg).into());
} }
let anthropic_resp: AnthropicResponse = serde_json::from_str(&body_text).map_err(|e| { let anthropic_resp: AnthropicResponse = match serde_json::from_str(&body_text) {
let err_msg = format!("decode error: {} | body: {}", e, &body_text); Ok(response) => response,
if let Some(ref storage) = self.storage { Err(e) => {
let name = self.name.clone(); let err_msg = format!("decode error: {} | body: {}", e, &body_text);
let model = self.model_id.clone(); if let Some(ref storage) = self.storage {
let req = req_body_str.clone(); let dur = start.elapsed().as_millis() as u64;
let resp_body = body_text.clone(); if let Err(error) = storage
let dur = start.elapsed().as_millis() as u64; .append_llm_call(
let err = err_msg.clone(); &self.name,
let s = storage.clone(); &self.model_id,
tokio::spawn(async move { &req_body_str,
let _ = s Some(&body_text),
.append_llm_call(&name, &model, &req, Some(&resp_body), Some(&err), dur) Some(&err_msg),
.await; dur,
}); )
.await
{
tracing::warn!("failed to persist LLM call (decode error): {}", error);
}
}
return Err(err_msg.into());
} }
err_msg };
})?;
let mut content = String::new(); let mut content = String::new();
let mut reasoning = None; let mut reasoning = None;

View File

@ -298,27 +298,29 @@ impl LLMProvider for OpenAIProvider {
return Err(error.into()); return Err(error.into());
} }
let openai_resp: OpenAIResponse = serde_json::from_str(&text).map_err(|e| { let openai_resp: OpenAIResponse = match serde_json::from_str(&text) {
let err_msg = format!("decode error: {} | body: {}", e, &text); Ok(response) => response,
if let Some(ref storage) = self.storage { Err(e) => {
let name = self.name.clone(); let err_msg = format!("decode error: {} | body: {}", e, &text);
let model = self.model_id.clone(); if let Some(ref storage) = self.storage {
let req = req_body_str.clone(); let dur = start.elapsed().as_millis() as u64;
let resp = text.clone(); if let Err(error) = storage
let dur = start.elapsed().as_millis() as u64; .append_llm_call(
let err = err_msg.clone(); &self.name,
let s = storage.clone(); &self.model_id,
tokio::spawn(async move { &req_body_str,
if let Err(e) = s Some(&text),
.append_llm_call(&name, &model, &req, Some(&resp), Some(&err), dur) Some(&err_msg),
dur,
)
.await .await
{ {
tracing::warn!("failed to persist LLM call (decode error): {}", e); tracing::warn!("failed to persist LLM call (decode error): {}", error);
} }
}); }
return Err(err_msg.into());
} }
err_msg };
})?;
let first_choice = openai_resp let first_choice = openai_resp
.choices .choices