mod decision_devnote;
mod file;
mod file_link;
mod gotcha;
mod reads;
mod record_import;
#[cfg(test)]
mod link_sync_tests;
pub(crate) use decision_devnote::{handle_decision_upsert, handle_dev_note_upsert};
pub(crate) use file::{handle_doc_capture, handle_file_enrich, handle_file_reparse};
pub(crate) use gotcha::{handle_gotcha_confirm, handle_gotcha_tombstone, handle_gotcha_upsert};
pub(crate) use reads::{handle_mem_bootstrap, handle_mem_get, handle_mem_query};
pub(crate) use record_import::handle_record_import;
use file_link::{apply_confirmation_propagation, compute_file_link_updates};
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use uuid::Uuid;
use crate::health::quality;
use crate::mcp::protocol::{self, AuditEntry, ErrorCode};
use crate::store::db::KnowledgeWriteOp;
use crate::store::record::{
Category, ConfidenceScore, FileRecord, GotchaRecord, Priority as StorePriority, QualityScore,
Record, RecordLifecycle, RecordSource, RecordVersion, StalenessScore, TombstoneReason,
};
use crate::store::Store;
use super::dispatch_v2::RequestContext;
fn now_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
fn audit_nanos_key(prefix: &str) -> String {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_nanos();
format!("{prefix}{nanos}")
}
const AUDIT_KNOWLEDGE_PREFIX: &str = "audit:knowledge:";
pub(crate) const AUDIT_SESSION_PREFIX: &str = "audit:session:";
fn map_priority(p: &protocol::Priority) -> StorePriority {
match p {
protocol::Priority::Critical => StorePriority::Critical,
protocol::Priority::High => StorePriority::High,
protocol::Priority::Normal => StorePriority::Normal,
protocol::Priority::Low => StorePriority::Low,
}
}
fn map_severity(s: &protocol::Severity) -> StorePriority {
match s {
protocol::Severity::Critical => StorePriority::Critical,
protocol::Severity::High => StorePriority::High,
protocol::Severity::Normal => StorePriority::Normal,
protocol::Severity::Low => StorePriority::Low,
}
}
pub(crate) fn make_audit_with_prefix(
ctx: &RequestContext,
request_id: Uuid,
command_kind: &str,
target_key: &str,
accepted: bool,
error_code: Option<ErrorCode>,
prefix: &str,
) -> Option<(String, Vec<u8>)> {
let entry = AuditEntry {
ts: now_secs(),
peer_uid: ctx.peer.uid,
peer_pid: ctx.peer.pid,
daemon_session: ctx.daemon_session,
request_id,
command_kind: command_kind.to_string(),
target_key: target_key.to_string(),
accepted,
error_code,
};
match rmp_serde::to_vec_named(&entry) {
Ok(bytes) => Some((audit_nanos_key(prefix), bytes)),
Err(e) => {
tracing::error!("audit serialization failed — this is a bug, audit entry skipped: {e}");
None
}
}
}
pub(crate) fn make_audit(
ctx: &RequestContext,
request_id: Uuid,
command_kind: &str,
target_key: &str,
accepted: bool,
error_code: Option<ErrorCode>,
) -> Option<(String, Vec<u8>)> {
make_audit_with_prefix(
ctx,
request_id,
command_kind,
target_key,
accepted,
error_code,
AUDIT_KNOWLEDGE_PREFIX,
)
}
pub(crate) fn make_session_audit(
ctx: &RequestContext,
request_id: Uuid,
command_kind: &str,
target_key: &str,
accepted: bool,
error_code: Option<ErrorCode>,
) -> Option<(String, Vec<u8>)> {
make_audit_with_prefix(
ctx,
request_id,
command_kind,
target_key,
accepted,
error_code,
AUDIT_SESSION_PREFIX,
)
}
type HandlerResult = std::result::Result<serde_json::Value, (ErrorCode, String)>;
const WRITE_CONFLICT_RETRIES: usize = 4;
async fn retry_on_write_conflict<T, F, Fut>(mut op: F) -> Result<T, (ErrorCode, String)>
where
F: FnMut() -> Fut,
Fut: std::future::Future<Output = Result<T, (ErrorCode, String)>>,
{
for attempt in 0..WRITE_CONFLICT_RETRIES {
match op().await {
Ok(value) => return Ok(value),
Err((_, ref msg))
if attempt + 1 < WRITE_CONFLICT_RETRIES
&& msg.to_lowercase().contains("write conflict") =>
{
tokio::time::sleep(std::time::Duration::from_millis(5u64 << attempt)).await;
}
Err(e) => return Err(e),
}
}
unreachable!("retry_on_write_conflict loop always returns within the body")
}