Skip to content

Commit 3b477b5

Browse files
fix(save_memory): bound concurrent enrichment LLM calls with semaphore
Background enrichment tasks from save_memory fired unbounded tokio::spawn per observation, risking LLM API flooding under burst saves. Add a shared Semaphore(3) to ObservationService that each enrichment task acquires before proceeding. Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
1 parent 79b5fb2 commit 3b477b5

2 files changed

Lines changed: 9 additions & 1 deletion

File tree

crates/service/src/observation_service/mod.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ use opencode_mem_embeddings::EmbeddingService;
1212
use opencode_mem_llm::LlmClient;
1313
use opencode_mem_storage::StorageBackend;
1414
use opencode_mem_storage::traits::ObservationStore;
15-
use tokio::sync::broadcast;
15+
use tokio::sync::{Semaphore, broadcast};
1616

1717
use crate::InfiniteMemoryService;
1818

@@ -32,6 +32,7 @@ pub struct ObservationService {
3232
pub(crate) injection_dedup_threshold: f32,
3333
pub(crate) project_filter: Option<opencode_mem_core::ProjectFilter>,
3434
pub(crate) low_value_filter: opencode_mem_core::LowValueFilter,
35+
pub(crate) enrichment_semaphore: Arc<Semaphore>,
3536
}
3637

3738
impl ObservationService {
@@ -86,6 +87,7 @@ impl ObservationService {
8687
injection_dedup_threshold,
8788
project_filter,
8889
low_value_filter,
90+
enrichment_semaphore: Arc::new(Semaphore::new(3)),
8991
}
9092
}
9193

crates/service/src/observation_service/save_memory.rs

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -104,8 +104,14 @@ impl ObservationService {
104104
let llm = self.llm.clone();
105105
let storage = self.storage.clone();
106106
let svc = self.clone();
107+
let semaphore = self.enrichment_semaphore.clone();
107108

108109
tokio::spawn(async move {
110+
let _permit = match semaphore.acquire().await {
111+
Ok(permit) => permit,
112+
Err(_) => return, // Semaphore closed — shut down gracefully
113+
};
114+
109115
let narrative = obs.narrative.as_deref().unwrap_or("");
110116
if narrative.is_empty() && obs.title.is_empty() {
111117
return;

0 commit comments

Comments
 (0)