use anyhow::{anyhow, bail, Result};
use super::candidates::{Cluster, MemoryCandidate};
use super::constants::DREAM_PROMPT;
#[derive(Debug)]
pub(super) enum MergeDecision {
Merge(MergeResult),
NoMerge {
reason: Option<String>,
},
Conflict {
conflicting_ids: Vec<i64>,
reason: Option<String>,
},
}
#[derive(Debug)]
pub(super) struct MergeResult {
pub topic_key: String,
pub memory_type: String,
pub title: String,
pub content: String,
pub superseded_ids: Vec<i64>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DecisionTag {
Conflict,
NoMerge,
Memory,
}
pub(super) async fn merge_cluster(
cluster: &Cluster,
project: &str,
host: Option<String>,
profile: Option<String>,
) -> Result<MergeDecision> {
let user_message = build_user_message(&cluster.members);
let response = crate::ai::call_ai(
DREAM_PROMPT,
&user_message,
crate::ai::UsageContext {
project: Some(project),
session_id: None,
operation: "dream",
host: profile.is_none().then_some(host.as_deref()).flatten(),
profile: profile.as_deref(),
},
)
.await?;
filter_superseded_ids(parse_response(&response)?, cluster)
}
fn filter_superseded_ids(decision: MergeDecision, cluster: &Cluster) -> Result<MergeDecision> {
let member_ids: std::collections::HashSet<i64> = cluster.members.iter().map(|m| m.id).collect();
match decision {
MergeDecision::Merge(mut result) => {
let before = result.superseded_ids.len();
result.superseded_ids.retain(|id| member_ids.contains(id));
let dropped = before - result.superseded_ids.len();
if dropped > 0 {
crate::log::warn(
"dream",
&format!(
"dropped {} hallucinated superseded_id(s) not in cluster",
dropped
),
);
}
if result.superseded_ids.is_empty() {
crate::log::warn(
"dream",
"rejecting merge with no valid superseded_id(s) after filtering",
);
return Ok(MergeDecision::NoMerge {
reason: Some("no valid superseded ids after filtering".to_string()),
});
}
Ok(MergeDecision::Merge(result))
}
MergeDecision::Conflict {
mut conflicting_ids,
reason,
} => {
let before = conflicting_ids.len();
conflicting_ids.retain(|id| member_ids.contains(id));
let dropped = before - conflicting_ids.len();
if dropped > 0 {
crate::log::warn(
"dream",
&format!(
"dropped {} hallucinated conflicting_id(s) not in cluster",
dropped
),
);
bail!("dream conflict referenced memory id outside cluster");
}
conflicting_ids.sort_unstable();
conflicting_ids.dedup();
if conflicting_ids.len() < 2 {
crate::log::warn(
"dream",
"rejecting conflict with fewer than two valid conflicting id(s) after filtering",
);
bail!("dream conflict requires at least two valid cluster ids");
}
Ok(MergeDecision::Conflict {
conflicting_ids,
reason,
})
}
MergeDecision::NoMerge { .. } => Ok(decision),
}
}
fn build_user_message(members: &[MemoryCandidate]) -> String {
let mut msg = String::from("Merge these memory entries:\n\n");
for m in members {
msg.push_str(&format!(
"<entry id=\"{}\" type=\"{}\" topic_key=\"{}\">\n<title>{}</title>\n<content>{}</content>\n</entry>\n\n",
m.id,
xml_escape(&m.memory_type),
xml_escape(m.topic_key.as_deref().unwrap_or("")),
xml_escape(&m.title),
xml_escape(&m.content),
));
}
msg
}
fn parse_response(response: &str) -> Result<MergeDecision> {
match first_decision_tag(response) {
Some(DecisionTag::Conflict) => {
return Ok(MergeDecision::Conflict {
conflicting_ids: extract_conflict_ids(response)?,
reason: extract_conflict_reason(response),
});
}
Some(DecisionTag::NoMerge) => {
return Ok(MergeDecision::NoMerge {
reason: extract_no_merge_reason(response),
});
}
Some(DecisionTag::Memory) | None => {}
}
let topic_key = require_tag(response, "topic_key")?;
let title = require_tag(response, "title")?;
let content = require_tag(response, "content")?;
let memory_type = extract_tag(response, "type")
.filter(|s| !s.trim().is_empty())
.unwrap_or_else(|| "discovery".to_owned());
let supersedes_raw = extract_tag(response, "supersedes").unwrap_or_default();
let superseded_ids: Vec<i64> = supersedes_raw
.split_whitespace()
.filter_map(|s| s.parse::<i64>().ok())
.collect();
if superseded_ids.is_empty() {
return Ok(MergeDecision::NoMerge {
reason: Some("merge response had no superseded ids".to_string()),
});
}
Ok(MergeDecision::Merge(MergeResult {
topic_key,
memory_type,
title,
content,
superseded_ids,
}))
}
fn first_decision_tag(response: &str) -> Option<DecisionTag> {
[
("<conflict", DecisionTag::Conflict),
("<no_merge", DecisionTag::NoMerge),
("<memory", DecisionTag::Memory),
]
.into_iter()
.filter_map(|(marker, tag)| response.find(marker).map(|index| (index, tag)))
.min_by_key(|(index, _)| *index)
.map(|(_, tag)| tag)
}
fn require_tag(response: &str, tag: &str) -> Result<String> {
extract_tag(response, tag)
.filter(|s| !s.trim().is_empty())
.ok_or_else(|| {
let response_sha256 = crate::db::content_identity_hash(response.as_bytes());
crate::log::error(
"dream",
&format!(
"merge_parse_error error_code=missing_required_tag field={} response_bytes={} response_sha256={}",
tag,
response.len(),
response_sha256
),
);
anyhow!("dream merge parse failed error_code=missing_required_tag field=<{tag}>")
})
}
fn extract_tag(text: &str, tag: &str) -> Option<String> {
let open = format!("<{}>", tag);
let close = format!("</{}>", tag);
let start = text.find(&open)? + open.len();
let end = text[start..].find(&close)? + start;
Some(text[start..end].trim().to_owned())
}
fn extract_no_merge_reason(text: &str) -> Option<String> {
let start = text.find("<no_merge")?;
let tag = &text[start..];
let end = tag.find('>')?;
extract_attr(&tag[..=end], "reason")
.map(|reason| reason.trim().to_string())
.filter(|reason| !reason.is_empty())
}
fn extract_conflict_ids(text: &str) -> Result<Vec<i64>> {
let Some(start) = text.find("<conflict") else {
bail!("conflict response missing <conflict> tag");
};
let tag = &text[start..];
let Some(end) = tag.find('>') else {
bail!("conflict response missing closing tag bracket");
};
let ids_raw = extract_attr(&tag[..=end], "ids")
.map(|ids| ids.trim().to_string())
.filter(|ids| !ids.is_empty())
.ok_or_else(|| anyhow!("conflict response missing ids attribute"))?;
let mut ids = Vec::new();
for token in ids_raw.split_whitespace() {
let id = token.parse::<i64>().map_err(|_| {
anyhow!(
"conflict response contains invalid memory id error_code=invalid_memory_id token_bytes={} token_sha256={}",
token.len(),
crate::db::content_identity_hash(token.as_bytes())
)
})?;
ids.push(id);
}
if ids.len() < 2 {
bail!("conflict response requires at least two memory ids");
}
Ok(ids)
}
fn extract_conflict_reason(text: &str) -> Option<String> {
let start = text.find("<conflict")?;
let tag = &text[start..];
let end = tag.find('>')?;
extract_attr(&tag[..=end], "reason")
.map(|reason| reason.trim().to_string())
.filter(|reason| !reason.is_empty())
}
fn extract_attr(tag: &str, name: &str) -> Option<String> {
let marker = format!("{name}=");
let start = tag.find(&marker)? + marker.len();
let quote = tag[start..].chars().next()?;
if quote != '"' && quote != '\'' {
return None;
}
let value_start = start + quote.len_utf8();
let value_end = tag[value_start..].find(quote)? + value_start;
Some(xml_unescape(&tag[value_start..value_end]))
}
fn xml_unescape(s: &str) -> String {
s.replace(""", "\"")
.replace("'", "'")
.replace("<", "<")
.replace(">", ">")
.replace("&", "&")
}
fn xml_escape(s: &str) -> String {
s.replace('&', "&")
.replace('<', "<")
.replace('>', ">")
.replace('"', """)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_parse_valid_merge() {
let response = r#"<memory>
<topic_key>auth-design</topic_key>
<type>decision</type>
<title>Auth middleware uses JWT</title>
<content>Use JWT for stateless auth. Previously: session cookies.</content>
<supersedes>42 17</supersedes>
</memory>"#;
match parse_response(response).expect("expected Ok") {
MergeDecision::Merge(r) => {
assert_eq!(r.topic_key, "auth-design");
assert_eq!(r.memory_type, "decision");
assert_eq!(r.superseded_ids, vec![42, 17]);
}
MergeDecision::NoMerge { .. } => panic!("expected Merge"),
MergeDecision::Conflict { .. } => panic!("expected Merge"),
}
}
#[test]
fn test_parse_merge_content_with_literal_conflict_tag() -> Result<()> {
let response = r#"<memory>
<topic_key>dream-contract</topic_key>
<type>decision</type>
<title>Dream conflict contract</title>
<content>The prompt may return <conflict ids="10 20" reason="values differ"/> when facts disagree.</content>
<supersedes>10 20</supersedes>
</memory>"#;
match parse_response(response)? {
MergeDecision::Merge(r) => {
assert_eq!(r.topic_key, "dream-contract");
assert_eq!(r.superseded_ids, vec![10, 20]);
assert!(r.content.contains("<conflict ids="));
}
MergeDecision::NoMerge { .. } => panic!("expected Merge"),
MergeDecision::Conflict { .. } => panic!("expected Merge"),
}
Ok(())
}
#[test]
fn test_parse_no_merge() {
let response = r#"<no_merge reason="entries cover different topics"/>"#;
assert!(matches!(
parse_response(response).expect("expected Ok"),
MergeDecision::NoMerge { .. }
));
}
#[test]
fn test_parse_no_merge_reason() {
let response = r#"<no_merge reason="entries cover different topics"/>"#;
match parse_response(response).expect("expected Ok") {
MergeDecision::NoMerge { reason } => {
assert_eq!(reason.as_deref(), Some("entries cover different topics"));
}
MergeDecision::Merge(_) => panic!("expected no merge"),
MergeDecision::Conflict { .. } => panic!("expected no merge"),
}
}
#[test]
fn test_parse_conflict_reason_and_ids() -> Result<()> {
let response = r#"<conflict ids="20 10" reason="same setting has incompatible values"/>"#;
match parse_response(response)? {
MergeDecision::Conflict {
conflicting_ids,
reason,
} => {
assert_eq!(conflicting_ids, vec![20, 10]);
assert_eq!(
reason.as_deref(),
Some("same setting has incompatible values")
);
}
MergeDecision::Merge(_) | MergeDecision::NoMerge { .. } => panic!("expected conflict"),
}
Ok(())
}
#[test]
fn test_parse_conflict_invalid_id_errors() {
let response = r#"<conflict ids="20 not-a-number" reason="bad id"/>"#;
let Err(err) = parse_response(response) else {
panic!("invalid conflict id must fail");
};
assert!(
err.to_string().contains("invalid memory id"),
"unexpected error: {err}"
);
}
#[test]
fn test_parse_missing_required_tag_errors() {
let response = "<memory><topic_key>k</topic_key></memory>";
let err = parse_response(response).expect_err("expected Err");
assert!(
err.to_string().contains("<title>"),
"error should name the missing tag, got: {err}"
);
}
#[test]
fn test_parse_whitespace_only_required_tag_errors() {
let response = r#"<memory>
<topic_key>k</topic_key>
<title> </title>
<content>C</content>
<supersedes>1</supersedes>
</memory>"#;
let err = parse_response(response).expect_err("expected Err");
assert!(err.to_string().contains("<title>"));
}
#[test]
fn test_parse_empty_topic_key_errors() {
let response = r#"<memory>
<topic_key></topic_key>
<title>T</title>
<content>C</content>
<supersedes>1</supersedes>
</memory>"#;
let err = parse_response(response).expect_err("expected Err");
assert!(err.to_string().contains("<topic_key>"));
}
#[test]
fn test_parse_empty_content_errors() {
let response = r#"<memory>
<topic_key>k</topic_key>
<title>T</title>
<content></content>
<supersedes>1</supersedes>
</memory>"#;
let err = parse_response(response).expect_err("expected Err");
assert!(err.to_string().contains("<content>"));
}
#[test]
fn test_parse_missing_type_defaults_to_discovery() {
let response = r#"<memory>
<topic_key>k</topic_key>
<title>T</title>
<content>C</content>
<supersedes>1</supersedes>
</memory>"#;
match parse_response(response).expect("expected Ok") {
MergeDecision::Merge(r) => assert_eq!(r.memory_type, "discovery"),
MergeDecision::NoMerge { .. } => panic!("expected Merge"),
MergeDecision::Conflict { .. } => panic!("expected Merge"),
}
}
#[test]
fn test_filter_superseded_ids_drops_hallucinated() -> Result<()> {
let cluster = Cluster {
members: vec![
MemoryCandidate {
id: 10,
version: 1,
topic_key: Some("k".into()),
title: "t".into(),
content: "c".into(),
memory_type: "decision".into(),
updated_at_epoch: 0,
},
MemoryCandidate {
id: 20,
version: 1,
topic_key: Some("k".into()),
title: "t".into(),
content: "c".into(),
memory_type: "decision".into(),
updated_at_epoch: 0,
},
],
};
let decision = MergeDecision::Merge(MergeResult {
topic_key: "k".into(),
memory_type: "decision".into(),
title: "T".into(),
content: "C".into(),
superseded_ids: vec![10, 99999, 20],
});
match filter_superseded_ids(decision, &cluster)? {
MergeDecision::Merge(r) => assert_eq!(r.superseded_ids, vec![10, 20]),
MergeDecision::NoMerge { .. } => panic!("expected Merge"),
MergeDecision::Conflict { .. } => panic!("expected Merge"),
}
Ok(())
}
#[test]
fn test_filter_conflict_ids_dedupes_and_sorts() -> Result<()> {
let cluster = Cluster {
members: vec![
MemoryCandidate {
id: 10,
version: 1,
topic_key: Some("k".into()),
title: "t".into(),
content: "c".into(),
memory_type: "decision".into(),
updated_at_epoch: 0,
},
MemoryCandidate {
id: 20,
version: 1,
topic_key: Some("k".into()),
title: "t".into(),
content: "c".into(),
memory_type: "decision".into(),
updated_at_epoch: 0,
},
],
};
let decision = MergeDecision::Conflict {
conflicting_ids: vec![20, 10, 20],
reason: Some("same state differs".into()),
};
match filter_superseded_ids(decision, &cluster)? {
MergeDecision::Conflict {
conflicting_ids,
reason,
} => {
assert_eq!(conflicting_ids, vec![10, 20]);
assert_eq!(reason.as_deref(), Some("same state differs"));
}
MergeDecision::Merge(_) | MergeDecision::NoMerge { .. } => panic!("expected conflict"),
}
Ok(())
}
#[test]
fn test_filter_conflict_ids_rejects_hallucinated() {
let cluster = Cluster {
members: vec![
MemoryCandidate {
id: 10,
version: 1,
topic_key: Some("k".into()),
title: "t".into(),
content: "c".into(),
memory_type: "decision".into(),
updated_at_epoch: 0,
},
MemoryCandidate {
id: 20,
version: 1,
topic_key: Some("k".into()),
title: "t".into(),
content: "c".into(),
memory_type: "decision".into(),
updated_at_epoch: 0,
},
],
};
let decision = MergeDecision::Conflict {
conflicting_ids: vec![20, 99999, 10],
reason: Some("same state differs".into()),
};
let Err(err) = filter_superseded_ids(decision, &cluster) else {
panic!("hallucinated conflict id must fail");
};
assert!(
err.to_string().contains("outside cluster"),
"unexpected error: {err}"
);
}
#[test]
fn test_filter_superseded_ids_rejects_empty_after_filter() -> Result<()> {
let cluster = Cluster {
members: vec![MemoryCandidate {
id: 10,
version: 1,
topic_key: Some("k".into()),
title: "t".into(),
content: "c".into(),
memory_type: "decision".into(),
updated_at_epoch: 0,
}],
};
let decision = MergeDecision::Merge(MergeResult {
topic_key: "k".into(),
memory_type: "decision".into(),
title: "T".into(),
content: "C".into(),
superseded_ids: vec![99999],
});
assert!(matches!(
filter_superseded_ids(decision, &cluster)?,
MergeDecision::NoMerge { .. }
));
Ok(())
}
#[test]
fn test_filter_superseded_ids_no_merge_passthrough() -> Result<()> {
let cluster = Cluster { members: vec![] };
assert!(matches!(
filter_superseded_ids(MergeDecision::NoMerge { reason: None }, &cluster)?,
MergeDecision::NoMerge { .. }
));
Ok(())
}
#[test]
fn test_parse_empty_supersedes_becomes_no_merge() {
let response = r#"<memory>
<topic_key>k</topic_key>
<type>decision</type>
<title>T</title>
<content>C</content>
<supersedes></supersedes>
</memory>"#;
assert!(matches!(
parse_response(response).expect("expected Ok"),
MergeDecision::NoMerge { .. }
));
}
#[test]
fn test_parse_error_log_contains_only_structured_response_metadata() -> Result<()> {
let sentinel = "RAW_MODEL_SENTINEL_sk-ABCDEFGH123";
let response = format!(
"<memory><topic_key>provider</topic_key><content>{sentinel}</content></memory>"
);
let log_dir = std::env::temp_dir().join(format!(
"remem-dream-parse-log-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)?
.as_nanos()
));
std::fs::create_dir_all(&log_dir)?;
let error = crate::log::with_log_dir(&log_dir, || {
parse_response(&response).expect_err("missing title must fail")
});
let log = std::fs::read_to_string(log_dir.join("remem.log"))?;
std::fs::remove_dir_all(&log_dir)?;
assert!(!error.to_string().contains(sentinel));
assert!(!log.contains(sentinel));
assert!(log.contains("error_code=missing_required_tag"));
assert!(log.contains("response_bytes="));
assert!(log.contains("response_sha256=sha256:content-v1:"));
Ok(())
}
}