use std::collections::{HashMap, HashSet};
use rusqlite::{Connection, TransactionBehavior, params};
use uuid::Uuid;
use crate::engine::memory::extract::{PendingSources, SourceCandidate};
use crate::engine::memory::types::{MemoryError, RememberInput, RememberReport};
use super::util::{now_millis, session_tombstone_key};
use super::write::ensure_project;
const EXTRACTION_LEASE_MILLIS: i64 = 30 * 60 * 1_000;
pub(super) struct ClaimedEvidence {
pub(super) pending: PendingSources,
pub(super) evidence_ids: HashMap<String, i64>,
pub(super) evidence_added: usize,
pub(super) owner: String,
}
impl ClaimedEvidence {
pub(super) fn empty_report(&self) -> RememberReport {
let mut report = self.pending.empty_report();
report.evidence_added = self.evidence_added;
report
}
}
pub(super) fn claim_evidence(
conn: &mut Connection,
input: &RememberInput<'_>,
mut pending: PendingSources,
) -> anyhow::Result<ClaimedEvidence> {
let owner = Uuid::new_v4().to_string();
let now = now_millis();
let lease_until = now.saturating_add(EXTRACTION_LEASE_MILLIS);
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let key = session_tombstone_key(input.provider.as_str(), input.session_id);
let forgotten: bool = tx.query_row(
"SELECT EXISTS(SELECT 1 FROM tombstones WHERE kind = 'session' AND key = ?1)",
params![key],
|row| row.get(0),
)?;
if forgotten {
pending.evidence.clear();
pending.context_forgotten = true;
tx.commit()?;
return Ok(ClaimedEvidence {
pending,
evidence_ids: HashMap::new(),
evidence_added: 0,
owner,
});
}
let project_id = ensure_project(&tx, &pending.project, now)?;
let mut evidence_ids = HashMap::new();
let mut evidence_added = 0;
for source in &pending.evidence {
evidence_added += insert_evidence(&tx, input, source, project_id, now)?;
let (evidence_id, completed_at, existing_project): (i64, Option<i64>, String) = tx
.query_row(
"SELECT evidence.id, evidence.extraction_completed_at, projects.path
FROM evidence
JOIN projects ON projects.id = evidence.project_id
WHERE evidence.provider = ?1 AND evidence.session_id = ?2
AND evidence.entry_id = ?3 AND evidence.content_hash = ?4",
params![
input.provider.as_str(),
input.session_id,
source.entry_id,
source.content_hash
],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
)?;
if existing_project != pending.project {
return Err(MemoryError::ProjectConflict {
session_id: input.session_id.to_string(),
entry_id: source.entry_id.clone(),
existing_project,
}
.into());
}
if completed_at.is_some() {
continue;
}
let claimed = tx.execute(
"UPDATE evidence
SET extraction_lease_owner = ?1, extraction_lease_until = ?2
WHERE id = ?3 AND extraction_completed_at IS NULL
AND (extraction_lease_owner = ?1
OR extraction_lease_until IS NULL
OR extraction_lease_until <= ?4)",
params![owner, lease_until, evidence_id, now],
)?;
if claimed > 0 {
evidence_ids.insert(source.prompt_id.clone(), evidence_id);
}
}
pending
.evidence
.retain(|source| evidence_ids.contains_key(&source.prompt_id));
tx.commit()?;
Ok(ClaimedEvidence {
pending,
evidence_ids,
evidence_added,
owner,
})
}
fn insert_evidence(
conn: &Connection,
input: &RememberInput<'_>,
source: &SourceCandidate,
project_id: i64,
now: i64,
) -> anyhow::Result<usize> {
Ok(conn.execute(
"INSERT OR IGNORE INTO evidence(
project_id, provider, session_id, entry_id, role, observed_at,
source_path, content_hash, content_json, text, created_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
params![
project_id,
input.provider.as_str(),
input.session_id,
source.entry_id,
source.role,
source.observed_at,
source.source_path.to_string_lossy(),
source.content_hash,
source.content_json,
source.text,
now,
],
)?)
}
pub(super) fn leased_evidence_ids(conn: &Connection, owner: &str) -> anyhow::Result<HashSet<i64>> {
let mut stmt = conn.prepare(
"SELECT id FROM evidence
WHERE extraction_completed_at IS NULL AND extraction_lease_owner = ?1",
)?;
Ok(stmt
.query_map(params![owner], |row| row.get(0))?
.collect::<rusqlite::Result<HashSet<_>>>()?)
}
pub(super) fn complete_evidence(
conn: &Connection,
owner: &str,
completed_at: i64,
) -> anyhow::Result<usize> {
Ok(conn.execute(
"UPDATE evidence
SET extraction_completed_at = ?2,
extraction_lease_owner = NULL,
extraction_lease_until = NULL
WHERE extraction_completed_at IS NULL AND extraction_lease_owner = ?1",
params![owner, completed_at],
)?)
}
pub(super) fn release_evidence_lease(conn: &Connection, owner: &str) -> anyhow::Result<usize> {
Ok(conn.execute(
"UPDATE evidence
SET extraction_lease_owner = NULL, extraction_lease_until = NULL
WHERE extraction_completed_at IS NULL AND extraction_lease_owner = ?1",
params![owner],
)?)
}