use std::sync::Arc;
use axum::{
Extension, Json, Router,
extract::{Path, Query},
http::StatusCode,
response::{IntoResponse, Response},
routing::{get, post},
};
use serde::{Deserialize, Serialize};
#[cfg(feature = "mining")]
use ijima_core::capabilities::MINING_TRIGGER;
use ijima_core::{
AcceptedExtraction, DiaryEntry, Embedder, EntityId, KnowledgeGraph, Memory, MemoryId,
NamespaceCount, NamespaceId, PalaceGraph, ProjectTaxon, QueuedExtraction, RepoDirectory, Room,
SearchHit, Session, SessionId, SessionTurn, Store, TunnelTraversal,
capabilities::{
ADMIN, KNOWLEDGE_READ, MEMORY_READ, MEMORY_WRITE, MINING_REVIEW, SESSION_INGEST,
TRUST_PROMOTE,
},
harness::Harness,
};
use crate::extractor::AuthPrincipal;
use crate::redaction::Redactor;
pub fn app(
auth: Arc<crate::IjimaAuth>,
store: Arc<dyn Store>,
kg: Arc<dyn KnowledgeGraph>,
embedder: Option<Arc<dyn Embedder>>,
redactor: Arc<Redactor>,
#[cfg(feature = "rate-limit")] rate_limiter: Option<crate::rate_limit::RateLimitState>,
) -> Router {
let router = Router::new()
.route("/health", get(health))
.route("/status", get(status))
.route("/memories", get(browse_memories).post(store_memory))
.route("/memories/check", post(check_duplicate))
.route("/memories/search", post(search_memories))
.route("/memories/stats", get(memory_stats))
.route("/memories/{id}", get(recall_memory).delete(delete_memory))
.route("/memories/{id}/promote", post(promote_memory))
.route("/rooms", get(list_rooms))
.route("/taxonomy", get(taxonomy))
.route("/palace/graph", get(palace_graph))
.route("/palace/tunnel", get(traverse_tunnel))
.route("/diaries", post(write_diary))
.route("/diaries/{agent}", get(read_diary))
.route("/repos", get(list_repos).post(register_repo))
.route("/repos/resolve", get(resolve_repo))
.route("/doctrine", post(ingest_doctrine))
.route("/wakeup", get(wakeup))
.route("/kg/triples", post(add_triple).get(find_triples))
.route("/kg/entities/{id}", get(query_entity))
.route("/kg/triples/{id}/invalidate", post(invalidate_triple))
.route("/kg/timeline", get(kg_timeline))
.route("/kg/stats", get(kg_stats))
.route(
"/sessions/{session_id}/turns",
post(ingest_turn).get(session_turns),
)
.route("/sessions", post(create_session).get(list_sessions))
.route("/sessions/{session_id}/end", post(end_session))
.route("/mining/queue", get(list_pending))
.route("/mining/queue/{id}/accept", post(accept_extraction))
.route("/mining/queue/{id}/reject", post(reject_extraction));
#[cfg(feature = "mining")]
let router = router.route("/sessions/{session_id}/mine", post(trigger_mine));
let router = router
.layer(Extension(auth))
.layer(Extension(store))
.layer(Extension(kg))
.layer(Extension(embedder))
.layer(Extension(redactor));
#[cfg(feature = "rate-limit")]
let router = match rate_limiter {
Some(rl) => router.layer(Extension(rl)),
None => router,
};
#[cfg(not(feature = "rate-limit"))]
let router = router;
router
}
#[derive(Debug)]
pub enum ApiError {
Forbidden,
NotFound,
BadRequest(String),
Conflict(String),
Internal(String),
}
impl IntoResponse for ApiError {
fn into_response(self) -> Response {
let (status, msg): (StatusCode, String) = match self {
ApiError::Forbidden => (StatusCode::FORBIDDEN, "forbidden".into()),
ApiError::NotFound => (StatusCode::NOT_FOUND, "not found".into()),
ApiError::BadRequest(m) => (StatusCode::BAD_REQUEST, m),
ApiError::Conflict(m) => (StatusCode::CONFLICT, m),
ApiError::Internal(m) => (StatusCode::INTERNAL_SERVER_ERROR, m),
};
(status, msg).into_response()
}
}
fn internal(e: ijima_core::IjimaError) -> ApiError {
match e {
ijima_core::IjimaError::Duplicate { detail } => ApiError::Conflict(detail),
other => ApiError::Internal(other.to_string()),
}
}
#[derive(Deserialize, Default)]
struct NsQuery {
namespace: Option<String>,
limit: Option<usize>,
}
fn resolve_ns(
principal: &AuthPrincipal,
requested: Option<&str>,
) -> Result<ijima_core::NamespaceId, ApiError> {
let own = format!("ns_{}_private", principal.0.principal.as_str());
match requested {
None => Ok(ijima_core::NamespaceId::new(own)),
Some(ns) if ns == own => Ok(ijima_core::NamespaceId::new(ns)),
Some(ns) if ns.ends_with("_private") => Err(ApiError::Forbidden),
Some(ns) => Ok(ijima_core::NamespaceId::new(ns)),
}
}
async fn health() -> impl IntoResponse {
Json(serde_json::json!({ "status": "ok" }))
}
#[derive(Serialize)]
struct StatusResponse {
memories: usize,
namespaces: Vec<NamespaceCount>,
entities: usize,
triples: usize,
}
async fn status(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Extension(kg): Extension<Arc<dyn KnowledgeGraph>>,
) -> Result<Json<StatusResponse>, ApiError> {
if !principal.0.may(ijima_core::capabilities::ADMIN) {
return Err(ApiError::Forbidden);
}
let store_stats = store.store_stats().await.map_err(internal)?;
let kg_stats = kg.kg_global_stats().await.map_err(internal)?;
Ok(Json(StatusResponse {
memories: store_stats.total_memories,
namespaces: store_stats.namespaces,
entities: kg_stats.entities,
triples: kg_stats.triples,
}))
}
#[derive(Serialize)]
struct IdResponse {
id: String,
}
async fn store_memory(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Json(memory): Json<Memory>,
) -> Result<Json<IdResponse>, ApiError> {
if !principal.0.may(MEMORY_WRITE) {
return Err(ApiError::Forbidden);
}
let ns = principal.0.personal_namespace();
let mut memory = memory;
if memory.created_at.is_empty() {
memory.created_at = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs().to_string())
.unwrap_or_default();
}
let id = store.store_memory(&ns, memory).await.map_err(internal)?;
Ok(Json(IdResponse { id: id.0 }))
}
#[derive(Deserialize)]
struct CheckDuplicateRequest {
content: String,
}
#[derive(Serialize)]
struct CheckDuplicateResponse {
duplicate: Option<String>,
}
async fn check_duplicate(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<NsQuery>,
Json(req): Json<CheckDuplicateRequest>,
) -> Result<Json<CheckDuplicateResponse>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let dup = store
.check_duplicate(&ns, &req.content)
.await
.map_err(internal)?;
Ok(Json(CheckDuplicateResponse {
duplicate: dup.map(|id| id.0),
}))
}
async fn recall_memory(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(id): Path<String>,
Query(q): Query<NsQuery>,
) -> Result<Json<Memory>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
match store
.recall_memory(&ns, &MemoryId(id))
.await
.map_err(internal)?
{
Some(memory) => Ok(Json(memory)),
None => Err(ApiError::NotFound),
}
}
async fn delete_memory(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(id): Path<String>,
Query(q): Query<NsQuery>,
) -> Result<StatusCode, ApiError> {
if !principal.0.may(MEMORY_WRITE) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
store
.delete_memory(&ns, &MemoryId(id))
.await
.map_err(internal)?;
Ok(StatusCode::NO_CONTENT)
}
#[derive(Deserialize)]
struct SearchRequest {
text: String,
limit: Option<usize>,
scope: Option<String>,
}
#[derive(Serialize)]
struct SearchResponse {
memories: Vec<SearchHit>,
}
async fn search_memories(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Extension(embedder): Extension<Option<Arc<dyn Embedder>>>,
Query(q): Query<NsQuery>,
Json(req): Json<SearchRequest>,
) -> Result<Json<SearchResponse>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let embedder = embedder
.ok_or_else(|| ApiError::Internal("search unavailable: daemon has no embedder".into()))?;
let query = embedder.embed(&req.text).map_err(internal)?;
let limit = req.limit.unwrap_or(10);
let hits = if req.scope.as_deref() == Some("visible") {
let own_ns = principal.0.personal_namespace();
let global_ns = NamespaceId::new("global");
let own_hits = store
.search_memories(&own_ns, &query, limit)
.await
.map_err(internal)?;
let global_hits = if own_ns == global_ns {
Vec::new()
} else {
store
.search_memories(&global_ns, &query, limit)
.await
.map_err(internal)?
};
merge_search_hits(own_hits, global_hits, limit)
} else {
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
store
.search_memories(&ns, &query, limit)
.await
.map_err(internal)?
};
Ok(Json(SearchResponse { memories: hits }))
}
fn merge_search_hits(a: Vec<SearchHit>, b: Vec<SearchHit>, limit: usize) -> Vec<SearchHit> {
use std::collections::HashSet;
let mut all: Vec<SearchHit> = a.into_iter().chain(b).collect();
all.sort_by(|x, y| {
y.similarity
.partial_cmp(&x.similarity)
.unwrap_or(std::cmp::Ordering::Equal)
});
let mut seen: HashSet<String> = HashSet::new();
all.retain(|h| seen.insert(h.memory.id.0.clone()));
all.truncate(limit);
all
}
#[derive(Deserialize)]
struct PromoteRequest {
target_namespace: String,
new_id: Option<String>,
}
#[derive(Serialize)]
struct PromoteResponse {
id: String,
original_id: String,
target_namespace: String,
redactions: Vec<crate::redaction::Redaction>,
}
async fn promote_memory(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Extension(redactor): Extension<Arc<Redactor>>,
Path(id): Path<String>,
Json(req): Json<PromoteRequest>,
) -> Result<Json<PromoteResponse>, ApiError> {
if !principal.0.may(TRUST_PROMOTE) {
return Err(ApiError::Forbidden);
}
let personal_ns = principal.0.personal_namespace();
let memory = store
.recall_memory(&personal_ns, &MemoryId(id.clone()))
.await
.map_err(internal)?
.ok_or(ApiError::NotFound)?;
let scrubbed = redactor.redact(&memory.content);
let new_id = req
.new_id
.clone()
.unwrap_or_else(|| format!("{id}__shared"));
let promoted = Memory {
id: MemoryId(new_id.clone()),
content: scrubbed.text,
project: memory.project,
topic: memory.topic,
source: ijima_core::memory::MemorySource::Explicit,
harness: memory.harness,
session_id: Some(id.clone()),
origin: memory.origin.clone(),
authority: memory.authority.clone(),
importance: memory.importance,
created_at: memory.created_at.clone(),
};
let target_ns = ijima_core::NamespaceId::new(&req.target_namespace);
store
.store_memory(&target_ns, promoted)
.await
.map_err(internal)?;
Ok(Json(PromoteResponse {
id: new_id,
original_id: id,
target_namespace: req.target_namespace,
redactions: scrubbed.redactions,
}))
}
#[derive(Deserialize)]
struct DoctrineRequest {
id: String,
content: String,
project: String,
topic: String,
}
async fn ingest_doctrine(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Json(req): Json<DoctrineRequest>,
) -> Result<Json<IdResponse>, ApiError> {
if !principal.0.may(ijima_core::capabilities::ADMIN) {
return Err(ApiError::Forbidden);
}
let ns = ijima_core::NamespaceId::new(ijima_core::namespace::DOCTRINE_NAMESPACE);
store
.delete_memory(&ns, &MemoryId(req.id.clone()))
.await
.map_err(internal)?;
let memory = Memory {
id: MemoryId(req.id.clone()),
content: req.content,
project: req.project,
topic: req.topic,
source: ijima_core::memory::MemorySource::Doctrine,
harness: ijima_core::harness::Harness::Other,
session_id: None,
origin: ijima_core::InstanceId::local(),
authority: ijima_core::AuthorityScope::local(),
importance: 1.0,
created_at: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs().to_string())
.unwrap_or_default(),
};
store.store_memory(&ns, memory).await.map_err(internal)?;
Ok(Json(IdResponse { id: req.id }))
}
const WAKEUP_PERSONAL_LIMIT: usize = 20;
const WAKEUP_DOCTRINE_LIMIT: usize = 50;
#[derive(Serialize)]
struct WakeupResponse {
identity: serde_json::Value,
personal_essentials: Vec<Memory>,
doctrine: Vec<Memory>,
}
async fn wakeup(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
) -> Result<Json<WakeupResponse>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let personal_ns = principal.0.personal_namespace();
let doctrine_ns = ijima_core::NamespaceId::new(ijima_core::namespace::DOCTRINE_NAMESPACE);
let (personal_essentials, doctrine) = tokio::join!(
store.list_memories(&personal_ns, WAKEUP_PERSONAL_LIMIT),
store.list_memories(&doctrine_ns, WAKEUP_DOCTRINE_LIMIT),
);
Ok(Json(WakeupResponse {
identity: serde_json::json!({ "principal": principal.0.principal.as_str() }),
personal_essentials: personal_essentials.map_err(internal)?,
doctrine: doctrine.map_err(internal)?,
}))
}
#[derive(Deserialize)]
struct AddTripleRequest {
subject: String,
predicate: String,
object: String,
valid_from: Option<String>,
confidence: Option<f32>,
source_memory_id: Option<String>,
}
async fn add_triple(
principal: AuthPrincipal,
Extension(kg): Extension<Arc<dyn KnowledgeGraph>>,
Extension(store): Extension<Arc<dyn Store>>,
Json(req): Json<AddTripleRequest>,
) -> Result<Json<ijima_core::Triple>, ApiError> {
if !principal.0.may(ijima_core::capabilities::KNOWLEDGE_WRITE) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, None)?;
let triple = kg
.add_triple(
&ns,
EntityId::new(req.subject),
&req.predicate,
EntityId::new(req.object),
req.valid_from.as_deref(),
req.confidence.unwrap_or(1.0),
req.source_memory_id.as_deref(),
)
.await
.map_err(internal)?;
let _ = store;
Ok(Json(triple))
}
async fn query_entity(
principal: AuthPrincipal,
Extension(kg): Extension<Arc<dyn KnowledgeGraph>>,
Path(id): Path<String>,
Query(q): Query<NsQuery>,
) -> Result<Json<ijima_core::EntityRecord>, ApiError> {
if !principal.0.may(KNOWLEDGE_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let rec = kg
.query_entity(&ns, &EntityId::new(id))
.await
.map_err(internal)?;
Ok(Json(rec))
}
async fn invalidate_triple(
principal: AuthPrincipal,
Extension(kg): Extension<Arc<dyn KnowledgeGraph>>,
Path(id): Path<String>,
Query(q): Query<NsQuery>,
) -> Result<StatusCode, ApiError> {
if !principal.0.may(ijima_core::capabilities::KNOWLEDGE_WRITE) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
kg.invalidate_triple(&ns, &id).await.map_err(internal)?;
Ok(StatusCode::NO_CONTENT)
}
#[derive(Deserialize, Default)]
struct FindTriplesQuery {
namespace: Option<String>,
subject: Option<String>,
predicate: Option<String>,
object: Option<String>,
}
async fn find_triples(
principal: AuthPrincipal,
Extension(kg): Extension<Arc<dyn KnowledgeGraph>>,
Query(q): Query<FindTriplesQuery>,
) -> Result<Json<Vec<ijima_core::Triple>>, ApiError> {
if !principal.0.may(KNOWLEDGE_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let triples = kg
.find_triples(
&ns,
q.subject.as_deref().map(EntityId::new).as_ref(),
q.predicate.as_deref(),
q.object.as_deref().map(EntityId::new).as_ref(),
)
.await
.map_err(internal)?;
Ok(Json(triples))
}
async fn kg_timeline(
principal: AuthPrincipal,
Extension(kg): Extension<Arc<dyn KnowledgeGraph>>,
Query(q): Query<NsQuery>,
) -> Result<Json<Vec<ijima_core::Triple>>, ApiError> {
if !principal.0.may(KNOWLEDGE_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let triples = kg
.kg_timeline(&ns, q.limit.unwrap_or(50))
.await
.map_err(internal)?;
Ok(Json(triples))
}
async fn kg_stats(
principal: AuthPrincipal,
Extension(kg): Extension<Arc<dyn KnowledgeGraph>>,
Query(q): Query<NsQuery>,
) -> Result<Json<ijima_core::KgStats>, ApiError> {
if !principal.0.may(KNOWLEDGE_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let stats = kg.knowledge_stats(&ns).await.map_err(internal)?;
Ok(Json(stats))
}
async fn ingest_turn(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(session_id): Path<String>,
Json(mut turn): Json<SessionTurn>,
) -> Result<StatusCode, ApiError> {
if !principal.0.may(SESSION_INGEST) {
return Err(ApiError::Forbidden);
}
let ns = principal.0.personal_namespace();
turn.session_id = SessionId::new(session_id);
store.ingest_turn(&ns, turn).await.map_err(internal)?;
Ok(StatusCode::NO_CONTENT)
}
#[derive(Serialize)]
struct TurnsResponse {
turns: Vec<SessionTurn>,
}
async fn session_turns(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(session_id): Path<String>,
Query(q): Query<NsQuery>,
) -> Result<Json<TurnsResponse>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let turns = store
.session_turns(&ns, &SessionId::new(session_id), q.limit.unwrap_or(50))
.await
.map_err(internal)?;
Ok(Json(TurnsResponse { turns }))
}
async fn create_session(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Json(mut session): Json<Session>,
) -> Result<Json<IdResponse>, ApiError> {
if !principal.0.may(SESSION_INGEST) {
return Err(ApiError::Forbidden);
}
let ns = principal.0.personal_namespace();
if session.started_at.is_empty() {
session.started_at = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs().to_string())
.unwrap_or_default();
}
session.ended_at = None;
let id = store.create_session(&ns, session).await.map_err(internal)?;
Ok(Json(IdResponse { id: id.0 }))
}
#[derive(Deserialize)]
struct SessionListQuery {
namespace: Option<String>,
harness: Option<String>,
limit: Option<usize>,
}
async fn list_sessions(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<SessionListQuery>,
) -> Result<Json<Vec<Session>>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let harness = q.harness.as_deref().map(Harness::from_wire_str);
let limit = q.limit.unwrap_or(50).min(500);
let sessions = store
.list_sessions(&ns, harness.as_ref(), limit)
.await
.map_err(internal)?;
Ok(Json(sessions))
}
#[derive(Deserialize)]
struct EndSessionRequest {
ended_at: String,
}
async fn end_session(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(session_id): Path<String>,
Json(req): Json<EndSessionRequest>,
) -> Result<StatusCode, ApiError> {
if !principal.0.may(SESSION_INGEST) {
return Err(ApiError::Forbidden);
}
let ns = principal.0.personal_namespace();
store
.end_session(&ns, &SessionId::new(session_id), req.ended_at)
.await
.map_err(internal)?;
Ok(StatusCode::NO_CONTENT)
}
async fn list_pending(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<NsQuery>,
) -> Result<Json<Vec<QueuedExtraction>>, ApiError> {
if !principal.0.may(MINING_REVIEW) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let limit = q.limit.unwrap_or(50).min(500);
let pending = store.list_pending(&ns, limit).await.map_err(internal)?;
Ok(Json(pending))
}
async fn accept_extraction(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(id): Path<String>,
) -> Result<Json<AcceptedExtraction>, ApiError> {
if !principal.0.may(MINING_REVIEW) {
return Err(ApiError::Forbidden);
}
let ns = principal.0.personal_namespace();
let accepted = store.accept_extraction(&ns, &id).await.map_err(internal)?;
Ok(Json(accepted))
}
async fn reject_extraction(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(id): Path<String>,
) -> Result<StatusCode, ApiError> {
if !principal.0.may(MINING_REVIEW) {
return Err(ApiError::Forbidden);
}
let ns = principal.0.personal_namespace();
store.reject_extraction(&ns, &id).await.map_err(internal)?;
Ok(StatusCode::NO_CONTENT)
}
#[cfg(feature = "mining")]
async fn trigger_mine(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(session_id): Path<String>,
Query(q): Query<NsQuery>,
) -> Result<Json<crate::mining_pipeline::MiningReport>, ApiError> {
use proserpina::backend::http::HttpAgent;
if !principal.0.may(MINING_TRIGGER) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let turns = store
.session_turns(&ns, &SessionId::new(session_id.clone()), 10_000)
.await
.map_err(internal)?;
let turn_texts: Vec<String> = turns.into_iter().map(|t| t.content).collect();
let ctx = crate::mining_pipeline::mining_context(&session_id, "general", Harness::Other);
let extractions = tokio::task::spawn_blocking(move || {
let mut agent: Option<HttpAgent> = build_mining_agent();
let agent_dyn: Option<&mut dyn proserpina::Agent> =
agent.as_mut().map(|a| a as &mut dyn proserpina::Agent);
ijima_miner::mine_all(&turn_texts, &ctx, agent_dyn)
})
.await
.map_err(|e| {
internal(ijima_core::IjimaError::Mining {
detail: format!("extraction task failed: {e}"),
})
})?
.map_err(internal)?;
let report = crate::mining_pipeline::ingest_extractions(store.as_ref(), &ns, extractions)
.await
.map_err(internal)?;
Ok(Json(report))
}
#[cfg(feature = "mining")]
fn build_mining_agent() -> Option<proserpina::backend::http::HttpAgent> {
use proserpina::{
AgentId, Persona,
backend::http::{HttpAgent, HttpConfig},
};
let base_url = std::env::var("IJIMA_LLM_BASE_URL")
.unwrap_or_else(|_| "https://api.deepseek.com/v1".to_string());
let model = std::env::var("IJIMA_LLM_MODEL").ok()?;
let api_key = std::env::var("IJIMA_LLM_API_KEY").ok()?;
let persona = Persona::new("Session Mining Extractor")
.with_framing(
"You mine session transcripts for durable facts and recurring \
patterns. Output one JSON object per line, each \
{\"content\",\"project\",\"topic\",\"confidence\"}. Omit all \
preamble. If nothing worth extracting, output nothing.",
)
.with_focus(
"decisions, chosen tools, stated constraints, measurements, recurring workflows",
);
Some(HttpAgent::new(
AgentId::new("ijima-miner"),
persona,
HttpConfig {
base_url,
model,
api_key,
},
))
}
#[derive(Deserialize)]
struct NamespaceQuery {
namespace: Option<String>,
}
#[derive(Deserialize)]
struct RoomsQuery {
namespace: Option<String>,
project: Option<String>,
limit: Option<usize>,
}
#[derive(Deserialize)]
struct TunnelQuery {
namespace: Option<String>,
topic: String,
project_a: String,
project_b: String,
limit: Option<usize>,
}
#[derive(Deserialize)]
struct DiaryQuery {
namespace: Option<String>,
limit: Option<usize>,
}
#[derive(Deserialize)]
struct MemoryBrowseQuery {
namespace: Option<String>,
project: Option<String>,
topic: Option<String>,
limit: Option<usize>,
}
#[derive(Deserialize)]
struct ResolveRepoQuery {
cwd: String,
}
async fn list_rooms(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<RoomsQuery>,
) -> Result<Json<Vec<Room>>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let limit = q.limit.unwrap_or(50).min(500);
let rooms = store
.list_rooms(&ns, q.project.as_deref(), limit)
.await
.map_err(internal)?;
Ok(Json(rooms))
}
async fn taxonomy(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<NamespaceQuery>,
) -> Result<Json<Vec<ProjectTaxon>>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
Ok(Json(store.taxonomy(&ns).await.map_err(internal)?))
}
async fn palace_graph(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<NamespaceQuery>,
) -> Result<Json<PalaceGraph>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
Ok(Json(store.palace_graph(&ns).await.map_err(internal)?))
}
async fn traverse_tunnel(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<TunnelQuery>,
) -> Result<Json<TunnelTraversal>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let limit = q.limit.unwrap_or(50).min(500);
Ok(Json(
store
.traverse_tunnel(&ns, &q.topic, &q.project_a, &q.project_b, limit)
.await
.map_err(internal)?,
))
}
async fn write_diary(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Json(entry): Json<DiaryEntry>,
) -> Result<StatusCode, ApiError> {
if !principal.0.may(MEMORY_WRITE) {
return Err(ApiError::Forbidden);
}
let ns = principal.0.personal_namespace();
store.write_diary(&ns, entry).await.map_err(internal)?;
Ok(StatusCode::NO_CONTENT)
}
async fn read_diary(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Path(agent): Path<String>,
Query(q): Query<DiaryQuery>,
) -> Result<Json<Vec<DiaryEntry>>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let limit = q.limit.unwrap_or(50).min(500);
Ok(Json(
store
.read_diary(&ns, &agent, limit)
.await
.map_err(internal)?,
))
}
async fn browse_memories(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<MemoryBrowseQuery>,
) -> Result<Json<Vec<Memory>>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let limit = q.limit.unwrap_or(50).min(500);
Ok(Json(
store
.list_memories_filtered(&ns, q.project.as_deref(), q.topic.as_deref(), limit)
.await
.map_err(internal)?,
))
}
#[derive(Serialize)]
struct NamespaceStats {
total: usize,
projects: Vec<ProjectCount>,
}
#[derive(Serialize)]
struct ProjectCount {
project: String,
count: usize,
}
async fn memory_stats(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<NamespaceQuery>,
) -> Result<Json<NamespaceStats>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
let ns = resolve_ns(&principal, q.namespace.as_deref())?;
let rooms = store.list_rooms(&ns, None, 1000).await.map_err(internal)?;
let total: usize = rooms.iter().map(|r| r.count).sum();
let mut by_project: std::collections::BTreeMap<String, usize> =
std::collections::BTreeMap::new();
for r in &rooms {
*by_project.entry(r.project.clone()).or_default() += r.count;
}
let projects = by_project
.into_iter()
.map(|(project, count)| ProjectCount { project, count })
.collect();
Ok(Json(NamespaceStats { total, projects }))
}
async fn register_repo(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Json(repo): Json<RepoDirectory>,
) -> Result<StatusCode, ApiError> {
if !principal.0.may(ADMIN) {
return Err(ApiError::Forbidden);
}
store.register_repo(repo).await.map_err(internal)?;
Ok(StatusCode::NO_CONTENT)
}
async fn list_repos(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
) -> Result<Json<Vec<RepoDirectory>>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
Ok(Json(store.list_repos().await.map_err(internal)?))
}
async fn resolve_repo(
principal: AuthPrincipal,
Extension(store): Extension<Arc<dyn Store>>,
Query(q): Query<ResolveRepoQuery>,
) -> Result<Json<RepoDirectory>, ApiError> {
if !principal.0.may(MEMORY_READ) {
return Err(ApiError::Forbidden);
}
match store.resolve_repo(&q.cwd).await.map_err(internal)? {
Some(repo) => Ok(Json(repo)),
None => Err(ApiError::NotFound),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::IjimaAuth;
use axum::body::Body;
use axum::http::{Request, StatusCode};
use ijima_core::{harness::Harness, memory::MemorySource};
use tower::ServiceExt;
async fn app_with_store() -> (Router, Arc<IjimaAuth>) {
let auth = Arc::new(IjimaAuth::from_embedded_policy().expect("policy"));
let store_inner = Arc::new(crate::SurrealStore::open_embedded().await.expect("open"));
let store: Arc<dyn Store> = store_inner.clone();
let kg: Arc<dyn KnowledgeGraph> = store_inner;
(
app(
auth.clone(),
store,
kg,
None,
Arc::new(crate::redaction::Redactor::new()),
#[cfg(feature = "rate-limit")]
None,
),
auth,
)
}
fn bearer(auth: &IjimaAuth, principal: &str, cap: &str) -> String {
format!(
"Bearer {}",
auth.issue_bearer(principal, cap).expect("issue")
)
}
fn sample_memory_json(id: &str) -> String {
serde_json::json!({
"id": id,
"content": "decided to wire the daemon",
"project": "ijima",
"topic": "api",
"source": "Explicit",
"harness": "Pi",
"session_id": "sess_1",
"importance": 0.5,
"created_at": "0",
})
.to_string()
}
#[tokio::test]
async fn health_is_public() {
let (app, _) = app_with_store().await;
let res = app
.oneshot(
Request::builder()
.uri("/health")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
}
#[tokio::test]
async fn recall_without_auth_is_401() {
let (app, _) = app_with_store().await;
let res = app
.oneshot(
Request::builder()
.uri("/memories/mem_1")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::UNAUTHORIZED);
}
#[tokio::test]
async fn store_then_recall_round_trips() {
let (app, auth) = app_with_store().await;
let write = bearer(&auth, "elliott", MEMORY_WRITE);
let read = bearer(&auth, "elliott", MEMORY_READ);
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", &write)
.header("content-type", "application/json")
.body(Body::from(sample_memory_json("mem_1")))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let res = app
.oneshot(
Request::builder()
.uri("/memories/mem_1")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let body = axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap();
let mem: Memory = serde_json::from_slice(&body).unwrap();
assert_eq!(mem.content, "decided to wire the daemon");
assert_eq!(mem.harness, Harness::Pi);
assert_eq!(mem.source, MemorySource::Explicit);
}
#[tokio::test]
async fn store_with_read_only_token_is_403() {
let (app, auth) = app_with_store().await;
let read = bearer(&auth, "elliott", MEMORY_READ);
let res = app
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", &read)
.header("content-type", "application/json")
.body(Body::from(sample_memory_json("mem_x")))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
}
#[tokio::test]
async fn namespace_isolation_across_principals() {
let (app, auth) = app_with_store().await;
let alice_write = bearer(&auth, "alice", MEMORY_WRITE);
let _ = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", &alice_write)
.header("content-type", "application/json")
.body(Body::from(sample_memory_json("mem_a")))
.unwrap(),
)
.await
.unwrap();
let bob_read = bearer(&auth, "bob", MEMORY_READ);
let res = app
.oneshot(
Request::builder()
.uri("/memories/mem_a")
.header("authorization", &bob_read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn promote_redacts_secrets_and_leaves_original_intact() {
let (app, auth) = app_with_store().await;
let write = bearer(&auth, "elliott", MEMORY_WRITE);
let read = bearer(&auth, "elliott", MEMORY_READ);
let promote = bearer(&auth, "elliott", TRUST_PROMOTE);
let body = serde_json::json!({
"id": "mem_secret",
"content": "deploy key sk-abcdefghijklmnopqrstuvwxyz1234567890 contact ops@test.com",
"project": "ijima",
"topic": "ops",
"source": "Explicit",
"harness": "Pi",
})
.to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", &write)
.header("content-type", "application/json")
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let promote_body = serde_json::json!({
"target_namespace": "ns_team_shared",
})
.to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories/mem_secret/promote")
.header("authorization", &promote)
.header("content-type", "application/json")
.body(Body::from(promote_body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let resp: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
let new_id = resp["id"].as_str().unwrap();
assert_eq!(new_id, "mem_secret__shared");
let cats: Vec<&str> = resp["redactions"]
.as_array()
.unwrap()
.iter()
.map(|r| r["category"].as_str().unwrap())
.collect();
assert!(cats.contains(&"api_key"));
assert!(cats.contains(&"email"));
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/memories/mem_secret")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
let orig: Memory = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert!(orig.content.contains("sk-abcdef"));
assert!(orig.content.contains("ops@test.com"));
let res = app
.oneshot(
Request::builder()
.uri("/memories/mem_secret__shared?namespace=ns_team_shared")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let shared: Memory = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert!(shared.content.contains("[REDACTED:api_key]"));
assert!(shared.content.contains("[REDACTED:email]"));
assert!(!shared.content.contains("sk-abcdef"));
assert!(!shared.content.contains("ops@test.com"));
assert_eq!(shared.session_id.as_deref(), Some("mem_secret"));
}
#[tokio::test]
async fn promote_requires_trust_promote_not_memory_write() {
let (app, auth) = app_with_store().await;
let write = bearer(&auth, "elliott", MEMORY_WRITE);
let body = serde_json::json!({
"id": "mem_p",
"content": "provenance tier test",
"project": "ijima",
"topic": "t",
"source": "Explicit",
"harness": "Pi",
})
.to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", &write)
.header("content-type", "application/json")
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let promote_body = serde_json::json!({ "target_namespace": "ns_team_shared" }).to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories/mem_p/promote")
.header("authorization", &write)
.header("content-type", "application/json")
.body(Body::from(promote_body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
let promote = bearer(&auth, "elliott", TRUST_PROMOTE);
let promote_body = serde_json::json!({ "target_namespace": "ns_team_shared" }).to_string();
let res = app
.oneshot(
Request::builder()
.method("POST")
.uri("/memories/mem_p/promote")
.header("authorization", &promote)
.header("content-type", "application/json")
.body(Body::from(promote_body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
}
#[tokio::test]
async fn cross_principal_personal_namespace_is_forbidden() {
let (app, auth) = app_with_store().await;
let alice_write = bearer(&auth, "alice", MEMORY_WRITE);
let _ = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", &alice_write)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"id": "mem_a",
"content": "alice only",
"project": "x",
"topic": "x",
"source": "Explicit",
"harness": "Pi",
})
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
let bob_read = bearer(&auth, "bob", MEMORY_READ);
let res = app
.oneshot(
Request::builder()
.uri("/memories/mem_a?namespace=ns_alice_private")
.header("authorization", &bob_read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
}
#[tokio::test]
async fn doctrine_ingest_requires_admin_and_is_readable_shared() {
let (app, auth) = app_with_store().await;
let admin = bearer(&auth, "ci", "admin");
let read = bearer(&auth, "anyone", MEMORY_READ);
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/doctrine")
.header("authorization", &read)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"id": "d1",
"content": "doctrine body",
"project": "ijima",
"topic": "arch",
})
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/doctrine")
.header("authorization", &admin)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"id": "d1",
"content": "doctrine body",
"project": "ijima",
"topic": "arch",
})
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let res = app
.oneshot(
Request::builder()
.uri("/memories/d1?namespace=ns_doctrine")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let mem: Memory = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert_eq!(mem.content, "doctrine body");
assert_eq!(mem.source, ijima_core::memory::MemorySource::Doctrine);
}
#[tokio::test]
async fn wakeup_composes_personal_and_doctrine() {
let (app, auth) = app_with_store().await;
let write = bearer(&auth, "elliott", MEMORY_WRITE);
let admin = bearer(&auth, "ci", "admin");
let read = bearer(&auth, "elliott", MEMORY_READ);
let _ = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", &write)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"id": "mem_p",
"content": "personal essential",
"project": "ijima",
"topic": "x",
"source": "Explicit",
"harness": "Pi",
})
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
let _ = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/doctrine")
.header("authorization", &admin)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"id": "doc_1",
"content": "doctrine baseline",
"project": "ijima",
"topic": "arch",
})
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
let res = app
.oneshot(
Request::builder()
.uri("/wakeup")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let body: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert_eq!(body["identity"]["principal"], "elliott");
assert_eq!(body["personal_essentials"].as_array().unwrap().len(), 1);
assert_eq!(
body["personal_essentials"][0]["content"],
"personal essential"
);
assert_eq!(body["doctrine"].as_array().unwrap().len(), 1);
assert_eq!(body["doctrine"][0]["content"], "doctrine baseline");
assert_eq!(body["doctrine"][0]["source"], "Doctrine");
}
#[tokio::test]
async fn knowledge_graph_add_query_invalidate() {
let (app, auth) = app_with_store().await;
let write = bearer(&auth, "elliott", "knowledge:write");
let read = bearer(&auth, "elliott", "knowledge:read");
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/kg/triples")
.header("authorization", &write)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"subject": "Ijima",
"predicate": "depends_on",
"object": "SurrealDB",
"confidence": 1.0,
})
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/kg/entities/Ijima")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let body: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert_eq!(body["outgoing"].as_array().unwrap().len(), 1);
assert_eq!(body["outgoing"][0]["object"], "SurrealDB");
assert!(body["incoming"].as_array().unwrap().is_empty());
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/kg/stats")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
let body: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert_eq!(body["entities"], 2);
assert_eq!(body["triples"], 1);
let res = app
.oneshot(
Request::builder()
.method("POST")
.uri("/kg/triples/Ijima:depends_on:SurrealDB/invalidate")
.header("authorization", &write)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::NO_CONTENT);
}
#[tokio::test]
async fn status_requires_admin_and_reports_counts() {
let (app, auth) = app_with_store().await;
let admin = bearer(&auth, "op", "admin");
let read = bearer(&auth, "user", MEMORY_READ);
let _ = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", &admin)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"id": "m1",
"content": "stat test",
"project": "x",
"topic": "x",
"source": "Explicit",
"harness": "Pi",
})
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/status")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
let res = app
.oneshot(
Request::builder()
.uri("/status")
.header("authorization", &admin)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let body: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert_eq!(body["memories"], 1);
assert!(!body["namespaces"].as_array().unwrap().is_empty());
}
#[tokio::test]
async fn sessions_create_list_end_via_http() {
let (app, auth) = app_with_store().await;
let ingest = bearer(&auth, "op", SESSION_INGEST);
let read = bearer(&auth, "op", MEMORY_READ);
for (id, harness) in [("sess_a", "Pi"), ("sess_b", "Sakamoto")] {
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/sessions")
.header("authorization", &ingest)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"id": id,
"harness": harness,
"channel": "thread-1",
"started_at": "2026-07-05T10:00:00Z",
})
.to_string(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
}
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/sessions")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let body: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
let arr = body.as_array().unwrap();
assert_eq!(arr.len(), 2);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/sessions?harness=pi")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
let body: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert_eq!(body.as_array().unwrap().len(), 1);
assert_eq!(body[0]["harness"], "Pi");
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/sessions/sess_a/end")
.header("authorization", &ingest)
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({ "ended_at": "2026-07-05T11:00:00Z" }).to_string(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::NO_CONTENT);
let res = app
.oneshot(
Request::builder()
.uri("/sessions?harness=pi")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
let body: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert_eq!(body[0]["ended_at"], "2026-07-05T11:00:00Z");
}
#[tokio::test]
async fn mining_queue_requires_review_capability() {
let (app, auth) = app_with_store().await;
let reviewer = bearer(&auth, "op", MINING_REVIEW);
let reader = bearer(&auth, "op", MEMORY_READ);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/mining/queue")
.header("authorization", &reader)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
let res = app
.oneshot(
Request::builder()
.uri("/mining/queue")
.header("authorization", &reviewer)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let body: serde_json::Value = serde_json::from_slice(
&axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap(),
)
.unwrap();
assert!(body.as_array().unwrap().is_empty());
}
fn hit_mem(id: &str, sim: f32) -> SearchHit {
SearchHit {
memory: Memory {
id: MemoryId(id.into()),
content: id.into(),
project: "p".into(),
topic: "t".into(),
source: ijima_core::MemorySource::Explicit,
harness: ijima_core::harness::Harness::Pi,
session_id: None,
origin: ijima_core::InstanceId::local(),
authority: ijima_core::AuthorityScope::local(),
importance: 0.5,
created_at: "0".into(),
},
similarity: sim,
}
}
#[test]
fn merge_search_hits_ranks_desc_dedups_and_truncates() {
let a = vec![hit_mem("a", 0.9), hit_mem("b", 0.5)];
let b = vec![hit_mem("c", 0.8), hit_mem("a", 0.7)]; let merged = merge_search_hits(a, b, 3);
assert_eq!(merged.len(), 3);
assert_eq!(merged[0].memory.id.0, "a");
assert_eq!((merged[0].similarity * 10.0).round() as i32, 9);
assert_eq!(merged[1].memory.id.0, "c");
assert_eq!(merged[2].memory.id.0, "b");
}
#[test]
fn merge_search_hits_respects_limit() {
let a = vec![hit_mem("a", 0.9), hit_mem("b", 0.8)];
let b = vec![hit_mem("c", 0.7), hit_mem("d", 0.6)];
let merged = merge_search_hits(a, b, 2);
assert_eq!(merged.len(), 2);
assert_eq!(merged[0].memory.id.0, "a");
assert_eq!(merged[1].memory.id.0, "b");
}
#[cfg(feature = "mining")]
#[tokio::test]
async fn trigger_requires_mining_trigger_capability() {
let (app, auth) = app_with_store().await;
let write = bearer(&auth, "elliott", MEMORY_WRITE);
let res = app
.oneshot(
Request::builder()
.method("POST")
.uri("/sessions/sess_x/mine")
.header("authorization", &write)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
}
#[cfg(feature = "mining")]
#[tokio::test]
async fn trigger_mines_decision_and_archives() {
let (app, auth) = app_with_store().await;
let ingest = bearer(&auth, "elliott", SESSION_INGEST);
let trigger = bearer(&auth, "elliott", MINING_TRIGGER);
let turn = serde_json::json!({
"session_id": "sess_mine",
"turn_index": 0,
"role": "User",
"content": "We decided to use SurrealDB for storage.",
"timestamp": "0",
});
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/sessions/sess_mine/turns")
.header("authorization", &ingest)
.header("content-type", "application/json")
.body(Body::from(turn.to_string()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::NO_CONTENT);
let res = app
.oneshot(
Request::builder()
.method("POST")
.uri("/sessions/sess_mine/mine")
.header("authorization", &trigger)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let body = axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap();
let report: crate::mining_pipeline::MiningReport = serde_json::from_slice(&body).unwrap();
assert!(
report.archived >= 1,
"rules tier should archive the decision: {report:?}"
);
}
async fn body_json(res: axum::response::Response) -> serde_json::Value {
let body = axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap();
serde_json::from_slice(&body).unwrap()
}
async fn seed_memory(app: &Router, auth: &IjimaAuth, id: &str, project: &str, topic: &str) {
let body = serde_json::json!({
"id": id,
"content": format!("{project}/{topic} note"),
"project": project,
"topic": topic,
"source": "Explicit",
"harness": "Pi",
"session_id": "sess_1",
"importance": 0.5,
"created_at": "0",
})
.to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/memories")
.header("authorization", bearer(auth, "elliott", MEMORY_WRITE))
.header("content-type", "application/json")
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK, "seed {id} failed");
}
#[tokio::test]
async fn rooms_taxonomy_stats_reflect_seeded_memories() {
let (app, auth) = app_with_store().await;
seed_memory(&app, &auth, "mem_a", "ijima", "api").await;
seed_memory(&app, &auth, "mem_b", "ijima", "auth").await;
let read = bearer(&auth, "elliott", MEMORY_READ);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/rooms")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let rooms = body_json(res).await;
let topics: std::collections::HashSet<&str> = rooms
.as_array()
.unwrap()
.iter()
.map(|r| r["topic"].as_str().unwrap())
.collect();
assert!(
topics.contains("api") && topics.contains("auth"),
"rooms: {rooms}"
);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/memories/stats")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let stats = body_json(res).await;
assert_eq!(stats["total"], 2, "stats: {stats}");
assert_eq!(stats["projects"][0]["project"], "ijima");
assert_eq!(stats["projects"][0]["count"], 2);
}
#[tokio::test]
async fn browse_memories_filters_by_project() {
let (app, auth) = app_with_store().await;
seed_memory(&app, &auth, "mem_a", "ijima", "api").await;
seed_memory(&app, &auth, "mem_b", "possum", "efficiency").await;
let read = bearer(&auth, "elliott", MEMORY_READ);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/memories?project=possum")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let mems = body_json(res).await;
let arr = mems.as_array().unwrap();
assert_eq!(arr.len(), 1);
assert_eq!(arr[0]["project"], "possum");
}
#[tokio::test]
async fn palace_graph_and_tunnel_link_shared_topic() {
let (app, auth) = app_with_store().await;
seed_memory(&app, &auth, "mem_a", "ijima", "efficiency").await;
seed_memory(&app, &auth, "mem_b", "possum", "efficiency").await;
let read = bearer(&auth, "elliott", MEMORY_READ);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/palace/graph")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let graph = body_json(res).await;
let projects: std::collections::HashSet<&str> = graph["projects"]
.as_array()
.unwrap()
.iter()
.map(|p| p.as_str().unwrap())
.collect();
assert!(
projects.contains("ijima") && projects.contains("possum"),
"graph: {graph}"
);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/palace/tunnel?topic=efficiency&project_a=ijima&project_b=possum")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let trav = body_json(res).await;
assert_eq!(trav["memories_a"].as_array().unwrap().len(), 1);
assert_eq!(trav["memories_b"].as_array().unwrap().len(), 1);
}
#[tokio::test]
async fn diary_write_then_read_round_trips() {
let (app, auth) = app_with_store().await;
let write = bearer(&auth, "elliott", MEMORY_WRITE);
let read = bearer(&auth, "elliott", MEMORY_READ);
let body = serde_json::json!({
"agent": "pi",
"content": "shipped the routes",
"topic": "ijima",
"timestamp": "2026-08-09T12:00:00Z"
})
.to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/diaries")
.header("authorization", &write)
.header("content-type", "application/json")
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::NO_CONTENT);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/diaries/pi")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let entries = body_json(res).await;
let arr = entries.as_array().unwrap();
assert_eq!(arr.len(), 1);
assert_eq!(arr[0]["content"], "shipped the routes");
}
#[tokio::test]
async fn diary_write_requires_memory_write_not_read() {
let (app, auth) = app_with_store().await;
let read = bearer(&auth, "elliott", MEMORY_READ);
let body = serde_json::json!({"agent": "pi", "content": "x", "timestamp": "t"}).to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/diaries")
.header("authorization", &read)
.header("content-type", "application/json")
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
}
#[tokio::test]
async fn repo_register_list_resolve_round_trips() {
let (app, auth) = app_with_store().await;
let admin = bearer(&auth, "elliott", ADMIN);
let read = bearer(&auth, "elliott", MEMORY_READ);
let body = serde_json::json!({
"name": "Ijima",
"path": "/home/x/Ijima",
"remote_url": "git@github.com:Industrial-Algebra/Ijima.git",
"role": "memory-service"
})
.to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/repos")
.header("authorization", &admin)
.header("content-type", "application/json")
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::NO_CONTENT);
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/repos")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let repos = body_json(res).await;
assert_eq!(repos[0]["name"], "Ijima");
assert_eq!(repos[0]["path"], "/home/x/Ijima");
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/repos/resolve?cwd=/home/x/Ijima/src")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::OK);
let repo = body_json(res).await;
assert_eq!(repo["name"], "Ijima");
let res = app
.clone()
.oneshot(
Request::builder()
.uri("/repos/resolve?cwd=/nowhere/here")
.header("authorization", &read)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn repo_register_requires_admin() {
let (app, auth) = app_with_store().await;
let read = bearer(&auth, "elliott", MEMORY_READ);
let body = serde_json::json!({
"name": "X", "path": "/x", "remote_url": "u", "role": "r"
})
.to_string();
let res = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/repos")
.header("authorization", &read)
.header("content-type", "application/json")
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(res.status(), StatusCode::FORBIDDEN);
}
}