mirror of
https://github.com/tinyhumansai/openhuman.git
synced 2026-07-30 23:14:37 +00:00
perf(learning): overlap reflection pattern/pref writes with join_all (#3216)
This commit is contained in:
@@ -256,33 +256,66 @@ impl ReflectionHook {
|
||||
);
|
||||
}
|
||||
|
||||
// Patterns and raw user-preference facts are independent best-effort
|
||||
// writes with no shared mutable state. Each `store(...)` does a
|
||||
// markdown write + SQLite tx + an embedding round-trip that does NOT
|
||||
// hold the SQLite connection mutex across its `.await` (the lock is
|
||||
// taken after the embed), so overlapping them with `join_all` lets
|
||||
// their network/disk waits run concurrently instead of paying one
|
||||
// document's full latency after another on the post-turn path.
|
||||
//
|
||||
// (`user_preferences` are also handled by `UserProfileHook`; the raw
|
||||
// text is stored here as before.)
|
||||
let mut writes: Vec<(&'static str, String, &str)> =
|
||||
Vec::with_capacity(output.patterns.len() + output.user_preferences.len());
|
||||
for pattern in &output.patterns {
|
||||
let slug = slugify(pattern);
|
||||
let key = format!("pat/{slug}");
|
||||
self.memory
|
||||
.store(
|
||||
"learning_patterns",
|
||||
&key,
|
||||
pattern,
|
||||
MemoryCategory::Custom("learning_patterns".into()),
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
writes.push((
|
||||
"learning_patterns",
|
||||
format!("pat/{}", slugify(pattern)),
|
||||
pattern.as_str(),
|
||||
));
|
||||
}
|
||||
for pref in &output.user_preferences {
|
||||
writes.push((
|
||||
"user_profile",
|
||||
format!("pref/{}", slugify(pref)),
|
||||
pref.as_str(),
|
||||
));
|
||||
}
|
||||
|
||||
// User preferences are handled by UserProfileHook, but store raw if present
|
||||
for pref in &output.user_preferences {
|
||||
let slug = slugify(pref);
|
||||
let key = format!("pref/{slug}");
|
||||
self.memory
|
||||
.store(
|
||||
"user_profile",
|
||||
&key,
|
||||
pref,
|
||||
MemoryCategory::Custom("user_profile".into()),
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
let results =
|
||||
futures::future::join_all(writes.into_iter().map(|(namespace, key, content)| {
|
||||
let memory = &self.memory;
|
||||
async move {
|
||||
memory
|
||||
.store(
|
||||
namespace,
|
||||
&key,
|
||||
content,
|
||||
MemoryCategory::Custom(namespace.into()),
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.map_err(|e| (key, e))
|
||||
}
|
||||
}))
|
||||
.await;
|
||||
|
||||
// Preserve the previous error contract: surface the first failure so
|
||||
// the caller still rolls back the session counter. Unlike the old
|
||||
// early-`?`, every independent write is now attempted regardless of a
|
||||
// sibling's failure (these are best-effort learning writes).
|
||||
let mut first_err = None;
|
||||
for result in results {
|
||||
if let Err((key, e)) = result {
|
||||
log::warn!("[learning] failed to persist reflection write {key}: {e}");
|
||||
if first_err.is_none() {
|
||||
first_err = Some(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Some(e) = first_err {
|
||||
return Err(e);
|
||||
}
|
||||
|
||||
// Reflection persistence is handled by the caller
|
||||
|
||||
@@ -273,6 +273,53 @@ async fn store_reflection_persists_all_categories() {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn store_reflection_persists_every_pattern_and_pref_concurrently() {
|
||||
// Regression guard for the `join_all` batching: every independent
|
||||
// pattern / preference write must land regardless of completion order
|
||||
// (the old code awaited them one at a time).
|
||||
let memory_impl = Arc::new(MockMemory::default());
|
||||
let memory: Arc<dyn Memory> = memory_impl.clone();
|
||||
let hook = ReflectionHook::new(
|
||||
reflection_config(),
|
||||
Arc::new(Config::default()),
|
||||
memory,
|
||||
None,
|
||||
);
|
||||
|
||||
hook.store_reflection(&ReflectionOutput {
|
||||
observations: Vec::new(),
|
||||
patterns: vec![
|
||||
"Pattern One".into(),
|
||||
"Pattern Two".into(),
|
||||
"Pattern Three".into(),
|
||||
],
|
||||
user_preferences: vec!["Pref One".into(), "Pref Two".into()],
|
||||
user_reflections: Vec::new(),
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let keys: Vec<String> = memory_impl.entries.lock().keys().cloned().collect();
|
||||
for expected in [
|
||||
"pat/pattern_one",
|
||||
"pat/pattern_two",
|
||||
"pat/pattern_three",
|
||||
"pref/pref_one",
|
||||
"pref/pref_two",
|
||||
] {
|
||||
assert!(
|
||||
keys.iter().any(|key| key == expected),
|
||||
"missing {expected}; got {keys:?}"
|
||||
);
|
||||
}
|
||||
assert_eq!(
|
||||
keys.len(),
|
||||
5,
|
||||
"all 5 independent writes must persist; got {keys:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn persist_reflection_writes_to_dedicated_namespace_and_category() {
|
||||
let memory_impl = Arc::new(MockMemory::default());
|
||||
|
||||
Reference in New Issue
Block a user