#![cfg(feature = "memory")]
use serde::{Deserialize, Serialize};
use crate::index::IndexDb;
use super::types_memory::{MemoryRecord, Provenance, VerifyState};
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct WireMemoryEntry {
pub key: String,
pub value: String,
pub tags: Vec<String>,
pub created_at: i64,
pub updated_at: i64,
}
pub(crate) const MEMORY_PREVIEW_CHARS: usize = 200;
#[derive(Debug, thiserror::Error)]
pub(crate) enum MemoryOpError {
#[cfg(all(feature = "comms", any(unix, windows)))]
#[error("memory_by_key index not available")]
IndexUnavailable,
#[error("fjall {op}: {source}")]
Fjall {
op: &'static str,
source: fjall::Error,
},
#[error("serialize memory record: {0}")]
Serialize(rmp_serde::encode::Error),
}
impl From<MemoryOpError> for rmcp::ErrorData {
fn from(error: MemoryOpError) -> Self {
rmcp::ErrorData::internal_error(error.to_string(), None)
}
}
fn read_record(
idx: &IndexDb,
scope: &str,
vis_byte: u8,
owner: &str,
key: &str,
) -> Result<Option<MemoryRecord>, MemoryOpError> {
let raw_key = crate::index::keys::memory_by_key(scope, vis_byte, owner, key);
let bytes = idx
.memory_by_key
.get(raw_key)
.map_err(|source| MemoryOpError::Fjall { op: "get", source })?;
Ok(bytes.and_then(|b| rmp_serde::from_slice(&b).ok()))
}
fn write_record(
idx: &IndexDb,
scope: &str,
vis_byte: u8,
owner: &str,
key: &str,
record: &MemoryRecord,
) -> Result<(), MemoryOpError> {
let raw_key = crate::index::keys::memory_by_key(scope, vis_byte, owner, key);
let bytes = rmp_serde::to_vec_named(record).map_err(MemoryOpError::Serialize)?;
idx.memory_by_key
.insert(raw_key, bytes)
.map_err(|source| MemoryOpError::Fjall { op: "insert", source })
}
fn remove_record(idx: &IndexDb, scope: &str, vis_byte: u8, owner: &str, key: &str) -> Result<bool, MemoryOpError> {
let raw_key = crate::index::keys::memory_by_key(scope, vis_byte, owner, key);
let existed = idx
.memory_by_key
.get(raw_key.as_slice())
.map_err(|source| MemoryOpError::Fjall { op: "get", source })?
.is_some();
if existed {
idx.memory_by_key
.remove(raw_key)
.map_err(|source| MemoryOpError::Fjall { op: "remove", source })?;
}
Ok(existed)
}
fn now_micros() -> i64 {
crate::lance::now_micros()
}
pub(crate) fn put_core(
idx: &IndexDb,
scope: &str,
vis_byte: u8,
owner: &str,
key: &str,
value: &str,
tags: &[String],
) -> Result<(i64, i64), MemoryOpError> {
let now = now_micros();
let existing = read_record(idx, scope, vis_byte, owner, key)?;
let created_at = existing.map(|r| r.created_at).unwrap_or(now);
let record = MemoryRecord {
value: value.to_string(),
tags: tags.to_vec(),
created_at,
updated_at: now,
provenance: Provenance::default(),
verified: VerifyState::Unverified,
last_verified: 0,
importance: 0.0,
};
write_record(idx, scope, vis_byte, owner, key, &record)?;
Ok((created_at, now))
}
pub(crate) fn get_core(
idx: &IndexDb,
scope: &str,
vis_byte: u8,
owner: &str,
key: &str,
) -> Result<Option<WireMemoryEntry>, MemoryOpError> {
Ok(read_record(idx, scope, vis_byte, owner, key)?.map(|r| WireMemoryEntry {
key: key.to_string(),
value: r.value,
tags: r.tags,
created_at: r.created_at,
updated_at: r.updated_at,
}))
}
pub(crate) fn delete_core(
idx: &IndexDb,
scope: &str,
vis_byte: u8,
owner: &str,
key: &str,
) -> Result<bool, MemoryOpError> {
remove_record(idx, scope, vis_byte, owner, key)
}
fn preview(value: &str) -> String {
if value.len() > MEMORY_PREVIEW_CHARS {
format!(
"{}…",
value
.char_indices()
.nth(MEMORY_PREVIEW_CHARS)
.map(|(i, _)| &value[..i])
.unwrap_or(value)
)
} else {
value.to_string()
}
}
pub(crate) struct ListQuery<'a> {
pub vis_byte: u8,
pub owner: &'a str,
pub key_prefix: &'a str,
pub tag: Option<&'a str>,
pub limit: usize,
pub cursor: Option<&'a [u8]>,
}
pub(crate) struct ListResult {
pub entries: Vec<WireMemoryEntry>,
pub total: u32,
pub truncated: bool,
pub next_cursor: Option<Vec<u8>>,
}
pub(crate) fn list_core(idx: &IndexDb, scope: &str, query: &ListQuery<'_>) -> Result<ListResult, MemoryOpError> {
use std::ops::Bound;
use super::cursor::prefix_upper_bound;
let limit = query.limit;
let ns_prefix = crate::index::keys::memory_by_key_ns_prefix(scope, query.vis_byte, query.owner);
let upper = prefix_upper_bound(&ns_prefix);
let lower: Bound<Vec<u8>> = match query.cursor {
Some(k) => Bound::Excluded(k.to_vec()),
None => Bound::Included(ns_prefix.clone()),
};
let upper_bound: Bound<Vec<u8>> = match upper {
Some(b) => Bound::Excluded(b),
None => Bound::Unbounded,
};
let scan_cap = limit.saturating_mul(8).max(2_000);
let mut entries: Vec<WireMemoryEntry> = Vec::with_capacity(limit.min(64));
let mut total: usize = 0;
let mut last_emitted_key: Option<Vec<u8>> = None;
let mut has_more = false;
for guard in idx.memory_by_key.range::<Vec<u8>, _>((lower, upper_bound)) {
let (raw_key, raw_val) = guard
.into_inner()
.map_err(|source| MemoryOpError::Fjall { op: "iter", source })?;
let Some(key) = crate::index::keys::parse_memory_key_only(&raw_key) else {
continue;
};
if !key.starts_with(query.key_prefix) {
continue;
}
let Ok(record): Result<MemoryRecord, _> = rmp_serde::from_slice(&raw_val) else {
continue;
};
if let Some(tag) = query.tag
&& !record.tags.iter().any(|t| t == tag)
{
continue;
}
total += 1;
if entries.len() < limit {
entries.push(WireMemoryEntry {
key: key.to_string(),
value: preview(&record.value),
tags: record.tags,
created_at: record.created_at,
updated_at: record.updated_at,
});
last_emitted_key = Some(raw_key.to_vec());
} else {
has_more = true;
if total > scan_cap {
break;
}
}
}
let next_cursor = if has_more { last_emitted_key } else { None };
Ok(ListResult {
entries,
total: total as u32,
truncated: total > limit,
next_cursor,
})
}
#[cfg(all(feature = "comms", any(unix, windows)))]
pub(crate) fn run_memory_op(
idx: &IndexDb,
scope: &str,
op: &crate::comms::memory_proto::MemoryOp,
) -> Result<crate::comms::memory_proto::MemoryOutcome, MemoryOpError> {
use crate::comms::memory_proto::{MemoryOp, MemoryOutcome};
match op {
MemoryOp::Get { vis_byte, owner, key } => Ok(MemoryOutcome::Got(get_core(idx, scope, *vis_byte, owner, key)?)),
MemoryOp::Put {
vis_byte,
owner,
key,
value,
tags,
} => {
let (created_at, updated_at) = put_core(idx, scope, *vis_byte, owner, key, value, tags)?;
Ok(MemoryOutcome::Put { created_at, updated_at })
}
MemoryOp::List {
vis_byte,
owner,
prefix,
tag,
limit,
cursor,
} => {
let result = list_core(
idx,
scope,
&ListQuery {
vis_byte: *vis_byte,
owner,
key_prefix: prefix,
tag: tag.as_deref(),
limit: *limit as usize,
cursor: cursor.as_deref(),
},
)?;
Ok(MemoryOutcome::Listed {
entries: result.entries,
total: result.total,
truncated: result.truncated,
next_cursor: result.next_cursor,
})
}
MemoryOp::Delete { vis_byte, owner, key } => Ok(MemoryOutcome::Deleted {
deleted: delete_core(idx, scope, *vis_byte, owner, key)?,
}),
}
}