use anyhow::Context;
use rusqlite::{params, Connection};
use std::io::{Read, Write};
use std::sync::{
atomic::{AtomicBool, AtomicUsize, Ordering},
Arc,
};
use super::super::hybrid_context::query_hybrid_context_memories;
use super::super::policy::{ContextLimits, ContextPolicy};
use super::super::query::{
load_context_data_with_policy, load_context_data_with_policy_local_only,
};
use super::super::sections::render_core_memory_with_limits;
use super::{insert_global_memory, insert_memory, setup_context_schema};
const EMBEDDING_ENV_KEYS: &[&str] = &[
"REMEM_CONFIG",
"REMEM_EMBEDDINGS_PROVIDER",
"REMEM_EMBEDDING_PROVIDER",
"REMEM_EMBEDDINGS_FALLBACK",
"REMEM_EMBEDDINGS_MODEL",
"REMEM_EMBEDDING_MODEL",
"REMEM_EMBEDDINGS_DIMENSIONS",
"REMEM_EMBEDDING_DIMENSIONS",
"REMEM_EMBEDDINGS_API_KEY",
"REMEM_EMBEDDING_API_KEY",
"REMEM_EMBEDDINGS_API_KEY_ENV",
"REMEM_EMBEDDINGS_BASE_URL",
"REMEM_EMBEDDING_BASE_URL",
"REMEM_EMBEDDINGS_TIMEOUT_SECS",
"REMEM_EMBEDDINGS_MODEL_DIR",
"OPENAI_API_KEY",
];
const RERANK_ENV_KEYS: &[&str] = &[
"REMEM_CONFIG",
"REMEM_RERANK_ENABLED",
"REMEM_RERANK_TOP_N",
"REMEM_RERANK_TOP_K",
];
struct ScopedApiFallbackEnv {
_guard: crate::runtime_config::TestEnvGuard,
saved: Vec<(&'static str, Option<String>)>,
}
impl ScopedApiFallbackEnv {
fn new(base_url: &str) -> Self {
let guard = crate::runtime_config::TEST_ENV_LOCK
.lock()
.expect("env lock should acquire");
let saved = EMBEDDING_ENV_KEYS
.iter()
.map(|key| (*key, std::env::var(key).ok()))
.collect::<Vec<_>>();
for key in EMBEDDING_ENV_KEYS {
unsafe { std::env::remove_var(key) };
}
unsafe {
std::env::set_var("REMEM_EMBEDDINGS_PROVIDER", "api");
std::env::set_var("REMEM_EMBEDDINGS_FALLBACK", "feature-hash");
std::env::set_var("REMEM_EMBEDDINGS_API_KEY", "test-key");
std::env::set_var("REMEM_EMBEDDINGS_BASE_URL", base_url);
}
Self {
_guard: guard,
saved,
}
}
fn unavailable_without_fallback() -> Self {
let guard = crate::runtime_config::TEST_ENV_LOCK
.lock()
.expect("env lock should acquire");
let saved = EMBEDDING_ENV_KEYS
.iter()
.map(|key| (*key, std::env::var(key).ok()))
.collect::<Vec<_>>();
for key in EMBEDDING_ENV_KEYS {
unsafe { std::env::remove_var(key) };
}
let isolated_config = std::env::temp_dir().join(format!(
"remem-context-local-only-api-{}-{}.toml",
std::process::id(),
chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default()
));
unsafe {
std::env::set_var("REMEM_CONFIG", isolated_config);
std::env::set_var("REMEM_EMBEDDINGS_PROVIDER", "api");
}
Self {
_guard: guard,
saved,
}
}
}
impl Drop for ScopedApiFallbackEnv {
fn drop(&mut self) {
for (key, value) in self.saved.drain(..) {
match value {
Some(value) => unsafe { std::env::set_var(key, value) },
None => unsafe { std::env::remove_var(key) },
}
}
}
}
struct ScopedInvalidRerankEnv {
_guard: crate::runtime_config::TestEnvGuard,
saved: Vec<(&'static str, Option<String>)>,
}
impl ScopedInvalidRerankEnv {
fn new() -> Self {
let guard = crate::runtime_config::TEST_ENV_LOCK
.lock()
.expect("env lock should acquire");
let saved = RERANK_ENV_KEYS
.iter()
.map(|key| (*key, std::env::var(key).ok()))
.collect::<Vec<_>>();
for key in RERANK_ENV_KEYS {
unsafe { std::env::remove_var(key) };
}
let isolated_config = std::env::temp_dir().join(format!(
"remem-context-local-only-rerank-{}-{}.toml",
std::process::id(),
chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default()
));
unsafe {
std::env::set_var("REMEM_CONFIG", isolated_config);
std::env::set_var("REMEM_RERANK_ENABLED", "true");
std::env::set_var("REMEM_RERANK_TOP_N", "1");
std::env::set_var("REMEM_RERANK_TOP_K", "2");
}
Self {
_guard: guard,
saved,
}
}
}
impl Drop for ScopedInvalidRerankEnv {
fn drop(&mut self) {
for (key, value) in self.saved.drain(..) {
match value {
Some(value) => unsafe { std::env::set_var(key, value) },
None => unsafe { std::env::remove_var(key) },
}
}
}
}
struct FailingEmbeddingServer {
base_url: String,
calls: Arc<AtomicUsize>,
stop: Arc<AtomicBool>,
handle: Option<std::thread::JoinHandle<anyhow::Result<()>>>,
}
impl FailingEmbeddingServer {
fn start(tracked_input: &'static str) -> anyhow::Result<Self> {
let listener = std::net::TcpListener::bind("127.0.0.1:0")?;
listener.set_nonblocking(true)?;
let addr = listener.local_addr()?;
let calls = Arc::new(AtomicUsize::new(0));
let stop = Arc::new(AtomicBool::new(false));
let calls_for_thread = Arc::clone(&calls);
let stop_for_thread = Arc::clone(&stop);
let handle = std::thread::spawn(move || -> anyhow::Result<()> {
while !stop_for_thread.load(Ordering::SeqCst) {
match listener.accept() {
Ok((mut stream, _)) => {
stream.set_nonblocking(false)?;
let mut buffer = [0u8; 8192];
let read = stream.read(&mut buffer)?;
let request = String::from_utf8_lossy(&buffer[..read]);
if request.contains(tracked_input) {
calls_for_thread.fetch_add(1, Ordering::SeqCst);
}
let body = "provider unavailable";
let response = format!(
"HTTP/1.1 500 Internal Server Error\r\ncontent-length: {}\r\n\r\n{}",
body.len(),
body
);
stream.write_all(response.as_bytes())?;
}
Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {
std::thread::sleep(std::time::Duration::from_millis(5));
}
Err(error) => return Err(error.into()),
}
}
Ok(())
});
Ok(Self {
base_url: format!("http://{addr}/v1"),
calls,
stop,
handle: Some(handle),
})
}
fn call_count(&self) -> usize {
self.calls.load(Ordering::SeqCst)
}
}
impl Drop for FailingEmbeddingServer {
fn drop(&mut self) {
self.stop.store(true, Ordering::SeqCst);
if let Some(handle) = self.handle.take() {
handle
.join()
.expect("embedding server thread should not panic")
.expect("embedding server should stop cleanly");
}
}
}
#[test]
fn load_context_data_uses_hybrid_retrieval_from_workstream_signal() {
let conn = Connection::open_in_memory().unwrap();
setup_context_schema(&conn);
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
let limits = ContextLimits {
candidate_fetch_limit: 3,
memory_index_limit: 10,
core_item_limit: 4,
..ContextLimits::default()
};
let policy = ContextPolicy::from_limits(limits);
for idx in 0..20 {
insert_memory(
&conn,
idx + 1,
project,
Some(&format!("recent-noise-{idx}")),
"discovery",
&format!("Recent unrelated note {idx}"),
"Recent context entry without the task terms.",
now - idx,
);
}
insert_memory(
&conn,
200,
project,
Some("sqlcipher-storage-decision"),
"decision",
"SQLCipher storage decision",
"Persist private data with SQLCipher encryption at rest.",
now - 10_000,
);
conn.execute(
"INSERT INTO workstreams
(id, project, title, status, next_action, created_at_epoch, updated_at_epoch)
VALUES (1, ?1, 'Private persistence', 'active',
'Fix SQLCipher recall for private persisted data', ?2, ?2)",
params![project, now],
)
.unwrap();
let loaded = load_context_data_with_policy(&conn, project, None, &policy, true);
assert!(loaded
.memories
.iter()
.any(|memory| memory.title == "SQLCipher storage decision"));
}
#[test]
fn hybrid_context_retrieval_still_excludes_global_non_preferences() {
let conn = Connection::open_in_memory().unwrap();
setup_context_schema(&conn);
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
let limits = ContextLimits {
candidate_fetch_limit: 1,
memory_index_limit: 10,
core_item_limit: 3,
..ContextLimits::default()
};
let policy = ContextPolicy::from_limits(limits);
insert_memory(
&conn,
1,
project,
Some("local-sqlcipher-decision"),
"decision",
"Local SQLCipher decision",
"Repository-local SQLCipher storage decision.",
now - 100,
);
insert_memory(
&conn,
2,
"global",
Some("global-sqlcipher-decision"),
"bugfix",
"Global SQLCipher note",
"Global SQLCipher note should not enter project startup context.",
now,
);
conn.execute(
"UPDATE memories SET scope = 'global', owner_scope = 'user', owner_key = 'manual'
WHERE id = 2",
[],
)
.unwrap();
conn.execute(
"INSERT INTO workstreams
(id, project, title, status, next_action, created_at_epoch, updated_at_epoch)
VALUES (1, ?1, 'SQLCipher recall', 'active',
'Find SQLCipher startup context decision', ?2, ?2)",
params![project, now],
)
.unwrap();
let loaded = load_context_data_with_policy(&conn, project, None, &policy, true);
assert!(loaded
.memories
.iter()
.any(|memory| memory.title == "Local SQLCipher decision"));
assert!(!loaded
.memories
.iter()
.any(|memory| memory.title == "Global SQLCipher note"));
}
#[test]
fn hybrid_context_temporal_retrieval_uses_reference_time() -> anyhow::Result<()> {
let conn = Connection::open_in_memory()?;
setup_context_schema(&conn);
let project = "/tmp/remem";
let wall_clock_epoch = chrono::Utc::now().timestamp();
let event_epoch = 1_600_000_000_i64;
insert_memory(
&conn,
1,
project,
Some("historical-reference-time"),
"decision",
"Historical reference time",
"Episode provenance was captured with a historical event time.",
wall_clock_epoch,
);
conn.execute(
"UPDATE memories
SET reference_time_epoch = ?1
WHERE id = 1",
params![event_epoch],
)?;
let memories = query_hybrid_context_memories(
&conn,
project,
"what happened on 2020-09-13",
None,
&[],
5,
true,
)?;
assert_eq!(memories.len(), 1);
assert_eq!(memories[0].title, "Historical reference time");
Ok(())
}
#[test]
fn hybrid_context_fact_retrieval_labels_validity_in_core_output() -> anyhow::Result<()> {
let conn = Connection::open_in_memory()?;
crate::migrate::run_migrations(&conn)?;
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
let valid_from = chrono::NaiveDate::from_ymd_opt(2026, 1, 2)
.and_then(|date| date.and_hms_opt(12, 0, 0))
.context("valid fact label date")?
.and_utc()
.timestamp();
let limits = ContextLimits {
candidate_fetch_limit: 2,
memory_index_limit: 4,
core_item_limit: 4,
core_char_limit: 1_200,
..ContextLimits::default()
};
let policy = ContextPolicy::from_limits(limits);
for idx in 0..8 {
insert_memory(
&conn,
idx + 1,
project,
Some(&format!("recent-noise-{idx}")),
"session_activity",
&format!("Recent unrelated note {idx}"),
"Recent context entry without the fact terms.",
now - idx,
);
}
insert_memory(
&conn,
100,
project,
Some("harbormint-signer-source"),
"decision",
"HarborMint signer source",
"Signer details live in the temporal fact layer. This memory body is intentionally long enough to exceed the core preview limit before the validity window would appear if the fact label were appended after the body. The rendered context must show temporal facts first so current fact validity is visible.",
now - 10_000,
);
conn.execute(
"INSERT INTO memory_facts
(project, subject, predicate, object, valid_from_epoch, valid_to_epoch,
learned_at_epoch, source_memory_id, source_observation_id, source_event_ids,
confidence, supersedes_fact_id, status, invalidated_at_epoch,
created_at_epoch, updated_at_epoch)
VALUES (?1, 'HarborMint', 'verified_by', 'Toma Reed', ?2, NULL, ?3, 100,
NULL, '[]', 0.95, NULL, 'active', NULL, ?3, ?3)",
params![project, valid_from, now - 9_000],
)?;
conn.execute(
"INSERT INTO memory_facts
(project, subject, predicate, object, valid_from_epoch, valid_to_epoch,
learned_at_epoch, source_memory_id, source_observation_id, source_event_ids,
confidence, supersedes_fact_id, status, invalidated_at_epoch,
created_at_epoch, updated_at_epoch)
VALUES
(?1, 'UnrelatedService', 'verified_by', 'Mira Lane', ?2, NULL, ?3, 100,
NULL, '[]', 0.95, NULL, 'active', NULL, ?3, ?3),
(?1, 'UnrelatedService', 'blocked_by', 'North Region', ?2, NULL, ?3, 100,
NULL, '[]', 0.95, NULL, 'active', NULL, ?3, ?3)",
params![project, now - 500, now - 400],
)?;
conn.execute(
"INSERT INTO workstreams
(id, project, title, status, next_action, created_at_epoch, updated_at_epoch)
VALUES (1, ?1, 'HarborMint signer', 'active',
'Who signs HarborMint with Toma Reed?', ?2, ?2)",
params![project, now],
)?;
let loaded = load_context_data_with_policy(&conn, project, None, &policy, true);
let memory = loaded
.memories
.iter()
.find(|memory| memory.id == 100)
.context("fact channel should retrieve source memory")?;
assert!(memory.text.contains("Temporal facts:"));
assert!(memory.text.contains("Toma Reed"));
assert!(memory.text.contains("valid_from=2026-01-02"));
assert!(memory.text.contains("valid_to=open"));
let mut output = String::new();
render_core_memory_with_limits(&mut output, &loaded.memories, &limits);
assert!(output.contains("HarborMint signer source"), "{output}");
assert!(output.contains("Temporal facts:"), "{output}");
assert!(output.contains("valid_from=2026-01-02"), "{output}");
Ok(())
}
#[test]
fn hybrid_context_fact_retrieval_filters_excluded_types_before_ranking() -> anyhow::Result<()> {
let conn = Connection::open_in_memory()?;
crate::migrate::run_migrations(&conn)?;
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
for (id, memory_type, title, age) in [
(1, "preference", "Preference fact source", 100),
(2, "lesson", "Lesson fact source", 200),
(3, "decision", "Decision fact source", 300),
] {
insert_memory(
&conn,
id,
project,
Some(title),
memory_type,
title,
"Opaque source body without signer terms.",
now - age,
);
conn.execute(
"INSERT INTO memory_facts
(project, subject, predicate, object, valid_from_epoch, valid_to_epoch,
learned_at_epoch, source_memory_id, source_observation_id, source_event_ids,
confidence, supersedes_fact_id, status, invalidated_at_epoch,
created_at_epoch, updated_at_epoch)
VALUES (?1, 'HarborMint', 'verified_by', 'Toma Reed', ?2, NULL, ?3, ?4,
NULL, '[]', 0.95, NULL, 'active', NULL, ?3, ?3)",
params![project, now - 1_000, now - age, id],
)?;
}
let memories = query_hybrid_context_memories(
&conn,
project,
"Who signs HarborMint with Toma Reed?",
None,
&["preference", "lesson"],
1,
true,
)?;
assert_eq!(memories.len(), 1);
assert_eq!(memories[0].title, "Decision fact source");
Ok(())
}
#[test]
fn hybrid_context_fact_retrieval_applies_owner_filter_before_ranking() -> anyhow::Result<()> {
let conn = Connection::open_in_memory()?;
crate::migrate::run_migrations(&conn)?;
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
for id in 1..=20 {
insert_memory(
&conn,
id,
project,
Some(&format!("foreign-owner-fact-{id}")),
"decision",
&format!("Foreign owner fact {id}"),
"Opaque source body without signer terms.",
now - id,
);
conn.execute(
"UPDATE memories SET owner_scope = 'repo', owner_key = '/tmp/other' WHERE id = ?1",
params![id],
)?;
conn.execute(
"INSERT INTO memory_facts
(project, subject, predicate, object, valid_from_epoch, valid_to_epoch,
learned_at_epoch, source_memory_id, source_observation_id, source_event_ids,
confidence, supersedes_fact_id, status, invalidated_at_epoch,
created_at_epoch, updated_at_epoch)
VALUES (?1, 'HarborMint', 'verified_by', 'Toma Reed', ?2, NULL, ?3, ?4,
NULL, '[]', 0.95, NULL, 'active', NULL, ?3, ?3)",
params![project, now - 1_000, now - id, id],
)?;
}
insert_memory(
&conn,
100,
project,
Some("repo-owner-fact"),
"decision",
"Repo owner fact",
"Opaque source body without signer terms.",
now - 10_000,
);
conn.execute(
"INSERT INTO memory_facts
(project, subject, predicate, object, valid_from_epoch, valid_to_epoch,
learned_at_epoch, source_memory_id, source_observation_id, source_event_ids,
confidence, supersedes_fact_id, status, invalidated_at_epoch,
created_at_epoch, updated_at_epoch)
VALUES (?1, 'HarborMint', 'verified_by', 'Toma Reed', ?2, NULL, ?3, 100,
NULL, '[]', 0.95, NULL, 'active', NULL, ?3, ?3)",
params![project, now - 1_000, now - 10_000],
)?;
let memories = query_hybrid_context_memories(
&conn,
project,
"Who signs HarborMint with Toma Reed?",
None,
&[],
1,
true,
)?;
assert_eq!(memories.len(), 1);
assert_eq!(memories[0].title, "Repo owner fact");
Ok(())
}
#[test]
fn hybrid_context_vector_recall_is_not_crowded_out_by_global_hits() -> anyhow::Result<()> {
let conn = Connection::open_in_memory()?;
setup_context_schema(&conn);
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
let limits = ContextLimits {
candidate_fetch_limit: 1,
memory_index_limit: 10,
core_item_limit: 3,
..ContextLimits::default()
};
let policy = ContextPolicy::from_limits(limits);
insert_memory(
&conn,
1,
project,
Some("credential-store"),
"architecture",
"Credential store",
"SQLCipher encrypts secrets at rest.",
now - 10_000,
);
crate::retrieval::vector::upsert_memory_embedding_for_row(&conn, 1)?;
for idx in 0..30 {
let id = idx + 2;
insert_global_memory(
&conn,
id,
"global",
Some(&format!("global-private-data-{idx}")),
"bugfix",
&format!("Global private data note {idx}"),
"Protect private persisted data with a global-only diagnostic note.",
now + idx,
);
crate::retrieval::vector::upsert_memory_embedding_for_row(&conn, id)?;
}
conn.execute(
"INSERT INTO workstreams
(id, project, title, status, next_action, created_at_epoch, updated_at_epoch)
VALUES (1, ?1, 'Private persistence', 'active',
'How do we protect private persisted data?', ?2, ?2)",
params![project, now],
)?;
let loaded = load_context_data_with_policy(&conn, project, None, &policy, true);
assert!(loaded
.memories
.iter()
.any(|memory| memory.title == "Credential store"));
assert!(!loaded
.memories
.iter()
.any(|memory| memory.title.starts_with("Global private data note")));
Ok(())
}
#[test]
fn hybrid_context_vector_channel_uses_fallback_without_status_probe() -> anyhow::Result<()> {
let server = FailingEmbeddingServer::start(r#""input":"SQLCipher encrypts secrets""#)?;
let _env = ScopedApiFallbackEnv::new(&server.base_url);
let conn = Connection::open_in_memory()?;
setup_context_schema(&conn);
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
insert_memory(
&conn,
1,
project,
Some("credential-store"),
"architecture",
"Credential store",
"SQLCipher encrypts secrets at rest.",
now,
);
crate::retrieval::vector::ensure_vec_table(&conn)?;
let embedding = crate::retrieval::vector::embed_memory_text(
"Credential store",
"SQLCipher encrypts secrets at rest.",
"architecture",
Some("credential-store"),
);
crate::retrieval::vector::upsert_embedding(&conn, 1, &embedding)?;
let memories = query_hybrid_context_memories(
&conn,
project,
"SQLCipher encrypts secrets",
None,
&[],
5,
true,
)?;
assert!(memories
.iter()
.any(|memory| memory.title == "Credential store"));
assert_eq!(server.call_count(), 1);
Ok(())
}
#[test]
fn local_only_hybrid_context_never_calls_configured_api_provider() -> anyhow::Result<()> {
let server = FailingEmbeddingServer::start(r#""input":"SQLCipher encrypts secrets""#)?;
let _env = ScopedApiFallbackEnv::new(&server.base_url);
let conn = Connection::open_in_memory()?;
setup_context_schema(&conn);
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
insert_memory(
&conn,
1,
project,
Some("credential-store"),
"architecture",
"Credential store",
"SQLCipher encrypts secrets at rest.",
now,
);
crate::retrieval::vector::ensure_vec_table(&conn)?;
let embedding = crate::retrieval::vector::embed_memory_text(
"Credential store",
"SQLCipher encrypts secrets at rest.",
"architecture",
Some("credential-store"),
);
crate::retrieval::vector::upsert_embedding(&conn, 1, &embedding)?;
let memories = query_hybrid_context_memories(
&conn,
project,
"SQLCipher encrypts secrets",
None,
&[],
5,
false,
)?;
assert!(memories
.iter()
.any(|memory| memory.title == "Credential store"));
assert_eq!(server.call_count(), 0);
Ok(())
}
#[test]
fn local_only_context_loader_never_resolves_or_applies_rerank() -> anyhow::Result<()> {
let _env = ScopedInvalidRerankEnv::new();
let conn = Connection::open_in_memory()?;
setup_context_schema(&conn);
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
insert_memory(
&conn,
1,
project,
Some("rerank-policy"),
"architecture",
"Context bundle retrieval policy",
"The experimental bundle uses its planned canonical order.",
now,
);
let policy = ContextPolicy::from_limits(ContextLimits::default());
let local_only = load_context_data_with_policy_local_only(&conn, project, None, &policy, false);
assert!(local_only.rerank.is_none());
assert!(local_only
.errors
.iter()
.all(|error| error.section != "rerank"));
let default = load_context_data_with_policy(&conn, project, None, &policy, false);
assert!(default.errors.iter().any(|error| error.section == "rerank"));
Ok(())
}
#[test]
fn local_only_context_loader_keeps_lexical_channels_when_api_key_is_missing() -> anyhow::Result<()>
{
let _env = ScopedApiFallbackEnv::unavailable_without_fallback();
let conn = Connection::open_in_memory()?;
setup_context_schema(&conn);
let project = "/tmp/remem";
let now = chrono::Utc::now().timestamp();
insert_memory(
&conn,
1,
project,
Some("sqlcipher-lexical"),
"architecture",
"SQLCipher lexical recovery",
"Use SQLCipher for encrypted local persistence.",
now,
);
conn.execute(
"INSERT INTO workstreams
(id, project, title, status, next_action, created_at_epoch, updated_at_epoch)
VALUES (1, ?1, 'SQLCipher recovery', 'active',
'Fix SQLCipher encrypted local persistence', ?2, ?2)",
params![project, now],
)?;
let policy = ContextPolicy::from_limits(ContextLimits::default());
let loaded = load_context_data_with_policy_local_only(&conn, project, None, &policy, false);
assert!(loaded.errors.is_empty(), "{:?}", loaded.errors);
assert!(loaded.memories.iter().any(|memory| memory.id == 1));
Ok(())
}