diff --git a/src/agent/context_compressor.rs b/src/agent/context_compressor.rs index 5d0b018..eb0eddf 100644 --- a/src/agent/context_compressor.rs +++ b/src/agent/context_compressor.rs @@ -536,22 +536,19 @@ impl ContextCompressor { summary ); let key = format!("ctx_compressed_{}", uuid::Uuid::new_v4()); - let mm = self.memory.clone(); - let sid = self.session_id.clone(); - tokio::spawn(async move { - if let Err(e) = mm - .store( - &key, - &timeline_content, - crate::memory::MemoryCategory::Timeline, - sid.as_deref(), - Some(0.3), - ) - .await - { - tracing::warn!(error = %e, "Failed to store compressed context as timeline"); - } - }); + if let Err(e) = self + .memory + .store( + &key, + &timeline_content, + crate::memory::MemoryCategory::Timeline, + self.session_id.as_deref(), + Some(0.3), + ) + .await + { + tracing::warn!(error = %e, "Failed to store compressed context as timeline"); + } // Add summary as a special user message new_messages.push(ChatMessage::user(format!( diff --git a/src/channels/feishu.rs b/src/channels/feishu.rs index 43a8e90..b3ec68d 100644 --- a/src/channels/feishu.rs +++ b/src/channels/feishu.rs @@ -1252,29 +1252,24 @@ impl FeishuChannel { forwarded_metadata.insert("feishu.parent_id".to_string(), pid.clone()); } - // Publish to bus asynchronously - let channel = self.clone(); - let bus = bus.clone(); - tokio::spawn(async move { + #[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"); + 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) = 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)] - 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 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"); - } - }); + tracing::debug!(open_id = %parsed.open_id, chat_id = %parsed.chat_id, "Message published to bus successfully"); + } } Ok(None) => {} Err(e) => { diff --git a/src/providers/anthropic.rs b/src/providers/anthropic.rs index 5a1f771..f925dea 100644 --- a/src/providers/anthropic.rs +++ b/src/providers/anthropic.rs @@ -330,24 +330,29 @@ impl LLMProvider for AnthropicProvider { 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 err_msg = format!("decode error: {} | body: {}", e, &body_text); - if let Some(ref storage) = self.storage { - let name = self.name.clone(); - let model = self.model_id.clone(); - let req = req_body_str.clone(); - let resp_body = body_text.clone(); - let dur = start.elapsed().as_millis() as u64; - let err = err_msg.clone(); - let s = storage.clone(); - tokio::spawn(async move { - let _ = s - .append_llm_call(&name, &model, &req, Some(&resp_body), Some(&err), dur) - .await; - }); + let anthropic_resp: AnthropicResponse = match serde_json::from_str(&body_text) { + Ok(response) => response, + Err(e) => { + let err_msg = format!("decode error: {} | body: {}", e, &body_text); + if let Some(ref storage) = self.storage { + let dur = start.elapsed().as_millis() as u64; + if let Err(error) = storage + .append_llm_call( + &self.name, + &self.model_id, + &req_body_str, + Some(&body_text), + Some(&err_msg), + 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 reasoning = None; diff --git a/src/providers/openai.rs b/src/providers/openai.rs index 6693590..918c16b 100644 --- a/src/providers/openai.rs +++ b/src/providers/openai.rs @@ -298,27 +298,29 @@ impl LLMProvider for OpenAIProvider { return Err(error.into()); } - let openai_resp: OpenAIResponse = serde_json::from_str(&text).map_err(|e| { - let err_msg = format!("decode error: {} | body: {}", e, &text); - if let Some(ref storage) = self.storage { - let name = self.name.clone(); - let model = self.model_id.clone(); - let req = req_body_str.clone(); - let resp = text.clone(); - let dur = start.elapsed().as_millis() as u64; - let err = err_msg.clone(); - let s = storage.clone(); - tokio::spawn(async move { - if let Err(e) = s - .append_llm_call(&name, &model, &req, Some(&resp), Some(&err), dur) + let openai_resp: OpenAIResponse = match serde_json::from_str(&text) { + Ok(response) => response, + Err(e) => { + let err_msg = format!("decode error: {} | body: {}", e, &text); + if let Some(ref storage) = self.storage { + let dur = start.elapsed().as_millis() as u64; + if let Err(error) = storage + .append_llm_call( + &self.name, + &self.model_id, + &req_body_str, + Some(&text), + Some(&err_msg), + dur, + ) .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 .choices