This repository has been archived on 2026-08-20. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
el-retired/engram/crates/engram-server/src/main.rs
T

258 lines
9.1 KiB
Rust

/// Engram Server — HTTP API for Engram with sync and swarm activation.
///
/// # Endpoints
///
/// ## Core
/// GET /stats — node/edge counts
/// POST /nodes — create a node
/// GET /nodes/{id} — get a node
/// POST /edges — create an edge
/// GET /nodes/{id}/edges — list edges from a node
/// POST /activate — spreading activation
/// POST /search — embedding search
/// POST /decay — apply salience decay
/// POST /consolidate — promote Episodic → Semantic
///
/// ## Sync (auth required)
/// GET /sync/delta?since={ms}&peer_id={uuid} — generate delta
/// POST /sync/push — receive incoming delta
/// POST /sync/peers — register peer
/// GET /sync/peers — list peers
/// DELETE /sync/peers/{id} — remove peer
///
/// ## Swarm
/// POST /swarm/activate — distributed activation (auth required)
/// GET /swarm/status — peer health (auth required)
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use axum::{
body::Body,
http::{header, StatusCode},
middleware,
response::{IntoResponse, Response},
routing::{delete, get, post},
Router,
};
use engram_core::EngramDb;
use engram_projection::registry::ProjectionRegistry;
use engram_sync::{SyncConfig, SyncEngine};
use engram_tx::TransactionEngine;
use mime_guess::from_path;
use rust_embed::RustEmbed;
use tokio::time::interval;
use tower_http::cors::CorsLayer;
use tracing::info;
use tracing_subscriber::EnvFilter;
#[derive(RustEmbed)]
#[folder = "../../studio/"]
struct Studio;
async fn serve_studio_index() -> impl IntoResponse {
match Studio::get("index.html") {
Some(content) => Response::builder()
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, "text/html; charset=utf-8")
.body(Body::from(content.data.into_owned()))
.unwrap(),
None => Response::builder()
.status(StatusCode::NOT_FOUND)
.body(Body::from("Studio not found"))
.unwrap(),
}
}
async fn serve_studio_asset(uri: axum::extract::Path<String>) -> impl IntoResponse {
let path = uri.0;
match Studio::get(&path) {
Some(content) => {
let mime = from_path(&path).first_or_octet_stream();
Response::builder()
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, mime.as_ref())
.body(Body::from(content.data.into_owned()))
.unwrap()
}
None => Response::builder()
.status(StatusCode::NOT_FOUND)
.body(Body::from("Not found"))
.unwrap(),
}
}
mod auth;
mod routes;
mod state;
use auth::require_auth;
use state::AppState;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
// Logging
tracing_subscriber::fmt()
.with_env_filter(
EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")),
)
.init();
// Configuration from environment (with sensible defaults)
let db_path = std::env::var("ENGRAM_DB_PATH").unwrap_or_else(|_| "./engram-data".to_string());
let bind_addr = std::env::var("ENGRAM_BIND").unwrap_or_else(|_| "0.0.0.0:8742".to_string());
let api_key = std::env::var("ENGRAM_API_KEY").unwrap_or_else(|_| {
let key = uuid::Uuid::new_v4().to_string();
eprintln!("No ENGRAM_API_KEY set — generated key: {}", key);
key
});
// Open database
let db = EngramDb::open(&PathBuf::from(&db_path))?;
let db = Arc::new(Mutex::new(db));
// Transaction engine (separate sled db alongside the main db)
let tx_log_path = format!("{}-tx-log", db_path);
let tx_log_db = sled::open(&tx_log_path)?;
let tx_engine = Arc::new(Mutex::new(TransactionEngine::new(
db.clone(),
tx_log_db,
Some(uuid::Uuid::new_v4()),
)));
// Projection registry
let projection_registry = Arc::new(Mutex::new(ProjectionRegistry::new()));
info!("Database opened at {}", db_path);
// Sync engine — wrapped in tokio::sync::Mutex so it can be held across .await
let sync_config = SyncConfig {
our_id: uuid::Uuid::new_v4(),
our_name: std::env::var("ENGRAM_PEER_NAME").unwrap_or_else(|_| "engram-local".to_string()),
api_key: api_key.clone(),
default_sync_tiers: vec![engram_core::types::MemoryTier::Semantic],
sync_interval_secs: std::env::var("ENGRAM_SYNC_INTERVAL_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(300),
};
let sync_interval_secs = sync_config.sync_interval_secs;
let sync_engine = Arc::new(tokio::sync::Mutex::new(SyncEngine::new(db.clone(), sync_config)));
{
let e = sync_engine.lock().await;
info!(
peer_name = e.our_name(),
peer_id = %e.our_id(),
sync_interval_secs,
"Sync engine ready"
);
}
// Background sync task — tokio::sync::Mutex guard is Send-safe
{
let engine_arc = sync_engine.clone();
tokio::spawn(async move {
let mut ticker = interval(Duration::from_secs(sync_interval_secs));
loop {
ticker.tick().await;
let report = {
let mut e = engine_arc.lock().await;
e.sync_all().await
};
if report.peers_synced > 0 || !report.errors.is_empty() {
info!(
peers_synced = report.peers_synced,
nodes_received = report.nodes_received,
nodes_sent = report.nodes_sent,
errors = report.errors.len(),
"Sync cycle complete"
);
}
}
});
}
// Shared state
let state = Arc::new(AppState {
db: db.clone(),
sync_engine: sync_engine.clone(),
api_key: api_key.clone(),
projection_registry,
tx_engine,
});
// Protected sync/swarm routes (auth middleware applied)
let sync_routes = Router::new()
.route("/sync/delta", get(routes::sync::get_delta))
.route("/sync/push", post(routes::sync::push_delta))
.route("/sync/peers", get(routes::sync::list_peers))
.route("/sync/peers", post(routes::sync::register_peer))
.route("/sync/peers/{id}", delete(routes::sync::delete_peer))
.route("/swarm/activate", post(routes::swarm::swarm_activate))
.route("/swarm/status", get(routes::swarm::swarm_status))
.layer(middleware::from_fn_with_state(
state.clone(),
require_auth,
));
// Open core routes (no auth)
let core_routes = Router::new()
.route("/stats", get(routes::core::get_stats))
.route("/nodes", post(routes::core::create_node))
.route("/nodes/{id}", get(routes::core::get_node))
.route("/edges", post(routes::core::create_edge))
.route("/nodes/{id}/edges", get(routes::core::get_edges_from))
.route("/activate", post(routes::core::activate))
.route("/search", post(routes::core::search_embedding))
.route("/decay", post(routes::core::decay))
.route("/consolidate", post(routes::core::consolidate));
// Projection routes (no auth)
let projection_routes = Router::new()
.route("/projections", post(routes::projection::register_projection))
.route("/projections", get(routes::projection::list_projections))
.route(
"/projections/{name}/schema",
get(routes::projection::get_projection_schema),
)
.route(
"/projections/{name}/query",
post(routes::projection::query_projection),
);
// Transaction routes (no auth — add auth layer if needed)
let tx_routes = Router::new()
.route("/tx/apply", post(routes::tx::tx_apply))
.route("/tx/rollback/{command_id}", post(routes::tx::tx_rollback))
.route("/tx/history", get(routes::tx::tx_history))
.route("/tx/chain/{command_id}", get(routes::tx::tx_causal_chain));
// Reasoning routes (no auth — graph-native inference)
let reasoning_routes = Router::new()
.route("/reason", post(routes::reasoning::reason))
.route("/reason/causal", post(routes::reasoning::causal))
.route("/reason/contradictions", post(routes::reasoning::contradictions));
let studio_routes = Router::new()
.route("/", get(serve_studio_index))
.route("/studio", get(serve_studio_index))
.route("/studio/{*path}", get(serve_studio_asset));
let app = Router::new()
.merge(studio_routes)
.merge(core_routes)
.merge(sync_routes)
.merge(projection_routes)
.merge(tx_routes)
.merge(reasoning_routes)
.layer(CorsLayer::permissive())
.with_state(state);
let listener = tokio::net::TcpListener::bind(&bind_addr).await?;
info!("Engram server listening on {}", bind_addr);
axum::serve(listener, app).await?;
Ok(())
}