use std::collections::{BTreeMap, BTreeSet};
use anyhow::{Context, Result};
use rusqlite::{params, Connection};
use crate::memory::dedup::{
canonical_observation_text, find_hash_duplicates, mark_duplicate_accessed,
};
const VECTOR_WINDOW_SECS: i64 = 15 * 60;
const VECTOR_CANDIDATE_LIMIT: i64 = 200;
const FEATURE_HASH_OBSERVATION_THRESHOLD: f32 = 0.82;
const REAL_EMBEDDING_OBSERVATION_THRESHOLD: f32 = 0.92;
pub fn check_duplicate(
conn: &Connection,
project: &str,
narrative: &str,
_embedding: Option<&[f32]>,
) -> Result<Option<i64>> {
let content_hash = crate::db::content_identity_hash(narrative.as_bytes());
let hash_dups = find_hash_duplicates(conn, project, &content_hash, VECTOR_WINDOW_SECS)?;
if !hash_dups.is_empty() {
crate::log::info(
"dedup",
&format!("hash duplicate found: {} matches", hash_dups.len()),
);
mark_duplicate_accessed(conn, &hash_dups)?;
return Ok(Some(hash_dups[0]));
}
let vector_dups = find_vector_duplicates(conn, project, narrative, VECTOR_WINDOW_SECS)?;
if !vector_dups.is_empty() {
crate::log::info(
"dedup",
&format!("vector duplicate found: {} matches", vector_dups.len()),
);
mark_duplicate_accessed(conn, &vector_dups)?;
return Ok(Some(vector_dups[0]));
}
Ok(None)
}
pub(crate) fn check_duplicate_texts(narrative: &str, candidate_texts: &[String]) -> Result<bool> {
if candidate_texts.is_empty() {
return Ok(false);
}
let content_hash = crate::db::content_identity_hash(narrative.as_bytes());
if candidate_texts.iter().any(|candidate| {
let candidate_hash = crate::db::content_identity_hash(candidate.as_bytes());
let legacy_candidate_hash = crate::db::legacy_content_identity_hash(candidate.as_bytes());
candidate_hash == content_hash || legacy_candidate_hash == content_hash
}) {
return Ok(true);
}
let duplicate_indexes = find_vector_duplicate_indexes(
narrative,
candidate_texts
.iter()
.enumerate()
.map(|(index, text)| (index, text.as_str())),
)?;
Ok(!duplicate_indexes.is_empty())
}
fn find_vector_duplicates(
conn: &Connection,
project: &str,
narrative: &str,
window_secs: i64,
) -> Result<Vec<i64>> {
let cutoff = chrono::Utc::now().timestamp() - window_secs;
let mut stmt = conn.prepare_cached(
"SELECT id, text, narrative, title, facts
FROM observations
WHERE project = ?1
AND status = 'active'
AND created_at_epoch > ?2
AND (
(text IS NOT NULL AND length(text) > 0)
OR (narrative IS NOT NULL AND length(narrative) > 0)
OR (title IS NOT NULL AND length(title) > 0)
OR (facts IS NOT NULL AND length(facts) > 0)
)
ORDER BY created_at_epoch DESC, id DESC
LIMIT ?3",
)?;
let rows = stmt.query_map(params![project, cutoff, VECTOR_CANDIDATE_LIMIT], |row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, Option<String>>(1)?,
row.get::<_, Option<String>>(2)?,
row.get::<_, Option<String>>(3)?,
row.get::<_, Option<String>>(4)?,
))
})?;
let candidates = crate::db::query::collect_rows(rows)?;
let candidates = candidates
.into_iter()
.filter_map(|(id, text, candidate_narrative, title, facts)| {
canonical_observation_text(
text.as_deref(),
candidate_narrative.as_deref(),
title.as_deref(),
facts.as_deref(),
)
.map(|candidate_text| (id, candidate_text))
})
.collect::<Vec<_>>();
let duplicate_indexes = find_vector_duplicate_indexes(
narrative,
candidates
.iter()
.enumerate()
.map(|(index, (_id, text))| (index, text.as_str())),
)?;
Ok(duplicate_indexes
.into_iter()
.map(|index| candidates[index].0)
.collect())
}
fn find_vector_duplicate_indexes<'a>(
narrative: &str,
candidates: impl IntoIterator<Item = (usize, &'a str)>,
) -> Result<Vec<usize>> {
let candidates = candidates
.into_iter()
.filter(|(_index, text)| !text.trim().is_empty())
.filter(|(_index, text)| !observation_text_conflicts(narrative, text))
.collect::<Vec<_>>();
if candidates.is_empty() {
return Ok(Vec::new());
}
let mut fallback_cache = crate::retrieval::embedding::EmbeddingFallbackCache::default();
let mut query_embedding = match embed_observation_dedup_text(narrative, &mut fallback_cache) {
Ok(embedding) => embedding,
Err(error) if crate::retrieval::embedding::is_embedding_provider_off_error(&error) => {
return Ok(Vec::new());
}
Err(error) => return Err(error),
};
let mut threshold = observation_similarity_threshold(query_embedding.model());
let mut max_distance = 1.0 - threshold;
let mut duplicates = Vec::new();
for (index, candidate_text) in candidates {
let candidate_embedding =
match embed_observation_dedup_text(candidate_text, &mut fallback_cache) {
Ok(embedding) => embedding,
Err(error)
if crate::retrieval::embedding::is_embedding_provider_off_error(&error) =>
{
return Ok(Vec::new());
}
Err(error) => {
return Err(error).with_context(|| {
format!("embed observation duplicate candidate index={index}")
});
}
};
if candidate_embedding.model() != query_embedding.model()
|| candidate_embedding.dimensions() != query_embedding.dimensions()
{
query_embedding = match embed_observation_dedup_text(narrative, &mut fallback_cache) {
Ok(embedding) => embedding,
Err(error)
if crate::retrieval::embedding::is_embedding_provider_off_error(&error) =>
{
return Ok(Vec::new());
}
Err(error) => {
return Err(error)
.context("re-embed observation query after provider fallback");
}
};
threshold = observation_similarity_threshold(query_embedding.model());
max_distance = 1.0 - threshold;
}
if candidate_embedding.model() != query_embedding.model()
|| candidate_embedding.dimensions() != query_embedding.dimensions()
{
continue;
}
let distance = crate::retrieval::vector::cosine_distance(
query_embedding.values(),
candidate_embedding.values(),
)
.with_context(|| format!("compare observation duplicate candidate index={index}"))?;
if distance <= max_distance {
duplicates.push(index);
}
}
Ok(duplicates)
}
fn embed_observation_dedup_text(
text: &str,
fallback_cache: &mut crate::retrieval::embedding::EmbeddingFallbackCache,
) -> Result<crate::retrieval::embedding::TextEmbedding> {
let text = normalize_observation_dedup_embedding_text(text);
crate::retrieval::embedding::embed_query_with_fallback_cache(&text, fallback_cache)
}
fn normalize_observation_dedup_embedding_text(text: &str) -> String {
let chars = text.chars().collect::<Vec<_>>();
let mut normalized = String::with_capacity(text.len());
for (index, ch) in chars.iter().enumerate() {
let previous = index.checked_sub(1).and_then(|prev| chars.get(prev));
let next = chars.get(index + 1);
if *ch == ',' {
if previous.is_some_and(|previous| previous.is_ascii_digit())
&& next.is_some_and(|next| next.is_ascii_digit())
{
continue;
}
} else if *ch == '/'
&& previous.is_some_and(|previous| previous.is_ascii_alphabetic())
&& next.is_some_and(|next| next.is_ascii_digit())
{
continue;
}
normalized.push(*ch);
}
normalized
}
fn observation_similarity_threshold(model: &str) -> f32 {
if model == crate::retrieval::embedding::FEATURE_HASH_EMBEDDING_MODEL {
FEATURE_HASH_OBSERVATION_THRESHOLD
} else {
REAL_EMBEDDING_OBSERVATION_THRESHOLD
}
}
fn observation_text_conflicts(incoming: &str, existing: &str) -> bool {
let incoming_tokens = observation_token_list(incoming);
let existing_tokens = observation_token_list(existing);
if numeric_tokens_differ(&incoming_tokens, &existing_tokens) {
return true;
}
let incoming = observation_token_set(&incoming_tokens);
let existing = observation_token_set(&existing_tokens);
const OPPOSITES: &[(&[&str], &[&str])] = &[
(
&[
"pass",
"passed",
"passes",
"passing",
"success",
"succeeded",
"successful",
],
&["fail", "failed", "fails", "failing", "failure", "broken"],
),
(
&["enable", "enabled", "enables", "on", "active"],
&["disable", "disabled", "disables", "off", "inactive"],
),
(
&["allow", "allowed", "allows", "accept", "accepted"],
&["block", "blocked", "blocks", "reject", "rejected"],
),
(
&["increase", "increased", "higher"],
&["decrease", "decreased", "lower"],
),
(&["true", "yes", "present"], &["false", "no", "absent"]),
];
OPPOSITES.iter().any(|(left, right)| {
(has_any(&incoming, left) && has_any(&existing, right))
|| (has_any(&incoming, right) && has_any(&existing, left))
|| negated_same_status_conflict(
&incoming_tokens,
&incoming,
&existing_tokens,
&existing,
left,
)
|| negated_same_status_conflict(
&incoming_tokens,
&incoming,
&existing_tokens,
&existing,
right,
)
})
}
fn observation_token_list(text: &str) -> Vec<String> {
let chars = text.chars().collect::<Vec<_>>();
let mut tokens = Vec::new();
let mut current = String::new();
for (index, ch) in chars.iter().enumerate() {
let previous = index.checked_sub(1).and_then(|prev| chars.get(prev));
let next = chars.get(index + 1);
if ch.is_alphanumeric()
|| *ch == '%'
|| is_numeric_token_separator(*ch, previous.copied(), next.copied())
{
current.extend(ch.to_lowercase());
} else if !current.is_empty() {
tokens.push(std::mem::take(&mut current));
}
}
if !current.is_empty() {
tokens.push(current);
}
tokens
}
fn is_numeric_token_separator(ch: char, previous: Option<char>, next: Option<char>) -> bool {
match ch {
'+' | '-' => {
next.is_some_and(|next| next.is_ascii_digit())
&& !previous.is_some_and(|previous| previous.is_ascii_digit())
}
',' | '.' => {
next.is_some_and(|next| next.is_ascii_digit())
&& (previous.is_none()
|| previous.is_some_and(|previous| {
previous.is_ascii_digit() || !previous.is_alphanumeric()
}))
}
'/' => {
previous.is_some_and(|previous| previous.is_alphanumeric())
&& next.is_some_and(|next| next.is_alphanumeric())
}
_ => false,
}
}
fn observation_token_set(tokens: &[String]) -> BTreeSet<String> {
tokens.iter().cloned().collect()
}
fn numeric_tokens_differ(incoming: &[String], existing: &[String]) -> bool {
let incoming = numeric_signature_counts(incoming);
let existing = numeric_signature_counts(existing);
(!incoming.is_empty() || !existing.is_empty()) && incoming != existing
}
#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
struct NumericSignature {
label: String,
role: String,
value: String,
unit: String,
}
#[derive(Clone, Debug, Eq, PartialEq)]
struct NumericTokenPart {
label: Option<String>,
value: String,
unit: String,
}
fn numeric_signature_counts(tokens: &[String]) -> BTreeMap<NumericSignature, usize> {
let mut counts = BTreeMap::new();
for (index, token) in tokens.iter().enumerate() {
for mut part in numeric_token_parts(token) {
if part.unit.is_empty() {
part.unit = tokens
.get(index + 1)
.filter(|next| is_numeric_unit_token(next))
.cloned()
.unwrap_or_default();
}
part.unit = normalize_numeric_unit(&part.unit);
let label = part
.label
.unwrap_or_else(|| nearby_numeric_label(tokens, index));
let signature = NumericSignature {
value: normalize_numeric_value(&part.value, &part.unit, &label),
label,
role: nearby_numeric_role(tokens, index),
unit: part.unit,
};
*counts.entry(signature).or_insert(0) += 1;
}
}
counts
}
fn numeric_token_parts(token: &str) -> Vec<NumericTokenPart> {
let chars = token.chars().collect::<Vec<_>>();
let mut parts = Vec::new();
let mut index = 0;
while index < chars.len() {
if !is_numeric_value_start(&chars, index) {
index += 1;
continue;
}
let value_start = index;
if matches!(chars[index], '+' | '-') {
index += 1;
}
if chars.get(index) == Some(&'.') {
index += 1;
}
while index < chars.len() {
if chars[index].is_ascii_digit()
|| (matches!(chars[index], ',' | '.')
&& index > value_start
&& chars
.get(index - 1)
.is_some_and(|previous| previous.is_ascii_digit())
&& chars
.get(index + 1)
.is_some_and(|next| next.is_ascii_digit()))
{
index += 1;
} else {
break;
}
}
let value_end = index;
let prefix = alphabetic_prefix_before(&chars, value_start);
let mut suffix_end = value_end;
while suffix_end < chars.len()
&& (chars[suffix_end].is_ascii_alphabetic() || chars[suffix_end] == '%')
{
suffix_end += 1;
}
let suffix = chars[value_end..suffix_end].iter().collect::<String>();
let mut unit = String::new();
let label = if is_numeric_unit_prefix(&prefix) {
unit.push_str(&prefix);
None
} else if prefix.is_empty() {
None
} else {
Some(prefix)
};
if is_numeric_unit_token(&suffix) {
unit.push_str(&suffix);
}
parts.push(NumericTokenPart {
label,
value: chars[value_start..value_end].iter().collect::<String>(),
unit,
});
}
parts
}
fn is_numeric_value_start(chars: &[char], index: usize) -> bool {
chars[index].is_ascii_digit()
|| (matches!(chars[index], '+' | '-')
&& chars
.get(index + 1)
.is_some_and(|next| next.is_ascii_digit()))
|| (chars[index] == '.'
&& chars
.get(index + 1)
.is_some_and(|next| next.is_ascii_digit()))
}
fn alphabetic_prefix_before(chars: &[char], index: usize) -> String {
let mut prefix = Vec::new();
let mut cursor = index;
while cursor > 0 {
let ch = chars[cursor - 1];
if ch.is_ascii_alphabetic() {
prefix.push(ch);
} else if matches!(ch, '/' | '-' | '_')
&& cursor > 1
&& chars[cursor - 2].is_ascii_alphabetic()
{
} else {
break;
}
cursor -= 1;
}
prefix.into_iter().rev().collect()
}
fn normalize_numeric_value(raw: &str, unit: &str, label: &str) -> String {
let raw = normalize_grouped_numeric_commas(raw);
let (sign, unsigned) = raw
.strip_prefix('-')
.map(|value| ("-", value))
.or_else(|| raw.strip_prefix('+').map(|value| ("", value)))
.unwrap_or(("", raw.as_str()));
let (integer, fraction) = unsigned.split_once('.').unwrap_or((unsigned, ""));
let integer = integer.trim_start_matches('0');
let integer = if integer.is_empty() { "0" } else { integer };
let fraction = if is_version_numeric_context(unit, label) {
fraction
} else {
fraction.trim_end_matches('0')
};
let value = if fraction.is_empty() {
integer.to_string()
} else {
format!("{integer}.{fraction}")
};
if value == "0" {
value
} else {
format!("{sign}{value}")
}
}
fn is_version_numeric_context(unit: &str, label: &str) -> bool {
unit == "v"
|| label == "release"
|| label == "version"
|| label.ends_with(":release")
|| label.ends_with(":version")
}
fn normalize_grouped_numeric_commas(raw: &str) -> String {
if raw.contains(',') && has_valid_thousands_grouping(raw) {
raw.replace(',', "")
} else {
raw.to_string()
}
}
fn has_valid_thousands_grouping(raw: &str) -> bool {
let unsigned = raw
.strip_prefix('-')
.or_else(|| raw.strip_prefix('+'))
.unwrap_or(raw);
let integer = unsigned
.split_once('.')
.map_or(unsigned, |(integer, _)| integer);
let mut groups = integer.split(',');
let Some(first) = groups.next() else {
return false;
};
(1..=3).contains(&first.len())
&& first.chars().all(|ch| ch.is_ascii_digit())
&& groups.clone().count() > 0
&& groups.all(|group| group.len() == 3 && group.chars().all(|ch| ch.is_ascii_digit()))
}
fn nearby_numeric_label(tokens: &[String], index: usize) -> String {
tokens[..index]
.iter()
.enumerate()
.rev()
.take(4)
.find(|(_, token)| is_numeric_label_token(token))
.map(|(label_index, token)| numeric_label_with_qualifier(tokens, label_index, token))
.or_else(|| {
tokens
.iter()
.enumerate()
.skip(index + 1)
.take(4)
.find(|(_, token)| is_numeric_label_token(token))
.map(|(label_index, token)| {
numeric_label_with_qualifier(tokens, label_index, token)
})
})
.unwrap_or_default()
}
fn numeric_label_with_qualifier(tokens: &[String], label_index: usize, label: &str) -> String {
let label = if numeric_label_repeats(tokens, label) {
nearby_numeric_entity(tokens, label_index, label)
.map(|entity| format!("{entity}:{label}"))
.unwrap_or_else(|| label.to_string())
} else {
label.to_string()
};
tokens[..label_index]
.iter()
.rev()
.take(2)
.find(|token| is_numeric_qualifier_token(token))
.map(|qualifier| format!("{qualifier}:{label}"))
.unwrap_or(label)
}
fn numeric_label_repeats(tokens: &[String], label: &str) -> bool {
tokens
.iter()
.filter(|token| token.as_str() == label)
.count()
> 1
}
fn nearby_numeric_entity(tokens: &[String], label_index: usize, label: &str) -> Option<String> {
tokens[..label_index]
.iter()
.rev()
.take(3)
.find(|token| token.as_str() != label && is_numeric_label_token(token))
.cloned()
}
fn is_numeric_qualifier_token(token: &str) -> bool {
matches!(
token,
"ceiling" | "floor" | "max" | "maximum" | "min" | "minimum"
)
}
fn nearby_numeric_role(tokens: &[String], index: usize) -> String {
if !has_numeric_transition_role_pair(tokens) {
return String::new();
}
tokens[..index]
.iter()
.rev()
.take(3)
.find(|token| is_numeric_role_token(token))
.cloned()
.unwrap_or_default()
}
fn has_numeric_transition_role_pair(tokens: &[String]) -> bool {
(tokens.iter().any(|token| token == "from")
&& tokens.iter().any(|token| {
matches!(
token.as_str(),
"into" | "onto" | "to" | "toward" | "towards"
)
}))
|| (tokens.iter().any(|token| token == "before")
&& tokens.iter().any(|token| token == "after"))
}
fn is_numeric_role_token(token: &str) -> bool {
matches!(
token,
"after" | "before" | "from" | "into" | "onto" | "to" | "toward" | "towards"
)
}
fn is_numeric_label_token(token: &str) -> bool {
token.chars().any(|ch| ch.is_ascii_alphabetic())
&& token.chars().all(|ch| !ch.is_ascii_digit())
&& !is_numeric_context_stopword(token)
&& !is_numeric_unit_token(token)
}
fn is_numeric_context_stopword(token: &str) -> bool {
matches!(
token,
"a" | "an"
| "and"
| "are"
| "at"
| "be"
| "been"
| "being"
| "by"
| "for"
| "from"
| "in"
| "is"
| "of"
| "on"
| "or"
| "set"
| "the"
| "to"
| "use"
| "uses"
| "using"
| "value"
| "was"
| "were"
| "with"
)
}
fn is_numeric_unit_token(token: &str) -> bool {
matches!(
token,
"b" | "byte"
| "bytes"
| "gb"
| "gib"
| "h"
| "hour"
| "hours"
| "hr"
| "hrs"
| "kb"
| "kib"
| "mb"
| "mib"
| "m"
| "min"
| "mins"
| "minute"
| "minutes"
| "millisecond"
| "milliseconds"
| "msec"
| "msecs"
| "ms"
| "pct"
| "percent"
| "percentage"
| "%"
| "s"
| "sec"
| "second"
| "seconds"
| "secs"
| "d"
| "day"
| "days"
)
}
fn is_numeric_unit_prefix(prefix: &str) -> bool {
matches!(prefix, "v")
}
fn normalize_numeric_unit(unit: &str) -> String {
match unit {
"" => String::new(),
"byte" | "bytes" => "b".to_string(),
"day" | "days" => "d".to_string(),
"hour" | "hours" | "hr" | "hrs" => "h".to_string(),
"minute" | "minutes" | "min" | "mins" => "m".to_string(),
"millisecond" | "milliseconds" | "msec" | "msecs" => "ms".to_string(),
"pct" | "percent" | "percentage" | "%" => "%".to_string(),
"sec" | "secs" | "second" | "seconds" => "s".to_string(),
_ => unit.to_string(),
}
}
fn has_any(tokens: &BTreeSet<String>, values: &[&str]) -> bool {
values.iter().any(|value| tokens.contains(*value))
}
fn negated_same_status_conflict(
incoming_tokens: &[String],
incoming: &BTreeSet<String>,
existing_tokens: &[String],
existing: &BTreeSet<String>,
values: &[&str],
) -> bool {
(has_any(incoming, values) && has_negated_any(existing_tokens, values))
|| (has_negated_any(incoming_tokens, values) && has_any(existing, values))
}
fn has_negated_any(tokens: &[String], values: &[&str]) -> bool {
tokens.iter().enumerate().any(|(index, token)| {
values.contains(&token.as_str()) && {
let start = index.saturating_sub(3);
tokens[start..index]
.iter()
.any(|candidate| is_negation_token(candidate))
}
})
}
fn is_negation_token(token: &str) -> bool {
matches!(
token,
"not"
| "never"
| "no"
| "without"
| "cannot"
| "cant"
| "don"
| "doesn"
| "didn"
| "isn"
| "wasn"
| "weren"
| "won"
)
}