mod budget;
mod dispute;
mod receipt;
mod render;
pub(crate) use render::claim_value;
use crate::connection::ConnectionRegistry;
use crate::contracts::{
PromptTerms, RecallBounds, RecallMode, RecallOutcomeKind, RecallReceipt, RecallRequest,
RetrievalPolicy, SchemaAvailability, recall, terms,
};
use saya_agent::ContextBlock;
use saya_store::{SchemaStore, SqliteStateStore};
use saya_types::{DatabaseObjectRef, ProfileIdentity};
pub(crate) const BLOCK_LABEL: &str = "database-contracts";
pub(crate) async fn recall_context_blocks(
prompt: &str,
system_prompt: Option<&str>,
allow_database_context: bool,
recall_mode: RecallMode,
bounds: RecallBounds,
registry: &ConnectionRegistry,
state_db: Option<&SqliteStateStore>,
) -> (Vec<ContextBlock>, RecallReceipt) {
if !allow_database_context {
return (Vec::new(), RecallReceipt::privacy_gate_closed());
}
let Some(store) = state_db else {
return (Vec::new(), RecallReceipt::ran_empty(false));
};
let Some((identities, schemas)) = resolve_profiles(registry, store).await else {
return (Vec::new(), RecallReceipt::ran_empty(false));
};
if identities.is_empty() || prompt.trim().is_empty() {
return (Vec::new(), RecallReceipt::ran_empty(false));
}
let PromptTerms { explicit, terms } = terms::extract(prompt);
let explicit_refs = build_refs(&identities, &explicit);
if explicit_refs.is_empty() && terms.is_empty() {
return (Vec::new(), RecallReceipt::ran_empty(false));
}
let request = RecallRequest {
profiles: &identities,
explicit_refs: &explicit_refs,
terms: &terms,
allow_database_context: true,
schemas: &schemas,
now_unix_ms: crate::contracts::now_unix_ms(),
bounds,
recall_mode,
admit_candidate: None,
policy: RetrievalPolicy::ForModel,
};
let outcome = recall(store, request).await;
if outcome.diagnostics.store_unavailable {
return (Vec::new(), RecallReceipt::ran_empty(true));
}
if outcome.contracts.is_empty() {
return (Vec::new(), RecallReceipt::ran_empty(false));
}
let name_of = render::name_by_identity(registry);
let count_truncated = outcome.contracts.iter().any(|c| c.truncated);
let (body, byte_truncated, kept) = budget::bound_body(
&outcome.contracts,
&name_of,
system_prompt,
prompt,
bounds.max_bytes,
);
let byte_dropped_claims: usize = outcome.contracts[kept..]
.iter()
.map(|c| c.claims.len())
.sum();
let dropped_by_bounds = outcome.diagnostics.excluded_by_count_bounds + byte_dropped_claims;
let receipt = RecallReceipt {
kind: RecallOutcomeKind::Ran {
store_unavailable: false,
},
supplied: receipt::supplied_contracts(&outcome.contracts[..kept], &name_of),
dropped_by_bounds,
};
if body.is_empty() {
return (
vec![ContextBlock {
label: BLOCK_LABEL.to_string(),
body,
truncated: true,
}],
receipt,
);
}
(
vec![ContextBlock {
label: BLOCK_LABEL.to_string(),
body,
truncated: count_truncated || byte_truncated,
}],
receipt,
)
}
async fn resolve_profiles(
registry: &ConnectionRegistry,
store: &SqliteStateStore,
) -> Option<(
Vec<ProfileIdentity>,
Vec<(ProfileIdentity, SchemaAvailability)>,
)> {
let mut identities = Vec::new();
let mut schemas = Vec::new();
for (_name, entry) in registry.entries() {
let Some(id_str) = entry.profile_id.as_deref() else {
continue;
};
let Ok(identity) = ProfileIdentity::parse(id_str) else {
continue;
};
let availability = match store.get_schema(identity.as_str()).await {
Ok(Some(cached)) => {
SchemaAvailability::available(cached.schema, cached.updated_unix_ms)
}
Ok(None) => SchemaAvailability::Missing,
Err(_) => SchemaAvailability::Unavailable,
};
identities.push(identity.clone());
schemas.push((identity, availability));
}
if identities.is_empty() {
None
} else {
Some((identities, schemas))
}
}
fn build_refs(
identities: &[ProfileIdentity],
explicit: &[crate::contracts::args::QualifiedName],
) -> Vec<DatabaseObjectRef> {
use saya_types::DatabaseObjectKind;
let mut out = Vec::new();
for id in identities {
for q in explicit {
if let Ok(obj) = DatabaseObjectRef::new(
id.clone(),
&q.catalog,
&q.schema,
&q.object,
DatabaseObjectKind::Table,
) {
out.push(obj);
}
}
}
out
}
#[cfg(test)]
#[path = "../recall_context_tests.rs"]
mod tests;