use super::types::{SaveMemoryNextStep, SaveMemoryRequest, SaveMemoryResult};
use crate::memory::activation::SupplementalSaveReceipt;
use crate::memory::claims::{claims_enabled, insert_memory_claim, ClaimWriteRequest};
use crate::memory::lesson::SaveLessonRequest;
use crate::memory::lifecycle::MemoryLifecycleOp;
use crate::memory::poisoning::{
scan_instruction_pattern, InstructionPatternMatch, SourceTrustClass,
};
use crate::memory::{MemoryType, MEMORY_TYPES};
use anyhow::{anyhow, Context, Result};
use rusqlite::Connection;
mod activation;
mod errors;
mod local_copy_plan;
pub use activation::SaveMemoryCaller;
pub use errors::{LocalCopyError, SaveMemoryIdempotencyConflictError, SaveMemoryValidationError};
use local_copy_plan::{
cleanup_local_copy, discard_local_copy_backup, prepare_local_copy, replay_local_copy,
write_local_copy,
};
pub fn save_memory(conn: &Connection, req: &SaveMemoryRequest) -> Result<SaveMemoryResult> {
save_memory_from_with_reference_time(conn, req, req.created_at_epoch, SaveMemoryCaller::RustApi)
}
pub(crate) fn save_memory_for_benchmark_fixture(
conn: &Connection,
req: &SaveMemoryRequest,
) -> Result<SaveMemoryResult> {
let reference_time_epoch = normalize_reference_time_epoch(req.created_at_epoch, None);
save_memory_inner(
conn,
req,
reference_time_epoch,
SaveMemoryCaller::RustApi,
Some(SourceTrustClass::UserPrompt),
)
}
pub fn save_memory_with_reference_time(
conn: &Connection,
req: &SaveMemoryRequest,
reference_time_epoch: Option<i64>,
) -> Result<SaveMemoryResult> {
save_memory_from_with_reference_time(conn, req, reference_time_epoch, SaveMemoryCaller::RustApi)
}
pub fn save_memory_from_with_reference_time(
conn: &Connection,
req: &SaveMemoryRequest,
reference_time_epoch: Option<i64>,
caller: SaveMemoryCaller,
) -> Result<SaveMemoryResult> {
let reference_time_epoch =
normalize_reference_time_epoch(req.created_at_epoch, reference_time_epoch);
save_memory_inner(conn, req, reference_time_epoch, caller, None)
}
fn normalize_reference_time_epoch(
created_at_epoch: Option<i64>,
reference_time_epoch: Option<i64>,
) -> Option<i64> {
reference_time_epoch.or(created_at_epoch)
}
fn save_memory_inner(
conn: &Connection,
req: &SaveMemoryRequest,
reference_time_epoch: Option<i64>,
caller: SaveMemoryCaller,
source_trust_override: Option<SourceTrustClass>,
) -> Result<SaveMemoryResult> {
let validated = validate_save_memory_request(req)?;
let project = req.project.as_deref().unwrap_or("manual");
let title = req.title.as_deref().unwrap_or("Memory");
let memory_type = validated.memory_type.as_str();
let source_trust = source_trust_override.unwrap_or_else(|| caller.source_trust());
let files_json = req
.files
.as_ref()
.and_then(|files| serde_json::to_string(files).ok());
let scope = validated.scope.as_str();
let effective_topic_key = effective_topic_key(req, memory_type);
let (operation_input, operation_plan) = crate::memory::operation::plan_direct_save(
conn,
"direct",
"save_memory",
project,
scope,
memory_type,
effective_topic_key.as_deref(),
title,
&req.text,
files_json.as_deref(),
req.branch.as_deref(),
None,
None,
)?;
let mut local_copy = None;
let mut activation_request = activation::build_request(
req,
caller,
project,
title,
memory_type,
scope,
effective_topic_key.as_deref(),
reference_time_epoch,
caller == SaveMemoryCaller::RustApi
&& req
.acknowledge_pattern
.as_deref()
.is_some_and(|pattern| !pattern.trim().is_empty()),
);
activation_request.source_trust = source_trust;
activation_request.result_source_trust = source_trust;
if source_trust_override.is_some() {
activation_request.actor_kind = crate::memory::activation::ActivationActorKind::Operator;
activation_request.source_operation = "coding_bench_fixture_seed".to_string();
activation_request.provenance_ref = "coding-bench:curated-fixture".to_string();
}
let early_replay =
crate::memory::activation::replay_supplemental_if_present(conn, &activation_request)
.map_err(|err| -> anyhow::Error {
if err.is::<crate::memory::activation::ActivationIdConflictError>() {
SaveMemoryIdempotencyConflictError::new(err.to_string()).into()
} else {
err
}
})?;
let preserves_existing_provenance = if early_replay.is_some() {
true
} else {
activation::bind_existing_target_provenance(
conn,
&mut activation_request,
&operation_plan,
memory_type,
&req.text,
)?
};
if early_replay.is_none()
&& memory_type != "lesson"
&& operation_plan.op == MemoryLifecycleOp::Noop
{
let memory_id = operation_plan
.target_memory_id
.ok_or_else(|| anyhow!("noop memory operation missing existing memory id"))?;
let stored_trust: String = conn.query_row(
"SELECT source_trust_class FROM memories WHERE id = ?1",
[memory_id],
|row| row.get(0),
)?;
activation_request.result_source_trust = SourceTrustClass::parse(&stored_trust)
.ok_or_else(|| {
anyhow!("existing memory has invalid source trust class: {stored_trust}")
})?;
}
let mut applied_operation = None;
let save_result = if let Some(replayed) = early_replay {
Ok(replayed)
} else {
crate::memory::activation::execute_supplemental_save(conn, &activation_request, |permit| {
let acknowledgement = direct_save_pattern_acknowledgement(
title,
&req.text,
req.acknowledge_pattern.as_deref(),
caller,
)?;
let mut prepared_local_copy =
prepare_local_copy(project, title, req).map_err(LocalCopyError::from)?;
write_local_copy(&mut prepared_local_copy).map_err(LocalCopyError::from)?;
local_copy = Some(prepared_local_copy);
let local_copy_receipt = local_copy
.as_mut()
.context("fresh supplemental save omitted its local-copy plan")?
.receipt()?;
let result = if memory_type == "lesson" {
crate::memory::operation::with_operation_savepoint(conn, || {
let id = crate::memory::lesson::save_lesson_with_reference_time_activated(
conn,
permit,
&SaveLessonRequest {
session_id: req.session_id.as_deref(),
project,
topic_key: req.topic_key.as_deref(),
title,
content: &req.text,
confidence: 0.7,
source_evidence: None,
files: files_json.as_deref(),
branch: req.branch.as_deref(),
scope,
created_at_epoch: req.created_at_epoch,
stale_after_epoch: None,
},
reference_time_epoch,
)?;
let mut logged_plan = operation_plan.clone();
logged_plan.target_memory_id = Some(id);
if logged_plan.op == MemoryLifecycleOp::Noop {
logged_plan.op = MemoryLifecycleOp::Update;
logged_plan.noop_reason = None;
logged_plan.reason =
"existing lesson memory was reinforced by direct save".to_string();
}
crate::memory::operation::insert_operation_log(
conn,
&operation_input,
&logged_plan,
Some(id),
)?;
mark_direct_save_poisoning_metadata(
conn,
id,
acknowledgement,
source_trust,
true,
)?;
Ok((id, logged_plan.op))
})
} else {
crate::memory::operation::with_operation_savepoint(conn, || {
if operation_plan.op == MemoryLifecycleOp::Noop {
let id = operation_plan.target_memory_id.ok_or_else(|| {
anyhow!("noop memory operation missing existing memory id")
})?;
crate::memory::operation::insert_operation_log(
conn,
&operation_input,
&operation_plan,
Some(id),
)?;
mark_direct_save_poisoning_metadata(
conn,
id,
acknowledgement,
source_trust,
false,
)?;
return Ok((id, MemoryLifecycleOp::Noop));
}
let previous_preference =
if memory_type == "preference" {
operation_plan
.target_memory_id
.map(|memory_id| {
conn.query_row(
"SELECT content FROM memories WHERE id = ?1",
[memory_id],
|row| row.get::<_, String>(0),
)
.with_context(|| {
format!("load preference before direct save update id={memory_id}")
})
})
.transpose()?
} else {
None
};
let result =
crate::memory::store::insert_memory_full_with_operation_log_activated(
conn,
permit,
req.session_id.as_deref(),
project,
req.topic_key.as_deref(),
title,
&req.text,
memory_type,
files_json.as_deref(),
req.branch.as_deref(),
scope,
req.created_at_epoch,
reference_time_epoch,
&operation_input,
&operation_plan,
)?;
if let Some(previous_text) = previous_preference {
crate::memory::preference::compilation::enqueue_for_memory_ids(
conn,
&[result.0],
)?;
crate::memory::preference::reinforcement::reconcile_in_place_preference_update(
conn,
result.0,
&previous_text,
&req.text,
)?;
}
mark_direct_save_poisoning_metadata(
conn,
result.0,
acknowledgement,
source_trust,
true,
)?;
Ok(result)
})
}?;
if !preserves_existing_provenance {
conn.execute(
"UPDATE memories
SET evidence_event_ids = NULL, source_candidate_id = NULL,
confidence = NULL
WHERE id = ?1",
[result.0],
)?;
}
applied_operation = Some(result.1);
let claim_receipt = write_claim_after_durable_save(conn, result.0, req)?;
Ok((result.0, claim_receipt, local_copy_receipt))
})
};
let activation = match save_result {
Ok(result) => result,
Err(err) => {
let err: anyhow::Error =
if err.is::<crate::memory::activation::ActivationIdConflictError>() {
SaveMemoryIdempotencyConflictError::new(err.to_string()).into()
} else {
err
};
if let Some(local_copy) = local_copy.as_ref() {
if let Err(cleanup_err) = cleanup_local_copy(local_copy) {
return Err(err.context(format!(
"database save failed and local copy cleanup failed: {cleanup_err}"
)));
}
}
return Err(err);
}
};
if activation.replayed {
let local_copy_receipt = activation
.supplemental_local_copy_receipt
.as_ref()
.ok_or_else(|| anyhow!("supplemental replay omitted its local-copy receipt"))?;
local_copy = Some(
replay_local_copy(project, title, &req.text, local_copy_receipt)
.map_err(LocalCopyError::from)?,
);
}
let id = activation.memory_id;
let operation = if activation.replayed {
MemoryLifecycleOp::Noop
} else {
applied_operation.unwrap_or(MemoryLifecycleOp::Noop)
};
let claim_result = activation
.supplemental_receipt
.ok_or_else(|| anyhow!("supplemental save completed without a durable claim receipt"))?;
let local_copy = local_copy
.as_ref()
.ok_or_else(|| anyhow!("supplemental save completed without a local-copy plan"))?;
discard_local_copy_backup(local_copy);
let durable = load_durable_write_details(conn, id)?;
let local_copy_result = local_copy.result();
Ok(SaveMemoryResult {
id,
status: "saved".to_string(),
memory_type: durable.memory_type,
project: durable.project,
scope: durable.scope,
topic_key: durable.topic_key,
branch: durable.branch,
operation: operation.as_str().to_string(),
created_at_epoch: durable.created_at_epoch,
reference_time_epoch: durable.reference_time_epoch,
updated_at_epoch: durable.updated_at_epoch,
upserted: req.topic_key.is_some(),
local_status: local_copy_result.status.clone(),
local_path: local_copy_result.path.clone(),
local_copy: local_copy_result,
claim_status: claim_result.status().to_string(),
claim_id: claim_result.claim_id(),
claim_error: claim_result.error().map(str::to_string),
next_step: SaveMemoryNextStep {
tool: "get_observations".to_string(),
ids: vec![id],
source: "memory".to_string(),
reason: format!(
"Verify the durable memory with get_observations(ids=[{id}], source='memory') or search(project='{}').",
durable_project_hint(project)
),
},
})
}
fn direct_save_pattern_acknowledgement(
title: &str,
text: &str,
acknowledged_pattern_id: Option<&str>,
caller: SaveMemoryCaller,
) -> Result<Option<InstructionPatternMatch>> {
let scan_text = format!("{title}\n{text}");
let matched = scan_instruction_pattern(&scan_text);
let acknowledged_pattern_id = acknowledged_pattern_id
.map(str::trim)
.filter(|value| !value.is_empty());
if matches!(
caller,
SaveMemoryCaller::McpAgent | SaveMemoryCaller::RestAgent
) && acknowledged_pattern_id.is_some()
{
return Err(SaveMemoryValidationError::new(
"agent-facing save_memory cannot issue human instruction-pattern acknowledgement; use a separately reviewed governance surface",
)
.into());
}
match (matched, acknowledged_pattern_id) {
(Some(matched), Some(acknowledged)) if acknowledged == matched.pattern_id => {
Ok(Some(matched))
}
(Some(matched), Some(acknowledged)) => {
Err(SaveMemoryValidationError::new(format!(
"save_memory acknowledged pattern {acknowledged} does not match instruction-pattern {}@v{}",
matched.pattern_id, matched.pattern_set_version
))
.into())
}
(Some(matched), None) => {
let review_surface = match caller {
SaveMemoryCaller::McpAgent | SaveMemoryCaller::RestAgent => {
"candidate review governance"
}
SaveMemoryCaller::RustApi => "a separately reviewed governance surface",
};
Err(SaveMemoryValidationError::new(format!(
"save_memory text matched instruction-pattern {}@v{}; use {review_surface} to review and acknowledge the pattern before saving",
matched.pattern_id, matched.pattern_set_version
))
.into())
}
(None, Some(acknowledged)) => Err(SaveMemoryValidationError::new(format!(
"save_memory acknowledge_pattern {acknowledged} was provided, but no instruction-pattern matched"
))
.into()),
(None, None) => Ok(None),
}
}
fn mark_direct_save_poisoning_metadata(
conn: &Connection,
memory_id: i64,
acknowledgement: Option<InstructionPatternMatch>,
source_trust: SourceTrustClass,
update_source_trust: bool,
) -> Result<()> {
if let Some(acknowledgement) = acknowledgement {
let now = chrono::Utc::now().timestamp();
if update_source_trust {
conn.execute(
"UPDATE memories
SET source_trust_class = ?1,
acknowledged_pattern_id = ?2,
acknowledged_pattern_version = ?3,
acknowledged_at_epoch = ?4
WHERE id = ?5",
rusqlite::params![
source_trust.as_str(),
acknowledgement.pattern_id,
acknowledgement.pattern_set_version,
now,
memory_id
],
)?;
} else {
conn.execute(
"UPDATE memories
SET acknowledged_pattern_id = ?1,
acknowledged_pattern_version = ?2,
acknowledged_at_epoch = ?3
WHERE id = ?4",
rusqlite::params![
acknowledgement.pattern_id,
acknowledgement.pattern_set_version,
now,
memory_id
],
)?;
}
} else if update_source_trust {
conn.execute(
"UPDATE memories SET source_trust_class = ?1 WHERE id = ?2",
rusqlite::params![source_trust.as_str(), memory_id],
)?;
}
Ok(())
}
struct ValidatedSaveMemoryRequest {
memory_type: String,
scope: String,
}
fn validate_save_memory_request(req: &SaveMemoryRequest) -> Result<ValidatedSaveMemoryRequest> {
if req.text.trim().is_empty() {
return Err(SaveMemoryValidationError::new("save_memory text must not be blank").into());
}
if req
.idempotency_key
.as_deref()
.is_some_and(|key| key.trim().is_empty() || key.contains('\0') || key.len() > 256)
{
return Err(SaveMemoryValidationError::new(
"save_memory idempotency_key must be nonblank, contain no NUL, and be at most 256 bytes",
)
.into());
}
let memory_type = match req.memory_type.as_deref() {
Some(value) => {
let normalized = value.trim().to_ascii_lowercase();
let parsed = MemoryType::parse(&normalized).ok_or_else(|| {
SaveMemoryValidationError::new(format!(
"save_memory memory_type must be one of: {}",
MEMORY_TYPES.join(", ")
))
})?;
parsed.as_str().to_string()
}
None => MemoryType::Discovery.as_str().to_string(),
};
let scope = match req.scope.as_deref() {
Some(value) => match value.trim().to_ascii_lowercase().as_str() {
"project" => "project".to_string(),
"global" => "global".to_string(),
_ => {
return Err(SaveMemoryValidationError::new(
"save_memory scope must be one of: project, global",
)
.into());
}
},
None => "project".to_string(),
};
Ok(ValidatedSaveMemoryRequest { memory_type, scope })
}
fn write_claim_after_durable_save(
conn: &Connection,
memory_id: i64,
req: &SaveMemoryRequest,
) -> Result<SupplementalSaveReceipt> {
if !claims_enabled(req.claim_enabled) {
return Ok(SupplementalSaveReceipt::Disabled);
}
let claim_source = req.claim_source.as_deref().unwrap_or("manual_save");
match insert_memory_claim(
conn,
&ClaimWriteRequest {
memory_id,
session_id: req.session_id.as_deref(),
host: req.host.as_deref(),
claim_source,
},
) {
Ok(claim_id) => SupplementalSaveReceipt::saved(claim_id),
Err(err) => {
let error = format!("{err:#}");
crate::log::error(
"memory-claim",
&format!("claim write failed memory_id={} error={}", memory_id, error),
);
SupplementalSaveReceipt::failed(error)
}
}
}
fn effective_topic_key(req: &SaveMemoryRequest, memory_type: &str) -> Option<String> {
if memory_type == "lesson" {
return req.topic_key.clone().or_else(|| {
Some(format!(
"lesson-{}",
crate::memory::slugify_for_topic(&req.text, 64)
))
});
}
req.topic_key.clone()
}
struct DurableWriteDetails {
project: String,
scope: String,
topic_key: Option<String>,
branch: Option<String>,
memory_type: String,
created_at_epoch: i64,
reference_time_epoch: i64,
updated_at_epoch: i64,
}
fn load_durable_write_details(conn: &Connection, id: i64) -> Result<DurableWriteDetails> {
conn.query_row(
"SELECT project, COALESCE(scope, 'project'), topic_key, branch, memory_type,
created_at_epoch, COALESCE(reference_time_epoch, created_at_epoch), updated_at_epoch
FROM memories
WHERE id = ?1",
[id],
|row| {
Ok(DurableWriteDetails {
project: row.get(0)?,
scope: row.get(1)?,
topic_key: row.get(2)?,
branch: row.get(3)?,
memory_type: row.get(4)?,
created_at_epoch: row.get(5)?,
reference_time_epoch: row.get(6)?,
updated_at_epoch: row.get(7)?,
})
},
)
.with_context(|| format!("load durable write details for memory {id}"))
}
fn durable_project_hint(project: &str) -> String {
project.replace('\'', "\\'")
}