use super::handle::PalaceHandle;
use super::tier_c;
use super::types::RememberOptions;
use crate::memory_core::filter::{FilterReject, check_secret, classify};
use crate::memory_core::palace::{Drawer, DrawerType, RoomType};
use crate::memory_core::room_identity::DEFAULT_WING_ID;
use crate::memory_core::store::l1_cache::L1Cache;
use crate::memory_core::store::rooms::resolve_or_create_room_in_wing;
use crate::memory_core::store::vector::VectorStore;
use crate::memory_core::timeouts;
use anyhow::{Context, Result};
use std::time::{Duration, Instant};
use uuid::Uuid;
const KG_STORE_FILENAME: &str = "kg.redb";
pub(super) async fn remember_within(
handle: &PalaceHandle,
content: String,
room: RoomType,
tags: Vec<String>,
importance: f32,
opts: RememberOptions,
budget: Duration,
) -> Result<Uuid> {
handle.touch();
if handle.is_read_only() {
return Err(anyhow::anyhow!(
"palace '{}' is read-only: HTTP daemon holds the write lock — \
route writes through the daemon's HTTP API or stop the daemon \
before retrying via stdio",
handle.id
));
}
let _write_guard = timeouts::lock_with_timeout(
&handle.write_mutex,
timeouts::write_lock_timeout(),
handle.id.as_str(),
)
.await?;
if budget.is_zero() {
return Err(over_budget_error(handle, budget, Duration::ZERO));
}
let started = Instant::now();
let outcome = tokio::time::timeout(
budget,
run_pipeline(handle, content, room, tags, importance, opts),
)
.await;
let elapsed = started.elapsed();
match outcome {
Ok(result) => {
if elapsed >= timeouts::slow_write_warn_threshold() {
tracing::warn!(
palace = %handle.id,
elapsed_ms = elapsed.as_millis(),
budget_ms = budget.as_millis(),
kg_redb_bytes = kg_store_bytes(handle),
"#6366: write held the palace write mutex far longer than \
a write should; every other writer on this palace waited \
for it. A large kg.redb makes commits slower — consider \
compacting or splitting this palace"
);
}
result
}
Err(_) => Err(over_budget_error(handle, budget, elapsed)),
}
}
fn over_budget_error(handle: &PalaceHandle, budget: Duration, elapsed: Duration) -> anyhow::Error {
anyhow::anyhow!(
"palace '{}' write pipeline exceeded its {:?} budget after {:?} \
(issue #6366); the palace write mutex has been released so other \
writers proceed. kg.redb is {} bytes — a large store makes commits \
slower; raise TRUSTY_WRITE_PIPELINE_TIMEOUT_SECS if writes on this \
palace are legitimately this slow",
handle.id,
budget,
elapsed,
kg_store_bytes(handle).map_or_else(|| "unknown".to_string(), |bytes| bytes.to_string()),
)
}
fn kg_store_bytes(handle: &PalaceHandle) -> Option<u64> {
let data_dir = handle.data_dir.as_ref()?;
std::fs::metadata(data_dir.join(KG_STORE_FILENAME))
.ok()
.map(|m| m.len())
}
async fn run_pipeline(
handle: &PalaceHandle,
content: String,
room: RoomType,
tags: Vec<String>,
importance: f32,
opts: RememberOptions,
) -> Result<Uuid> {
if !opts.force {
opts.filter
.apply(&content, opts.enforce_min_tokens)
.map_err(|reject| match reject {
FilterReject::TooShort { .. }
| FilterReject::NoisePattern { .. }
| FilterReject::NonAlphabetic { .. }
| FilterReject::PotentialSecret { .. } => anyhow::anyhow!("{reject}"),
})?;
} else if !opts.allow_secret_like {
check_secret(&content).map_err(|reject: FilterReject| anyhow::anyhow!("{reject}"))?;
}
let wing_id = opts.wing_id.unwrap_or(DEFAULT_WING_ID);
let room_id = resolve_or_create_room_in_wing(&handle.kg, &room, wing_id).await;
let mut drawer = Drawer::new(room_id, content.clone());
drawer.tags = tags;
drawer.importance = importance.clamp(0.0, 1.0);
let final_type = match opts.classify_as {
Some(t) => t,
None => classify(&content, DrawerType::Unknown),
};
drawer = drawer.with_type(final_type);
tier_c::apply_admission(&mut drawer, &opts, &handle.id);
let id = drawer.id;
if !opts.defer_embedding {
let embedder = super::embedder::shared_embedder()
.await
.context("acquire shared embedder for remember")?;
let embed_timeout = timeouts::embed_batch_timeout();
let vecs = tokio::time::timeout(
embed_timeout,
embedder.embed_batch(std::slice::from_ref(&content)),
)
.await
.map_err(|_| {
anyhow::anyhow!(
"embed_batch timed out after {:?} on remember path (issue #906); \
increase TRUSTY_EMBED_BATCH_TIMEOUT_SECS if batches legitimately \
take longer on this host",
embed_timeout
)
})?
.context("embed drawer content")?;
if let Some(v) = vecs.into_iter().next() {
handle
.vector_store
.upsert(id, v)
.await
.context("upsert drawer vector")?;
}
}
let order = handle.commit_mutex.clone().lock_owned().await;
tokio::spawn(tier_c::commit_and_mirror(
handle.commit_ctx(),
drawer,
order,
))
.await
.context("write pipeline commit task join")??;
if opts.defer_embedding {
handle.spawn_deferred_embed(id, content);
}
if let Some(data_dir) = handle.data_dir.as_ref() {
let snap = handle.drawers.read().clone();
L1Cache::save_l1_cache(&snap, data_dir).context("save L1 snapshot")?;
}
handle.rebuild_closets();
Ok(id)
}