pub mod transport;
use std::collections::HashSet;
use std::sync::Arc;
use async_trait::async_trait;
use rusqlite::OptionalExtension;
use uuid::Uuid;
use khive_storage::attachment::AttachmentSubstrate;
use khive_storage::error::{StorageError, WriterTaskRequestState};
use khive_storage::note::{
FilterOp, Note, NoteFilter, NoteInstantSeekAfter, NoteKeyCursor, NoteSeekAfter, NoteTagMode,
SortDir,
};
use khive_storage::types::{
BatchWriteSummary, BoundedCount, DeleteMode, Page, PageRequest, SeekCursor, SeekPage,
SqlStatement, SqlValue,
};
use khive_storage::NoteStore;
use khive_storage::{StorageCapability, StorageResult};
use crate::error::SqliteError;
use crate::pool::ConnectionPool;
use crate::sql_bridge::bind_params;
use crate::stores::attachment::delete_record_attachments_statement;
use crate::writer_task::{execute_wrapped_transaction, WriterTaskHandle};
fn map_err(e: rusqlite::Error, op: &'static str) -> StorageError {
StorageError::driver(StorageCapability::Notes, op, e)
}
fn map_sqlite_err(e: SqliteError, op: &'static str) -> StorageError {
StorageError::driver(StorageCapability::Notes, op, e)
}
const NAMESPACE_COUNT_CHUNK_SIZE: usize = 500;
pub fn note_key_prefix_successor(prefix: &str) -> Option<String> {
let mut chars: Vec<char> = prefix.chars().collect();
while let Some(last) = chars.pop() {
if last == char::MAX {
continue;
}
let next = if last == '\u{d7ff}' {
'\u{e000}'
} else {
char::from_u32(u32::from(last) + 1).expect("incremented non-max scalar")
};
chars.push(next);
return Some(chars.into_iter().collect());
}
None
}
pub const NOTE_UPSERT_SQL: &str = "INSERT INTO notes \
(id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
properties, created_at, updated_at, deleted_at, key) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14) \
ON CONFLICT(id) DO UPDATE SET \
namespace = excluded.namespace, \
kind = excluded.kind, \
status = excluded.status, \
name = excluded.name, \
content = excluded.content, \
salience = excluded.salience, \
decay_factor = excluded.decay_factor, \
expires_at = excluded.expires_at, \
properties = excluded.properties, \
updated_at = excluded.updated_at, \
deleted_at = excluded.deleted_at";
pub const NOTE_INSERT_IF_ABSENT_SQL: &str = "INSERT INTO notes \
(id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
properties, created_at, updated_at, deleted_at, key) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14) \
ON CONFLICT(id) DO NOTHING";
pub fn note_insert_if_absent_statement(note: &Note) -> SqlStatement {
let mut statement = note_upsert_statement(note);
statement.sql = NOTE_INSERT_IF_ABSENT_SQL.to_string();
statement.label = Some("note-insert-if-absent".to_string());
statement
}
pub fn note_insert_keyed_statement(note: &Note) -> SqlStatement {
let mut statement = note_upsert_statement(note);
statement.sql = "INSERT INTO notes \
(id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
properties, created_at, updated_at, deleted_at, key) \
VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14) \
ON CONFLICT(namespace,kind,key) WHERE key IS NOT NULL AND deleted_at IS NULL DO NOTHING"
.into();
statement.label = Some("note-keyed-create".into());
statement
}
pub fn note_upsert_statement(note: &Note) -> SqlStatement {
let properties_str = note
.properties
.as_ref()
.map(|v| serde_json::to_string(v).unwrap_or_default());
SqlStatement {
sql: NOTE_UPSERT_SQL.to_string(),
params: vec![
SqlValue::Text(note.id.to_string()),
SqlValue::Text(note.namespace.clone()),
SqlValue::Text(note.kind.to_string()),
SqlValue::Text(note.status.clone()),
match ¬e.name {
Some(n) => SqlValue::Text(n.clone()),
None => SqlValue::Null,
},
SqlValue::Text(note.content.clone()),
match note.salience {
Some(s) => SqlValue::Float(s),
None => SqlValue::Null,
},
match note.decay_factor {
Some(d) => SqlValue::Float(d),
None => SqlValue::Null,
},
match note.expires_at {
Some(e) => SqlValue::Integer(e),
None => SqlValue::Null,
},
match properties_str {
Some(p) => SqlValue::Text(p),
None => SqlValue::Null,
},
SqlValue::Integer(note.created_at),
SqlValue::Integer(note.updated_at),
match note.deleted_at {
Some(d) => SqlValue::Integer(d),
None => SqlValue::Null,
},
match ¬e.key {
Some(key) => SqlValue::Text(key.clone()),
None => SqlValue::Null,
},
],
label: Some("note-upsert".to_string()),
}
}
pub fn note_replace_if_unchanged_statement(
note: &Note,
expected_updated_at: i64,
expected_deleted_at: Option<i64>,
) -> SqlStatement {
let properties_str = note
.properties
.as_ref()
.map(|v| serde_json::to_string(v).unwrap_or_default());
SqlStatement {
sql: "UPDATE notes SET \
namespace = ?1, kind = ?2, status = ?3, name = ?4, content = ?5, \
salience = ?6, decay_factor = ?7, expires_at = ?8, properties = ?9, \
updated_at = ?10, deleted_at = ?11 \
WHERE id = ?12 AND updated_at = ?13 AND deleted_at IS ?14 \
AND ?10 > updated_at"
.to_string(),
params: vec![
SqlValue::Text(note.namespace.clone()),
SqlValue::Text(note.kind.to_string()),
SqlValue::Text(note.status.clone()),
match ¬e.name {
Some(name) => SqlValue::Text(name.clone()),
None => SqlValue::Null,
},
SqlValue::Text(note.content.clone()),
match note.salience {
Some(value) => SqlValue::Float(value),
None => SqlValue::Null,
},
match note.decay_factor {
Some(value) => SqlValue::Float(value),
None => SqlValue::Null,
},
match note.expires_at {
Some(value) => SqlValue::Integer(value),
None => SqlValue::Null,
},
match properties_str {
Some(value) => SqlValue::Text(value),
None => SqlValue::Null,
},
SqlValue::Integer(note.updated_at),
match note.deleted_at {
Some(value) => SqlValue::Integer(value),
None => SqlValue::Null,
},
SqlValue::Text(note.id.to_string()),
SqlValue::Integer(expected_updated_at),
match expected_deleted_at {
Some(value) => SqlValue::Integer(value),
None => SqlValue::Null,
},
],
label: Some("note-replace-if-unchanged".to_string()),
}
}
pub fn note_metadata_replace_if_unchanged_statement(
note: &Note,
expected_updated_at: i64,
expected_deleted_at: Option<i64>,
) -> SqlStatement {
let mut statement =
note_replace_if_unchanged_statement(note, expected_updated_at, expected_deleted_at);
statement.sql = "UPDATE notes SET status=?3, name=?4, salience=?6, decay_factor=?7, expires_at=?8, updated_at=?10 \
WHERE id=?12 AND updated_at=?13 AND deleted_at IS ?14 AND ?10 > updated_at \
AND namespace=?1 AND kind=?2 AND content=?5 AND properties IS ?9 AND deleted_at IS ?11".into();
statement.label = Some("stream-note-metadata-cas".into());
statement
}
pub fn note_update_properties_statement(
id: Uuid,
properties: &Option<serde_json::Value>,
updated_at: i64,
) -> SqlStatement {
let properties_str = properties
.as_ref()
.map(|v| serde_json::to_string(v).unwrap_or_default());
SqlStatement {
sql: "UPDATE notes SET properties = ?1, updated_at = ?2 \
WHERE id = ?3 AND deleted_at IS NULL"
.to_string(),
params: vec![
match properties_str {
Some(p) => SqlValue::Text(p),
None => SqlValue::Null,
},
SqlValue::Integer(updated_at),
SqlValue::Text(id.to_string()),
],
label: Some("note-update-properties".to_string()),
}
}
pub fn note_set_property_statement(
id: Uuid,
key: &str,
value: &serde_json::Value,
updated_at: i64,
) -> Result<SqlStatement, StorageError> {
if key.contains('\0') {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "set_note_property".into(),
message: "property key must not contain U+0000".to_string(),
});
}
let path = format!("$.{}", serde_json::Value::String(key.to_string()));
Ok(SqlStatement {
sql: "UPDATE notes \
SET properties = json_set(COALESCE(properties, '{}'), ?1, json(?2)), \
updated_at = ?3 \
WHERE id = ?4 AND deleted_at IS NULL \
AND (properties IS NULL OR json_type(properties) = 'object')"
.to_string(),
params: vec![
SqlValue::Text(path),
SqlValue::Text(value.to_string()),
SqlValue::Integer(updated_at),
SqlValue::Text(id.to_string()),
],
label: Some("note-set-property".to_string()),
})
}
pub fn note_soft_delete_statement(id: Uuid, deleted_at: i64) -> SqlStatement {
SqlStatement {
sql: "UPDATE notes SET status = 'deleted', deleted_at = ?1 \
WHERE id = ?2 AND deleted_at IS NULL"
.to_string(),
params: vec![
SqlValue::Integer(deleted_at),
SqlValue::Text(id.to_string()),
],
label: Some("note-delete-soft".to_string()),
}
}
pub fn note_hard_delete_statement(id: Uuid) -> SqlStatement {
SqlStatement {
sql: "DELETE FROM notes WHERE id = ?1".to_string(),
params: vec![SqlValue::Text(id.to_string())],
label: Some("note-delete-hard".to_string()),
}
}
pub struct SqlNoteStore {
pool: Arc<ConnectionPool>,
writer_task: Option<WriterTaskHandle>,
}
impl SqlNoteStore {
pub fn new(pool: Arc<ConnectionPool>, _is_file_backed: bool) -> Self {
let writer_task = pool.writer_task_handle().ok().flatten();
Self { pool, writer_task }
}
fn current_writer_task(
&self,
operation: &'static str,
) -> Result<Option<WriterTaskHandle>, StorageError> {
self.pool
.writer_task_for_write(self.writer_task.as_ref(), operation)
}
async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
where
F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
R: Send + 'static,
{
if let Some(writer_task) = self.current_writer_task(op)? {
return writer_task
.send_bounded(move |conn| f(conn).map_err(|e| map_err(e, op)))
.await;
}
self.pool
.record_direct_route(crate::timeout_sink::Site::DirectRouteNote);
let pool = Arc::clone(&self.pool);
tokio::task::spawn_blocking(move || {
let guard = pool.try_writer().map_err(|e| map_sqlite_err(e, op))?;
f(guard.conn()).map_err(|e| map_err(e, op))
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
}
async fn with_writer_tx<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
where
F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
R: Send + 'static,
{
self.with_writer_tx_storage(op, move |conn| f(conn).map_err(|error| map_err(error, op)))
.await
}
async fn with_writer_tx_storage<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
where
F: FnOnce(&rusqlite::Connection) -> Result<R, StorageError> + Send + 'static,
R: Send + 'static,
{
if let Some(writer_task) = self.current_writer_task(op)? {
return writer_task.send_bounded(f).await;
}
self.pool
.record_direct_route(crate::timeout_sink::Site::DirectRouteNote);
let pool = Arc::clone(&self.pool);
tokio::task::spawn_blocking(move || {
let guard = pool.try_writer().map_err(|e| map_sqlite_err(e, op))?;
let conn = guard.conn();
if !conn.is_autocommit() {
pool.retire_pooled_writer(conn);
return Err(StorageError::WriterTaskTerminated {
request_state: WriterTaskRequestState::SideEffectsUnknown,
});
}
if let Err(begin_error) = conn.execute_batch("BEGIN IMMEDIATE") {
if !conn.is_autocommit() {
pool.retire_pooled_writer(conn);
return Err(StorageError::WriterTaskTerminated {
request_state: WriterTaskRequestState::SideEffectsUnknown,
});
}
return Err(map_err(begin_error, op));
}
let (result, terminal_state) = execute_wrapped_transaction(conn, op, f);
if terminal_state.is_some() {
pool.retire_pooled_writer(conn);
}
result
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
}
async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
where
F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
R: Send + 'static,
{
super::run_pooled_store_read(
Arc::clone(&self.pool),
StorageCapability::Notes,
op,
move |conn| f(conn).map_err(|error| map_err(error, op)),
)
.await
}
}
fn read_note(row: &rusqlite::Row<'_>) -> Result<Note, rusqlite::Error> {
let id_str: String = row.get(0)?;
let namespace: String = row.get(1)?;
let kind: String = row.get(2)?;
let status: String = row.get(3)?;
let name: Option<String> = row.get(4)?;
let content: String = row.get(5)?;
let salience: Option<f64> = row.get(6)?;
let decay_factor: Option<f64> = row.get(7)?;
let expires_at: Option<i64> = row.get(8)?;
let properties_str: Option<String> = row.get(9)?;
let created_at: i64 = row.get(10)?;
let updated_at: i64 = row.get(11)?;
let deleted_at: Option<i64> = row.get(12)?;
let key: Option<String> = row.get(13)?;
let version: i64 = row.get(14)?;
let id = parse_uuid(&id_str)?;
let properties = properties_str
.map(|s| {
serde_json::from_str(&s).map_err(|e| {
rusqlite::Error::FromSqlConversionFailure(
9,
rusqlite::types::Type::Text,
Box::new(e),
)
})
})
.transpose()?;
Ok(Note {
id,
namespace,
kind,
status,
name,
content,
salience,
decay_factor,
expires_at,
properties,
created_at,
updated_at,
deleted_at,
key,
version,
})
}
fn parse_uuid(s: &str) -> Result<Uuid, rusqlite::Error> {
Uuid::parse_str(s).map_err(|e| {
rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(e))
})
}
fn query_note_page_snapshot(
conn: &rusqlite::Connection,
operation: &'static str,
namespace: &str,
count_sql: &str,
count_params: &[Box<dyn rusqlite::types::ToSql>],
data_sql: &str,
data_params: &[Box<dyn rusqlite::types::ToSql>],
) -> Result<Page<Note>, rusqlite::Error> {
let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Deferred)?;
let total: i64 = {
let mut stmt = tx.prepare(count_sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
count_params.iter().map(|param| param.as_ref()).collect();
stmt.query_row(param_refs.as_slice(), |row| row.get(0))?
};
#[cfg(test)]
tests::page_snapshot_seam::hook(operation, namespace);
#[cfg(not(test))]
let _ = (operation, namespace);
let items = {
let mut stmt = tx.prepare(data_sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
data_params.iter().map(|param| param.as_ref()).collect();
let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
rows.collect::<Result<Vec<_>, _>>()?
};
tx.commit()?;
Ok(Page {
items,
total: Some(total as u64),
})
}
fn batch_upsert_notes(
conn: &rusqlite::Connection,
notes: &[Note],
attempted: u64,
) -> Result<BatchWriteSummary, rusqlite::Error> {
let mut summary = BatchWriteSummary {
attempted,
..BatchWriteSummary::default()
};
let mut stmt = conn.prepare_cached(NOTE_UPSERT_SQL)?;
for (index, note) in notes.iter().enumerate() {
let id_str = note.id.to_string();
let kind_str = note.kind.to_string();
let status_str = note.status.clone();
let properties_str = note
.properties
.as_ref()
.map(|v| serde_json::to_string(v).unwrap_or_default());
match stmt.execute(rusqlite::params![
id_str,
¬e.namespace,
kind_str,
status_str,
¬e.name,
note.content,
note.salience,
note.decay_factor,
note.expires_at,
properties_str,
note.created_at,
note.updated_at,
note.deleted_at,
note.key,
]) {
Ok(_) => {
assign_note_seq(conn, &id_str)?;
summary.affected = summary.affected.saturating_add(1);
}
Err(e) => {
let (class, retryability) = super::classify_batch_sqlite_error(&e);
summary.record_failure(index, Some(id_str), class, retryability, e.to_string());
}
}
}
Ok(summary)
}
fn assign_note_seq(conn: &rusqlite::Connection, note_id: &str) -> Result<(), rusqlite::Error> {
conn.execute(
"INSERT OR IGNORE INTO notes_seq (note_id) VALUES (?1)",
rusqlite::params![note_id],
)?;
Ok(())
}
fn build_note_where(
namespace: &str,
kind: Option<&str>,
) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
let mut conditions: Vec<String> = vec![
"namespace = ?1".to_string(),
"deleted_at IS NULL".to_string(),
];
let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![Box::new(namespace.to_string())];
if let Some(k) = kind {
params.push(Box::new(k.to_string()));
conditions.push(format!("kind = ?{}", params.len()));
}
let clause = format!(" WHERE {}", conditions.join(" AND "));
(clause, params)
}
fn build_note_where_for_namespaces(
namespaces: &[String],
kind: Option<&str>,
) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = namespaces
.iter()
.map(|namespace| -> Box<dyn rusqlite::types::ToSql> { Box::new(namespace.clone()) })
.collect();
let namespace_condition = match namespaces.len() {
0 => "0".to_string(),
1 => "namespace = ?1".to_string(),
_ => {
let placeholders: Vec<String> =
(1..=namespaces.len()).map(|i| format!("?{i}")).collect();
format!("namespace IN ({})", placeholders.join(", "))
}
};
let mut conditions = vec![namespace_condition, "deleted_at IS NULL".to_string()];
if let Some(kind) = kind {
params.push(Box::new(kind.to_string()));
conditions.push(format!("kind = ?{}", params.len()));
}
let clause = format!(" WHERE {}", conditions.join(" AND "));
(clause, params)
}
fn validate_json_path(path: &str) -> Result<(), StorageError> {
let valid = path.starts_with("$.")
&& path[2..].split('.').all(|part| {
!part.is_empty() && part.chars().all(|c| c.is_ascii_alphanumeric() || c == '_')
});
if valid {
Ok(())
} else {
Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered".into(),
message: format!("invalid JSON path for note filter: {path:?}"),
})
}
}
fn json_extract_expr(path: &str) -> String {
format!("json_extract(properties, '{path}')")
}
fn json_type_expr(path: &str) -> String {
format!("json_type(properties, '{path}')")
}
fn note_filter_page_order_clause(filter: &NoteFilter) -> String {
if filter.unordered {
return String::new();
}
match &filter.order_by {
Some((path, dir)) => {
if filter.order_by_instant {
let expr = json_extract_expr(path);
return format!(" ORDER BY khive_rfc3339_key({expr}) ASC, {expr} ASC, id ASC");
}
let dir_str = match dir {
SortDir::Asc => "ASC",
SortDir::Desc => "DESC",
};
format!(
" ORDER BY {} {dir_str}, id {dir_str}",
json_extract_expr(path)
)
}
None => " ORDER BY created_at DESC, id ASC".to_string(),
}
}
fn json_type_literal(value: &SqlValue) -> Result<&str, rusqlite::Error> {
const JSON_TYPES: [&str; 8] = [
"true", "false", "integer", "real", "text", "array", "object", "null",
];
match value {
SqlValue::Text(s) if JSON_TYPES.contains(&s.as_str()) => Ok(s.as_str()),
other => Err(rusqlite::Error::InvalidParameterName(format!(
"json_type comparison value must be one of SQLite's json_type strings \
({JSON_TYPES:?}), got {other:?}"
))),
}
}
fn text_prefix_upper_bound(prefix: &str) -> Option<String> {
let mut chars: Vec<char> = prefix.chars().collect();
while let Some(last) = chars.pop() {
let mut next = u32::from(last) + 1;
if (0xD800..=0xDFFF).contains(&next) {
next = 0xE000;
}
if let Some(next) = char::from_u32(next) {
chars.push(next);
return Some(chars.into_iter().collect());
}
}
None
}
fn sql_value_param(value: &SqlValue) -> Result<Box<dyn rusqlite::types::ToSql>, rusqlite::Error> {
Ok(match value {
SqlValue::Null => Box::new(Option::<String>::None),
SqlValue::Bool(v) => Box::new(*v as i64),
SqlValue::Integer(v) => Box::new(*v),
SqlValue::Float(v) => Box::new(*v),
SqlValue::Text(v) => Box::new(v.clone()),
SqlValue::Blob(v) => Box::new(v.clone()),
SqlValue::Json(v) => Box::new(
serde_json::to_string(v)
.map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?,
),
SqlValue::Uuid(v) => Box::new(v.to_string()),
SqlValue::Timestamp(v) => Box::new(v.timestamp_micros()),
})
}
fn build_note_filter_where(
namespace: &str,
filter: &NoteFilter,
) -> Result<(String, Vec<Box<dyn rusqlite::types::ToSql>>), rusqlite::Error> {
let (ns_condition, ns_params): (String, Vec<Box<dyn rusqlite::types::ToSql>>) =
if !filter.namespaces.is_empty() {
let placeholders: Vec<String> = (1..=filter.namespaces.len())
.map(|i| format!("?{i}"))
.collect();
let params: Vec<Box<dyn rusqlite::types::ToSql>> = filter
.namespaces
.iter()
.map(|ns| -> Box<dyn rusqlite::types::ToSql> { Box::new(ns.clone()) })
.collect();
(
format!("namespace IN ({})", placeholders.join(", ")),
params,
)
} else {
(
"namespace = ?1".to_string(),
vec![Box::new(namespace.to_string())],
)
};
let mut conditions = vec![ns_condition, "deleted_at IS NULL".to_string()];
let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = ns_params;
if let Some(kind) = &filter.kind {
params.push(Box::new(kind.clone()));
conditions.push(format!("kind = ?{}", params.len()));
}
if let Some(since) = filter.min_updated_at {
params.push(Box::new(since));
conditions.push(format!("updated_at >= ?{}", params.len()));
}
if !filter.tags.is_empty() {
let mut tag_predicates = Vec::new();
for tag in &filter.tags {
params.push(Box::new(tag.clone()));
tag_predicates.push(format!(
"EXISTS (SELECT 1 FROM json_each(CASE WHEN json_type(properties,'$.tags')='array' \
THEN json_extract(properties,'$.tags') ELSE '[]' END) AS tag \
WHERE tag.type='text' AND tag.value = ?{} COLLATE NOCASE)",
params.len()
));
}
let join = match filter.tag_mode {
NoteTagMode::Any => " OR ",
NoteTagMode::All => " AND ",
};
conditions.push(format!("({})", tag_predicates.join(join)));
}
for pf in &filter.property_filters {
match &pf.op {
FilterOp::Rfc3339Valid => {
let expr = json_extract_expr(&pf.json_path);
conditions.push(format!("khive_rfc3339_key({expr}) IS NOT NULL"));
}
FilterOp::Rfc3339Gte | FilterOp::Rfc3339Lte => {
let instant = match &pf.value {
SqlValue::Timestamp(instant) => *instant,
SqlValue::Text(text) => {
text.parse::<chrono::DateTime<chrono::Utc>>()
.map_err(|error| {
rusqlite::Error::ToSqlConversionFailure(Box::new(error))
})?
}
_ => {
return Err(rusqlite::Error::ToSqlConversionFailure(
"RFC 3339 filters require a timestamp or text value".into(),
));
}
};
let expr = json_extract_expr(&pf.json_path);
let op = if matches!(&pf.op, FilterOp::Rfc3339Gte) {
">="
} else {
"<="
};
params.push(Box::new(crate::pool::rfc3339_instant_key(instant)));
conditions.push(format!("khive_rfc3339_key({expr}) {op} ?{}", params.len()));
}
FilterOp::EqOrMissing => {
let expr = json_extract_expr(&pf.json_path);
params.push(sql_value_param(&pf.value)?);
conditions.push(format!(
"({expr} = ?{n} OR {expr} IS NULL)",
n = params.len()
));
}
FilterOp::EqOrMissingIndexed => {
let expr = json_extract_expr(&pf.json_path);
params.push(sql_value_param(&pf.value)?);
conditions.push(format!("ifnull({expr}, '') = ?{}", params.len()));
}
FilterOp::TextEqOrNonText => {
let expr = json_extract_expr(&pf.json_path);
let type_expr = json_type_expr(&pf.json_path);
params.push(sql_value_param(&pf.value)?);
let n = params.len();
conditions.push(format!(
"CASE WHEN {type_expr} = 'text' THEN {expr} ELSE ?{n} END = ?{n}"
));
}
FilterOp::TextInOrNonText(values) => {
let expr = json_extract_expr(&pf.json_path);
let type_expr = json_type_expr(&pf.json_path);
let mut placeholders = Vec::with_capacity(values.len());
for value in values {
params.push(sql_value_param(value)?);
placeholders.push(format!("?{}", params.len()));
}
let text_match = if placeholders.is_empty() {
"0".to_string()
} else {
format!("{expr} IN ({})", placeholders.join(", "))
};
conditions.push(format!(
"CASE WHEN {type_expr} = 'text' THEN {text_match} ELSE 1 END"
));
}
FilterOp::JsonTypeEq => {
let type_expr = json_type_expr(&pf.json_path);
params.push(sql_value_param(&pf.value)?);
conditions.push(format!("{type_expr} = ?{}", params.len()));
}
FilterOp::JsonTypeMissing => {
let type_expr = json_type_expr(&pf.json_path);
conditions.push(format!("{type_expr} IS NULL"));
}
FilterOp::JsonTypeMissingOrNullIndexed => {
let expr = json_extract_expr(&pf.json_path);
let type_expr = json_type_expr(&pf.json_path);
conditions.push(format!(
"ifnull({expr}, '') = '' AND ({type_expr} IS NULL OR {type_expr} = 'null')"
));
}
FilterOp::EqOrLegacyIndexed => {
let expr = json_extract_expr(&pf.json_path);
let type_expr = json_type_expr(&pf.json_path);
params.push(sql_value_param(&pf.value)?);
let n = params.len();
conditions.push(format!(
"ifnull({expr}, '') IN (?{n}, '') AND \
({type_expr} IS NULL OR {type_expr} = 'null' OR ifnull({expr}, '') != '')"
));
}
FilterOp::JsonTypeNeMissing => {
let type_expr = json_type_expr(&pf.json_path);
let literal = json_type_literal(&pf.value)?;
conditions.push(format!(
"({type_expr} IS NULL OR {type_expr} != '{literal}')"
));
}
FilterOp::In(values) => {
let expr = json_extract_expr(&pf.json_path);
if values.is_empty() {
conditions.push("0".to_string());
continue;
}
let mut placeholders = Vec::with_capacity(values.len());
for v in values {
params.push(sql_value_param(v)?);
placeholders.push(format!("?{}", params.len()));
}
conditions.push(format!("{expr} IN ({})", placeholders.join(", ")));
}
FilterOp::TextStartsWithIndexed => {
let expr = json_extract_expr(&pf.json_path);
let SqlValue::Text(prefix) = &pf.value else {
return Err(rusqlite::Error::ToSqlConversionFailure(
"TextStartsWithIndexed takes a text prefix in PropertyFilter.value".into(),
));
};
params.push(Box::new(prefix.clone()));
let lower = params.len();
match text_prefix_upper_bound(prefix) {
Some(upper) => {
params.push(Box::new(upper));
conditions.push(format!(
"({expr} >= ?{lower} AND {expr} < ?{})",
params.len()
));
}
None => {
let type_expr = json_type_expr(&pf.json_path);
conditions.push(format!("({type_expr} = 'text' AND {expr} >= ?{lower})"));
}
}
}
FilterOp::NotInOrMissing(values) => {
let expr = json_extract_expr(&pf.json_path);
if values.is_empty() {
continue;
}
let mut placeholders = Vec::with_capacity(values.len());
for v in values {
params.push(sql_value_param(v)?);
placeholders.push(format!("?{}", params.len()));
}
conditions.push(format!(
"({expr} IS NULL OR {expr} NOT IN ({}))",
placeholders.join(", ")
));
}
_ => {
let expr = json_extract_expr(&pf.json_path);
let op = match pf.op {
FilterOp::Eq => "=",
FilterOp::Ne => "!=",
FilterOp::Lt => "<",
FilterOp::Lte => "<=",
FilterOp::Gt => ">",
FilterOp::Gte => ">=",
FilterOp::EqOrMissing
| FilterOp::EqOrMissingIndexed
| FilterOp::TextEqOrNonText
| FilterOp::TextInOrNonText(_)
| FilterOp::JsonTypeEq
| FilterOp::JsonTypeMissing
| FilterOp::JsonTypeMissingOrNullIndexed
| FilterOp::EqOrLegacyIndexed
| FilterOp::JsonTypeNeMissing
| FilterOp::In(_)
| FilterOp::NotInOrMissing(_)
| FilterOp::TextStartsWithIndexed => {
unreachable!()
}
FilterOp::Rfc3339Valid | FilterOp::Rfc3339Gte | FilterOp::Rfc3339Lte => {
unreachable!()
}
};
params.push(sql_value_param(&pf.value)?);
conditions.push(format!("{expr} {op} ?{}", params.len()));
}
}
}
if let Some(min_ts) = filter.min_created_at {
params.push(Box::new(min_ts));
conditions.push(format!("created_at >= ?{}", params.len()));
}
Ok((format!(" WHERE {}", conditions.join(" AND ")), params))
}
fn comm_filter_index_clause(filter: &NoteFilter, where_sql: &str) -> &'static str {
if filter.kind.as_deref() != Some("message") {
return "";
}
let Some(predicate) = where_sql.strip_prefix(" WHERE ") else {
return "";
};
let terms: Vec<_> = predicate.split(" AND ").collect();
let numbered_param = |value: &str| {
value.strip_prefix('?').is_some_and(|number| {
!number.is_empty() && number.bytes().all(|byte| byte.is_ascii_digit())
})
};
let equality = |prefix: &str| {
terms
.iter()
.any(|term| term.strip_prefix(prefix).is_some_and(numbered_param))
};
let namespace = equality("namespace = ")
|| terms.iter().any(|term| {
term.strip_prefix("namespace IN (")
.and_then(|term| term.strip_suffix(')'))
.is_some_and(|values| values.split(", ").all(numbered_param))
});
let recipient = equality("ifnull(json_extract(properties, '$.to_actor'), '') = ")
|| terms.contains(&"ifnull(json_extract(properties, '$.to_actor'), '') = ''")
|| terms.iter().any(|term| {
term.strip_prefix("ifnull(json_extract(properties, '$.to_actor'), '') IN (")
.and_then(|term| term.strip_suffix(", '')"))
.is_some_and(numbered_param)
});
if !namespace
|| !terms.contains(&"deleted_at IS NULL")
|| !equality("kind = ")
|| !equality("json_extract(properties, '$.direction') = ")
|| !recipient
{
return "";
}
if terms.contains(
&"(json_type(properties, '$.read') IS NULL OR json_type(properties, '$.read') != 'true')",
) {
let typed_exact_recipient =
equality("ifnull(json_extract(properties, '$.to_actor'), '') = ")
&& equality("json_type(properties, '$.to_actor') = ")
&& filter.property_filters.iter().any(|property| {
property.json_path == "$.to_actor"
&& matches!(property.op, FilterOp::JsonTypeEq)
&& matches!(&property.value, SqlValue::Text(value) if value == "text")
});
if typed_exact_recipient {
" INDEXED BY idx_notes_unread_probe_recipient_type_direction"
} else {
" INDEXED BY idx_notes_unread_probe_recipient_direction"
}
} else {
" INDEXED BY idx_notes_message_recipient_direction"
}
}
fn build_note_filter_read_clause(
namespace: &str,
filter: &NoteFilter,
) -> Result<(String, Vec<Box<dyn rusqlite::types::ToSql>>), rusqlite::Error> {
let (where_sql, params) = build_note_filter_where(namespace, filter)?;
let index_clause = comm_filter_index_clause(filter, &where_sql);
Ok((format!("{index_clause}{where_sql}"), params))
}
const NOTE_COLUMNS: &str = "id, namespace, kind, status, name, content, salience, decay_factor, \
expires_at, properties, created_at, updated_at, deleted_at, key, version";
fn fetch_notes_after_instant(
conn: &rusqlite::Connection,
namespace: &str,
base_filter: &NoteFilter,
after: &NoteInstantSeekAfter,
limit: i64,
) -> Result<Vec<Note>, rusqlite::Error> {
let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
params.push(Box::new(after.value.clone()));
let value_idx = params.len();
params.push(Box::new(after.id.to_string()));
let id_idx = params.len();
params.push(Box::new(limit));
let limit_idx = params.len();
let (path, _) = base_filter
.order_by
.as_ref()
.expect("instant cursor requires order_by");
let expr = json_extract_expr(path);
let order_clause = note_filter_page_order_clause(base_filter);
let sql = format!(
"SELECT {NOTE_COLUMNS} FROM notes{where_sql} \
AND (khive_rfc3339_key({expr}), {expr}, id) > \
(khive_rfc3339_key(?{value_idx}), ?{value_idx}, ?{id_idx}) \
{order_clause} LIMIT ?{limit_idx}"
);
let mut stmt = conn.prepare_cached(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|param| param.as_ref()).collect();
let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
rows.collect()
}
fn fetch_notes_after(
conn: &rusqlite::Connection,
namespace: &str,
base_filter: &NoteFilter,
after: &NoteSeekAfter,
limit: i64,
) -> Result<Vec<Note>, rusqlite::Error> {
let mut items = Vec::new();
if limit <= 0 {
return Ok(items);
}
{
let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
params.push(Box::new(after.created_at));
let ts_idx = params.len();
params.push(Box::new(after.id.to_string()));
let id_idx = params.len();
params.push(Box::new(limit));
let limit_idx = params.len();
let sql = format!(
"SELECT {NOTE_COLUMNS} FROM notes{where_sql} AND created_at = ?{ts_idx} \
AND id > ?{id_idx} ORDER BY id ASC LIMIT ?{limit_idx}"
);
let mut stmt = conn.prepare_cached(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|p| p.as_ref()).collect();
let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
for row in rows {
items.push(row?);
}
}
let remaining = limit - items.len() as i64;
if remaining > 0 {
let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
params.push(Box::new(after.created_at));
let ts_idx = params.len();
params.push(Box::new(remaining));
let limit_idx = params.len();
let sql = format!(
"SELECT {NOTE_COLUMNS} FROM notes{where_sql} AND created_at < ?{ts_idx} \
ORDER BY created_at DESC, id ASC LIMIT ?{limit_idx}"
);
let mut stmt = conn.prepare_cached(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|p| p.as_ref()).collect();
let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
for row in rows {
items.push(row?);
}
}
Ok(items)
}
fn execute_filtered_note_property_patch(
conn: &rusqlite::Connection,
id: Uuid,
namespace: &str,
filter: &NoteFilter,
json_path: &str,
value_json: &str,
updated_at: i64,
) -> Result<usize, rusqlite::Error> {
let (where_clause, mut params) = build_note_filter_where(namespace, filter)?;
let base = params.len();
let sql = format!(
"UPDATE notes SET properties = json_set(COALESCE(properties, '{{}}'), ?{p1}, json(?{p2})), \
updated_at = ?{p3} {where_clause} \
AND (properties IS NULL OR json_type(properties) = 'object') AND id = ?{p4}",
p1 = base + 1,
p2 = base + 2,
p3 = base + 3,
p4 = base + 4,
);
params.push(Box::new(json_path.to_string()));
params.push(Box::new(value_json.to_string()));
params.push(Box::new(updated_at));
params.push(Box::new(id.to_string()));
let mut stmt = conn.prepare_cached(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|param| param.as_ref()).collect();
stmt.execute(param_refs.as_slice())
}
#[async_trait]
impl NoteStore for SqlNoteStore {
async fn get_live_notes_by_key(
&self,
namespace: &str,
key: &str,
kind: Option<&str>,
) -> StorageResult<Vec<Note>> {
let namespace = namespace.to_owned();
let key = key.to_owned();
let kind = kind.map(str::to_owned);
self.with_reader("get_live_notes_by_key", move |conn| {
let (mut clause, mut params) = build_note_where(&namespace, kind.as_deref());
params.push(Box::new(key));
clause.push_str(&format!(" AND key = ?{}", params.len()));
let sql = format!("SELECT {NOTE_COLUMNS} FROM notes{clause} ORDER BY kind ASC");
let mut stmt = conn.prepare(&sql)?;
let params: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|p| p.as_ref()).collect();
let rows = stmt.query_map(params.as_slice(), read_note)?.collect();
rows
})
.await
}
async fn query_keyed_notes(
&self,
namespace: &str,
filter: &NoteFilter,
prefix: &str,
after: Option<&NoteKeyCursor>,
page: PageRequest,
) -> StorageResult<(Vec<Note>, Option<NoteKeyCursor>)> {
if !filter.namespaces.is_empty()
|| filter.order_by.is_some()
|| filter.after.is_some()
|| (after.is_some() && page.offset != 0)
|| page.limit == 0
{
return Err(StorageError::InvalidInput { capability: StorageCapability::Notes,
operation: "query_keyed_notes".into(), message: "keyed paging requires primary namespace, keyed order and a positive limit; cursor excludes offset".into() });
}
for property in &filter.property_filters {
validate_json_path(&property.json_path)?;
}
let offset = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_keyed_notes".into(),
message: "offset exceeds the supported integer range".into(),
})?;
let namespace = namespace.to_owned();
let filter = filter.clone();
let prefix = prefix.to_owned();
let after = after.cloned();
self.with_reader("query_keyed_notes", move |conn| {
let (mut clause, mut params) = build_note_filter_where(&namespace, &filter)?;
params.push(Box::new(prefix.clone()));
clause.push_str(&format!(" AND key IS NOT NULL AND key >= ?{}", params.len()));
if let Some(upper) = note_key_prefix_successor(&prefix) {
params.push(Box::new(upper));
clause.push_str(&format!(" AND key < ?{}", params.len()));
}
if let Some(after) = after {
params.push(Box::new(after.updated_at)); let u = params.len();
params.push(Box::new(after.key)); let k = params.len();
params.push(Box::new(after.id.to_string())); let id = params.len();
clause.push_str(&format!(" AND (updated_at < ?{u} OR (updated_at = ?{u} AND key < ?{k}) \
OR (updated_at = ?{u} AND key = ?{k} AND id > ?{id}))"));
}
params.push(Box::new(i64::from(page.limit) + 1)); let limit = params.len();
params.push(Box::new(offset)); let offset = params.len();
let sql = format!("SELECT {NOTE_COLUMNS} FROM notes{clause} ORDER BY updated_at DESC, key DESC, id ASC LIMIT ?{limit} OFFSET ?{offset}");
let mut stmt = conn.prepare(&sql)?;
let params: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
let mut notes = stmt.query_map(params.as_slice(), read_note)?.collect::<Result<Vec<_>, _>>()?;
let has_more = notes.len() > page.limit as usize;
notes.truncate(page.limit as usize);
let next = if has_more { notes.last().map(NoteKeyCursor::from) } else { None };
Ok((notes, next))
}).await
}
async fn upsert_note(&self, note: Note) -> Result<(), StorageError> {
let id_str = note.id.to_string();
let statement = note_upsert_statement(¬e);
self.with_writer_tx("upsert_note", move |conn| {
let mut stmt = conn.prepare_cached(&statement.sql)?;
bind_params(&mut stmt, &statement.params)?;
stmt.raw_execute()?;
assign_note_seq(conn, &id_str)?;
Ok(())
})
.await
}
async fn insert_note_if_absent(&self, note: Note) -> Result<bool, StorageError> {
let id_str = note.id.to_string();
let statement = note_insert_if_absent_statement(¬e);
self.with_writer_tx("insert_note_if_absent", move |conn| {
let mut stmt = conn.prepare_cached(&statement.sql)?;
bind_params(&mut stmt, &statement.params)?;
let inserted = stmt.raw_execute()? > 0;
if inserted {
assign_note_seq(conn, &id_str)?;
}
Ok(inserted)
})
.await
}
async fn replace_note_if_unchanged(
&self,
note: Note,
expected_updated_at: i64,
expected_deleted_at: Option<i64>,
) -> Result<bool, StorageError> {
let statement =
note_replace_if_unchanged_statement(¬e, expected_updated_at, expected_deleted_at);
self.with_writer("replace_note_if_unchanged", move |conn| {
let mut stmt = conn.prepare(&statement.sql)?;
bind_params(&mut stmt, &statement.params)?;
Ok(stmt.raw_execute()? > 0)
})
.await
}
async fn update_note_properties(
&self,
id: Uuid,
properties: Option<serde_json::Value>,
updated_at: i64,
) -> Result<bool, StorageError> {
let statement = note_update_properties_statement(id, &properties, updated_at);
self.with_writer("update_note_properties", move |conn| {
let mut stmt = conn.prepare(&statement.sql)?;
bind_params(&mut stmt, &statement.params)?;
Ok(stmt.raw_execute()? > 0)
})
.await
}
async fn set_note_property(
&self,
id: Uuid,
key: &str,
value: serde_json::Value,
updated_at: i64,
) -> Result<bool, StorageError> {
let statement = note_set_property_statement(id, key, &value, updated_at)?;
self.with_writer("set_note_property", move |conn| {
let mut stmt = conn.prepare(&statement.sql)?;
bind_params(&mut stmt, &statement.params)?;
Ok(stmt.raw_execute()? > 0)
})
.await
}
async fn try_patch_note_property(
&self,
id: Uuid,
namespace: &str,
filter: &NoteFilter,
json_path: &str,
value: serde_json::Value,
updated_at: i64,
) -> Result<bool, StorageError> {
let namespace = namespace.to_string();
let filter = filter.clone();
let value_json = serde_json::to_string(&value).map_err(|e| {
StorageError::driver(StorageCapability::Notes, "try_patch_note_property", e)
})?;
let json_path = json_path.to_string();
self.with_writer("try_patch_note_property", move |conn| {
execute_filtered_note_property_patch(
conn,
id,
&namespace,
&filter,
&json_path,
&value_json,
updated_at,
)
.map(|rows| rows > 0)
})
.await
}
async fn patch_note_property_atomic(
&self,
mut ids: Vec<Uuid>,
namespace: &str,
filter: &NoteFilter,
json_path: &str,
value: serde_json::Value,
updated_at: i64,
) -> Result<(), StorageError> {
let mut seen = HashSet::with_capacity(ids.len());
ids.retain(|id| seen.insert(*id));
if ids.is_empty() {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "patch_note_property_atomic".into(),
message: "at least one note id is required".to_string(),
});
}
let namespace = namespace.to_string();
let filter = filter.clone();
let value_json = serde_json::to_string(&value).map_err(|e| {
StorageError::driver(StorageCapability::Notes, "patch_note_property_atomic", e)
})?;
let json_path = json_path.to_string();
self.with_writer_tx_storage("patch_note_property_atomic", move |conn| {
for id in ids {
let rows = execute_filtered_note_property_patch(
conn,
id,
&namespace,
&filter,
&json_path,
&value_json,
updated_at,
)
.map_err(|error| map_err(error, "patch_note_property_atomic"))?;
if rows != 1 {
return Err(StorageError::Conflict {
capability: StorageCapability::Notes,
operation: "patch_note_property_atomic".into(),
message: format!(
"precondition failed for note {id}: guarded update changed {rows} rows; expected 1"
),
});
}
}
Ok(())
})
.await
}
async fn try_insert_note(&self, note: Note) -> Result<bool, StorageError> {
let namespace = note.namespace.clone();
let id_str = note.id.to_string();
let kind_str = note.kind.to_string();
let status_str = note.status.clone();
let properties_str = note
.properties
.as_ref()
.map(|v| serde_json::to_string(v).unwrap_or_default());
let ext_id_opt: Option<String> = note
.properties
.as_ref()
.and_then(|v| v.get("external_id"))
.and_then(|v| v.as_str())
.filter(|s| !s.is_empty())
.map(|s| s.to_string());
self.with_writer_tx("try_insert_note", move |conn| {
let rows = conn.execute(
"INSERT OR IGNORE INTO notes \
(id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
properties, created_at, updated_at, deleted_at, key) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)",
rusqlite::params![
id_str,
namespace,
kind_str,
status_str,
note.name,
note.content,
note.salience,
note.decay_factor,
note.expires_at,
properties_str,
note.created_at,
note.updated_at,
note.deleted_at,
note.key,
],
)?;
if rows > 0 {
assign_note_seq(conn, &id_str)?;
return Ok(true);
}
if let Some(ref ext_id) = ext_id_opt {
let is_dedup: bool = conn.query_row(
"SELECT COUNT(*) > 0 FROM notes \
WHERE namespace = ?1 \
AND kind = ?2 \
AND json_extract(properties, '$.external_id') = ?3 \
AND deleted_at IS NULL",
rusqlite::params![namespace, kind_str, ext_id],
|row| row.get(0),
)?;
if is_dedup {
return Ok(false);
}
}
Err(rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
Some(
"try_insert_note: INSERT ignored for a constraint other than \
external_id dedup; not masking as deduplication"
.to_string(),
),
))
})
.await
}
async fn upsert_notes(&self, notes: Vec<Note>) -> Result<BatchWriteSummary, StorageError> {
let attempted = notes.len() as u64;
let origin = self.pool.origin();
self.with_writer_tx("upsert_notes", move |conn| {
let _tx_handle = khive_storage::tx_registry::register_scoped(
Some("note_upsert_batch".to_string()),
origin,
);
batch_upsert_notes(conn, ¬es, attempted)
})
.await
}
async fn get_note(&self, id: Uuid) -> Result<Option<Note>, StorageError> {
let id_str = id.to_string();
self.with_reader("get_note", move |conn| {
let mut stmt = conn.prepare(
"SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
properties, created_at, updated_at, deleted_at, key, version \
FROM notes WHERE id = ?1 AND deleted_at IS NULL",
)?;
let mut rows = stmt.query(rusqlite::params![id_str])?;
match rows.next()? {
Some(row) => Ok(Some(read_note(row)?)),
None => Ok(None),
}
})
.await
}
async fn get_note_including_deleted(&self, id: Uuid) -> Result<Option<Note>, StorageError> {
let id_str = id.to_string();
self.with_reader("get_note_including_deleted", move |conn| {
let mut stmt = conn.prepare(
"SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
properties, created_at, updated_at, deleted_at, key, version \
FROM notes WHERE id = ?1",
)?;
let mut rows = stmt.query(rusqlite::params![id_str])?;
match rows.next()? {
Some(row) => Ok(Some(read_note(row)?)),
None => Ok(None),
}
})
.await
}
async fn note_sequence(&self, id: Uuid) -> Result<Option<i64>, StorageError> {
let id = id.to_string();
self.with_reader("note_sequence", move |conn| {
conn.query_row(
"SELECT seq FROM notes_seq WHERE note_id = ?1",
rusqlite::params![id],
|row| row.get(0),
)
.optional()
})
.await
}
async fn get_notes_batch(&self, ids: &[Uuid]) -> Result<Vec<Note>, StorageError> {
if ids.is_empty() {
return Ok(vec![]);
}
const CHUNK: usize = 900;
let id_strings: Vec<String> = ids.iter().map(|id| id.to_string()).collect();
let mut result = Vec::with_capacity(ids.len());
for chunk in id_strings.chunks(CHUNK) {
let chunk_owned = chunk.to_vec();
let notes = self
.with_reader("get_notes_batch", move |conn| {
let placeholders: String = (1..=chunk_owned.len())
.map(|i| format!("?{i}"))
.collect::<Vec<_>>()
.join(", ");
let sql = format!(
"SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
properties, created_at, updated_at, deleted_at, key, version \
FROM notes WHERE id IN ({placeholders}) AND deleted_at IS NULL"
);
let mut stmt = conn.prepare(&sql)?;
let params: Vec<&dyn rusqlite::types::ToSql> = chunk_owned
.iter()
.map(|s| s as &dyn rusqlite::types::ToSql)
.collect();
let rows = stmt.query_map(params.as_slice(), read_note)?;
let mut notes = Vec::new();
for row in rows {
notes.push(row?);
}
Ok(notes)
})
.await?;
result.extend(notes);
}
Ok(result)
}
async fn delete_note(&self, id: Uuid, mode: DeleteMode) -> Result<bool, StorageError> {
match mode {
DeleteMode::Soft => {
let now = chrono::Utc::now().timestamp_micros();
let statement = note_soft_delete_statement(id, now);
self.with_writer("delete_note_soft", move |conn| {
let mut stmt = conn.prepare(&statement.sql)?;
bind_params(&mut stmt, &statement.params)?;
Ok(stmt.raw_execute()? > 0)
})
.await
}
DeleteMode::Hard => {
let note_statement = note_hard_delete_statement(id);
let attachment_statement =
delete_record_attachments_statement(id, AttachmentSubstrate::Note);
self.with_writer_tx("delete_note_hard", move |conn| {
let mut note_stmt = conn.prepare(¬e_statement.sql)?;
bind_params(&mut note_stmt, ¬e_statement.params)?;
let deleted = note_stmt.raw_execute()? > 0;
drop(note_stmt);
if deleted {
let mut attachment_stmt = conn.prepare(&attachment_statement.sql)?;
bind_params(&mut attachment_stmt, &attachment_statement.params)?;
attachment_stmt.raw_execute()?;
}
Ok(deleted)
})
.await
}
}
}
async fn query_notes(
&self,
namespace: &str,
kind: Option<&str>,
page: PageRequest,
) -> Result<Page<Note>, StorageError> {
let namespace = namespace.to_string();
let kind = kind.map(|k| k.to_string());
let limit_i64 = i64::from(page.limit);
let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes".into(),
message: format!(
"PageRequest: offset must be <= i64::MAX, got {}",
page.offset
),
})?;
self.with_reader("query_notes", move |conn| {
let (count_sql, count_params) = build_note_where(&namespace, kind.as_deref());
let count_sql = format!("SELECT COUNT(*) FROM notes{count_sql}");
let (where_sql, mut data_params) = build_note_where(&namespace, kind.as_deref());
data_params.push(Box::new(limit_i64));
data_params.push(Box::new(offset_i64));
let limit_idx = data_params.len() - 1;
let offset_idx = data_params.len();
let data_sql = format!(
"SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
properties, created_at, updated_at, deleted_at, key, version \
FROM notes{} ORDER BY created_at DESC, id ASC LIMIT ?{} OFFSET ?{}",
where_sql, limit_idx, offset_idx,
);
query_note_page_snapshot(
conn,
"query_notes",
&namespace,
&count_sql,
&count_params,
&data_sql,
&data_params,
)
})
.await
}
async fn query_notes_count_free(
&self,
namespace: &str,
kind: Option<&str>,
page: PageRequest,
) -> Result<Page<Note>, StorageError> {
let namespace = namespace.to_string();
let kind = kind.map(str::to_string);
let limit_i64 = i64::from(page.limit);
let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_count_free".into(),
message: format!(
"PageRequest: offset must be <= i64::MAX, got {}",
page.offset
),
})?;
self.with_reader("query_notes_count_free", move |conn| {
let (where_sql, mut params) = build_note_where(&namespace, kind.as_deref());
params.push(Box::new(limit_i64));
params.push(Box::new(offset_i64));
let limit_idx = params.len() - 1;
let offset_idx = params.len();
let sql = format!(
"SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
expires_at, properties, created_at, updated_at, deleted_at, key, version \
FROM notes{where_sql} ORDER BY created_at DESC, id ASC \
LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
);
let mut stmt = conn.prepare(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|param| param.as_ref()).collect();
let mut rows = stmt.query(param_refs.as_slice())?;
let mut items = Vec::new();
while let Some(row) = rows.next()? {
items.push(read_note(row)?);
}
Ok(Page { items, total: None })
})
.await
}
async fn query_notes_filtered(
&self,
namespace: &str,
filter: &NoteFilter,
page: PageRequest,
) -> Result<Page<Note>, StorageError> {
if filter.unordered {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered".into(),
message: "NoteFilter.unordered is supported only by \
query_notes_filtered_count_free"
.into(),
});
}
for pf in &filter.property_filters {
validate_json_path(&pf.json_path)?;
}
if let Some((path, _)) = &filter.order_by {
validate_json_path(path)?;
}
if filter.order_by_instant && !matches!(filter.order_by.as_ref(), Some((_, SortDir::Asc))) {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered".into(),
message: "order_by_instant requires an ascending property order".into(),
});
}
if filter.after.is_some() || filter.after_instant.is_some() {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered".into(),
message: "NoteFilter.after or after_instant (keyset pagination) is not supported by this \
method: it computes an exact COUNT(*) total over the whole \
matching set, which has no defined meaning paired with a seek \
boundary; use query_notes_filtered_count_free instead, which \
seeks and returns total: None"
.into(),
});
}
let namespace = namespace.to_string();
let filter = filter.clone();
let limit_i64 = i64::from(page.limit);
let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered".into(),
message: format!(
"PageRequest: offset must be <= i64::MAX, got {}",
page.offset
),
})?;
self.with_reader("query_notes_filtered", move |conn| {
let (count_sql, count_params) = build_note_filter_read_clause(&namespace, &filter)?;
let count_sql = format!("SELECT COUNT(*) FROM notes{count_sql}");
let (where_sql, mut data_params) = build_note_filter_read_clause(&namespace, &filter)?;
data_params.push(Box::new(limit_i64));
data_params.push(Box::new(offset_i64));
let order_clause = note_filter_page_order_clause(&filter);
let limit_idx = data_params.len() - 1;
let offset_idx = data_params.len();
let data_sql = format!(
"SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
expires_at, properties, created_at, updated_at, deleted_at, key, version \
FROM notes{}{order_clause} LIMIT ?{} OFFSET ?{}",
where_sql, limit_idx, offset_idx,
);
query_note_page_snapshot(
conn,
"query_notes_filtered",
&namespace,
&count_sql,
&count_params,
&data_sql,
&data_params,
)
})
.await
}
async fn query_notes_filtered_count_free(
&self,
namespace: &str,
filter: &NoteFilter,
page: PageRequest,
) -> Result<Page<Note>, StorageError> {
for property_filter in &filter.property_filters {
validate_json_path(&property_filter.json_path)?;
}
if let Some((path, _)) = &filter.order_by {
validate_json_path(path)?;
}
if filter.order_by_instant && !matches!(filter.order_by.as_ref(), Some((_, SortDir::Asc))) {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_count_free".into(),
message: "order_by_instant requires an ascending property order".into(),
});
}
if filter.order_by_instant && filter.unordered {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_count_free".into(),
message: "order_by_instant is incompatible with unordered pages".into(),
});
}
if filter.after_instant.is_some() && !filter.order_by_instant {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_count_free".into(),
message: "after_instant requires order_by_instant".into(),
});
}
if filter.after.is_some() && filter.after_instant.is_some() {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_count_free".into(),
message: "after and after_instant are mutually exclusive".into(),
});
}
if filter.after.is_some() && filter.order_by.is_some() {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_count_free".into(),
message: "NoteFilter.after is incompatible with a custom order_by; it is \
defined only over the default created_at DESC, id ASC order"
.into(),
});
}
if (filter.after.is_some() || filter.after_instant.is_some()) && page.offset != 0 {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_count_free".into(),
message: "NoteFilter.after or after_instant and a non-zero PageRequest.offset \
are mutually exclusive pagination strategies; pass offset: 0"
.into(),
});
}
let namespace = namespace.to_string();
let filter = filter.clone();
let limit_i64 = i64::from(page.limit);
let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_count_free".into(),
message: format!(
"PageRequest: offset must be <= i64::MAX, got {}",
page.offset
),
})?;
self.with_reader("query_notes_filtered_count_free", move |conn| {
if let Some(after) = &filter.after_instant {
let mut base_filter = filter.clone();
base_filter.after_instant = None;
let items =
fetch_notes_after_instant(conn, &namespace, &base_filter, after, limit_i64)?;
return Ok(Page { items, total: None });
}
if let Some(after) = &filter.after {
let mut base_filter = filter.clone();
base_filter.after = None;
let items = fetch_notes_after(conn, &namespace, &base_filter, after, limit_i64)?;
return Ok(Page { items, total: None });
}
let (where_sql, mut params) = build_note_filter_read_clause(&namespace, &filter)?;
params.push(Box::new(limit_i64));
params.push(Box::new(offset_i64));
let limit_idx = params.len() - 1;
let offset_idx = params.len();
let order_clause = note_filter_page_order_clause(&filter);
let sql = format!(
"SELECT {NOTE_COLUMNS} FROM notes{where_sql}{order_clause} \
LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
);
let mut stmt = conn.prepare(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|param| param.as_ref()).collect();
let mut rows = stmt.query(param_refs.as_slice())?;
let mut items = Vec::new();
while let Some(row) = rows.next()? {
items.push(read_note(row)?);
#[cfg(test)]
if items.len() == 1 {
tests::page_snapshot_seam::hook("query_notes_filtered_count_free", &namespace);
}
}
Ok(Page { items, total: None })
})
.await
}
async fn count_notes_filtered_in_snapshot(
&self,
namespace: &str,
filters: &[NoteFilter],
) -> Result<Vec<u64>, StorageError> {
for filter in filters {
for pf in &filter.property_filters {
validate_json_path(&pf.json_path)?;
}
}
let namespace = namespace.to_string();
let filters = filters.to_vec();
self.with_reader("count_notes_filtered_in_snapshot", move |conn| {
let tx = rusqlite::Transaction::new_unchecked(
conn,
rusqlite::TransactionBehavior::Deferred,
)?;
let mut counts = Vec::with_capacity(filters.len());
for filter in filters.iter() {
#[cfg(test)]
if !counts.is_empty() {
tests::page_snapshot_seam::hook("count_notes_filtered_in_snapshot", &namespace);
}
let (where_sql, params) = build_note_filter_read_clause(&namespace, filter)?;
let sql = format!("SELECT COUNT(*) FROM notes{where_sql}");
let mut stmt = tx.prepare(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|param| param.as_ref()).collect();
let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
counts.push(count as u64);
}
tx.commit()?;
Ok(counts)
})
.await
}
async fn count_notes_filtered_bounded_in_snapshot(
&self,
namespace: &str,
filters: &[NoteFilter],
cap: u32,
) -> Result<Vec<BoundedCount>, StorageError> {
for filter in filters {
for property_filter in &filter.property_filters {
validate_json_path(&property_filter.json_path)?;
}
}
let namespace = namespace.to_string();
let filters = filters.to_vec();
let cap_u64 = u64::from(cap);
let probe_limit_i64 = i64::from(cap) + 1;
self.with_reader("count_notes_filtered_bounded_in_snapshot", move |conn| {
let tx = rusqlite::Transaction::new_unchecked(
conn,
rusqlite::TransactionBehavior::Deferred,
)?;
let mut counts = Vec::with_capacity(filters.len());
for filter in &filters {
#[cfg(test)]
if !counts.is_empty() {
tests::page_snapshot_seam::hook(
"count_notes_filtered_bounded_in_snapshot",
&namespace,
);
}
let (where_sql, mut params) = build_note_filter_read_clause(&namespace, filter)?;
params.push(Box::new(probe_limit_i64));
let limit_idx = params.len();
let sql = format!(
"SELECT COUNT(*) FROM (SELECT 1 FROM notes{where_sql} LIMIT ?{limit_idx})"
);
let mut stmt = tx.prepare(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|param| param.as_ref()).collect();
let observed: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
let observed = observed as u64;
counts.push(BoundedCount {
count: observed.min(cap_u64),
cap: cap_u64,
saturated: observed > cap_u64,
});
}
tx.commit()?;
Ok(counts)
})
.await
}
async fn query_notes_filtered_after(
&self,
namespace: &str,
filter: &NoteFilter,
after: Option<SeekCursor>,
limit: u32,
) -> Result<SeekPage<Note>, StorageError> {
if limit == 0 {
return Ok(SeekPage::default());
}
if filter.order_by.is_some() {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_after".into(),
message: "custom order_by is not compatible with insertion-sequence pagination"
.into(),
});
}
for property_filter in &filter.property_filters {
validate_json_path(&property_filter.json_path)?;
}
let namespace = namespace.to_string();
let filter = filter.clone();
let limit_usize = limit as usize;
let probe_limit_i64 = i64::from(limit) + 1;
self.with_reader("query_notes_filtered_after", move |conn| {
let (mut where_sql, mut params) = build_note_filter_where(&namespace, &filter)?;
if let Some(cursor) = after {
params.push(Box::new(cursor.sequence));
where_sql.push_str(&format!(" AND notes_seq.seq > ?{}", params.len()));
}
params.push(Box::new(probe_limit_i64));
let limit_idx = params.len();
let sql = format!(
"SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
expires_at, properties, created_at, updated_at, deleted_at, key, version, notes_seq.seq \
FROM notes_seq CROSS JOIN notes ON notes.id = notes_seq.note_id{where_sql} \
ORDER BY notes_seq.seq ASC LIMIT ?{limit_idx}"
);
let mut stmt = conn.prepare(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|param| param.as_ref()).collect();
let rows = stmt.query_map(param_refs.as_slice(), |row| {
Ok((read_note(row)?, row.get::<_, i64>(15)?))
})?;
let mut entries = rows.collect::<Result<Vec<_>, _>>()?;
let has_more = entries.len() > limit_usize;
if has_more {
entries.truncate(limit_usize);
}
let next_after = if has_more {
entries.last().map(|(note, sequence)| SeekCursor {
sequence: *sequence,
id: note.id,
})
} else {
None
};
let items = entries.into_iter().map(|(note, _)| note).collect();
Ok(SeekPage { items, next_after })
})
.await
}
async fn query_notes_filtered_bounded(
&self,
namespace: &str,
filter: &NoteFilter,
max_rows: u32,
) -> Result<Vec<Note>, StorageError> {
if filter.unordered {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "query_notes_filtered_bounded".into(),
message: "NoteFilter.unordered is supported only by \
query_notes_filtered_count_free"
.into(),
});
}
for pf in &filter.property_filters {
validate_json_path(&pf.json_path)?;
}
if let Some((path, _)) = &filter.order_by {
validate_json_path(path)?;
}
let namespace = namespace.to_string();
let filter = filter.clone();
let limit_i64 = i64::from(max_rows) + 1;
self.with_reader("query_notes_filtered_bounded", move |conn| {
let (where_sql, mut data_params) = build_note_filter_read_clause(&namespace, &filter)?;
data_params.push(Box::new(limit_i64));
let limit_idx = data_params.len();
let order_clause = match &filter.order_by {
Some((path, dir)) => {
let dir_str = match dir {
SortDir::Asc => "ASC",
SortDir::Desc => "DESC",
};
format!(" ORDER BY {} {dir_str}, id ASC", json_extract_expr(path))
}
None => " ORDER BY created_at DESC, id ASC".to_string(),
};
let data_sql = format!(
"SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
expires_at, properties, created_at, updated_at, deleted_at, key, version \
FROM notes{where_sql}{order_clause} LIMIT ?{limit_idx}",
);
let mut stmt = conn.prepare(&data_sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
data_params.iter().map(|p| p.as_ref()).collect();
let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
let mut items = Vec::new();
for row in rows {
items.push(row?);
}
Ok(items)
})
.await
}
async fn count_notes(&self, namespace: &str, kind: Option<&str>) -> Result<u64, StorageError> {
let namespace = namespace.to_string();
let kind = kind.map(|k| k.to_string());
self.with_reader("count_notes", move |conn| {
let (where_sql, params) = build_note_where(&namespace, kind.as_deref());
let sql = format!("SELECT COUNT(*) FROM notes{}", where_sql);
let mut stmt = conn.prepare(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|p| p.as_ref()).collect();
let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
Ok(count as u64)
})
.await
}
async fn count_notes_in_namespaces(
&self,
namespaces: &[String],
kind: Option<&str>,
) -> Result<u64, StorageError> {
let namespaces: Vec<String> = namespaces
.iter()
.cloned()
.collect::<HashSet<_>>()
.into_iter()
.collect();
let kind = kind.map(str::to_string);
self.with_reader("count_notes_in_namespaces", move |conn| {
let mut total = 0;
for chunk in namespaces.chunks(NAMESPACE_COUNT_CHUNK_SIZE) {
let (where_sql, params) = build_note_where_for_namespaces(chunk, kind.as_deref());
let sql = format!("SELECT COUNT(*) FROM notes{where_sql}");
let mut stmt = conn.prepare(&sql)?;
let param_refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|p| p.as_ref()).collect();
let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
total += count as u64;
}
Ok(total)
})
.await
}
}
const NOTES_DDL: &str = include_str!("../../sql/notes-ddl.sql");
const NOTES_SEQ_REPAIR_DDL: &str = include_str!("../../sql/008-notes-seq-repair.sql");
pub(crate) fn ensure_notes_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
conn.execute_batch(NOTES_DDL)
}
pub(crate) fn repair_notes_seq(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
conn.execute_batch(NOTES_SEQ_REPAIR_DDL)
}
#[cfg(test)]
#[path = "note_tests.rs"]
mod tests;
#[cfg(test)]
#[path = "comm_filter_plan_tests.rs"]
mod comm_filter_plan_tests;
#[cfg(test)]
#[path = "note_list_plan_tests.rs"]
mod note_list_plan_tests;