docs: update README and docs for v0.6.0 architecture and deployment, preserve current server/shared/wasm updates
This commit is contained in:
1 parent
6067746898
commit
2b8afd54e0
27 files changed
+1413
-2567
No files matched your search
@@ -30,3 +30,4 @@ hex = "0.4"
|
||||
base64 = "0.22"
|
||||
rand = "0.8"
|
||||
ed25519-dalek = "2"
|
||||
valkey = "0.0.0-alpha5"
|
||||
@@ -5,16 +5,10 @@ pub async fn cleanup_loop(state: Arc<AppState>) {
|
||||
loop {
|
||||
tokio::time::sleep(std::time::Duration::from_secs(60)).await;
|
||||
|
||||
// Evict expired sessions from SQLite.
|
||||
// Evict expired sessions from the configured storage backend.
|
||||
{
|
||||
if let Ok(conn) = state.db_pool.get() {
|
||||
let now = crate::storage::current_time_ms();
|
||||
let _ = conn.execute(
|
||||
"DELETE FROM sessions WHERE expires_at < ?1",
|
||||
rusqlite::params![now],
|
||||
);
|
||||
} else {
|
||||
tracing::error!("Failed to get database connection from pool for cleanup");
|
||||
if let Err(err) = state.db_pool.delete_expired_sessions() {
|
||||
tracing::error!("Failed to evict expired sessions: {}", err);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -123,6 +123,10 @@ pub struct RunArgs {
|
||||
/// Optional structured JSON log file.
|
||||
#[arg(long, env = "CHRONOSEAL_LOG_FILE")]
|
||||
pub log_file: Option<PathBuf>,
|
||||
|
||||
/// Number of mutation rounds to execute per program.
|
||||
#[arg(long, env = "CHRONOSEAL_MUTATION_ROUNDS")]
|
||||
pub mutation_rounds: Option<u8>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Args)]
|
||||
|
||||
@@ -46,6 +46,7 @@ pub struct Config {
|
||||
pub min_pause_count: u32,
|
||||
pub require_mouse_activity: bool,
|
||||
pub gene_size: usize,
|
||||
pub mutation_rounds: u8,
|
||||
}
|
||||
|
||||
impl Default for Config {
|
||||
@@ -68,6 +69,7 @@ impl Default for Config {
|
||||
min_pause_count: 1,
|
||||
require_mouse_activity: true,
|
||||
gene_size: shared::constants::DEFAULT_GENE_SIZE,
|
||||
mutation_rounds: shared::constants::DEFAULT_MUTATION_ROUNDS,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -118,6 +120,9 @@ impl Config {
|
||||
if let Some(log_file) = &args.log_file {
|
||||
self.log_file = Some(log_file.clone());
|
||||
}
|
||||
if let Some(mutation_rounds) = args.mutation_rounds {
|
||||
self.mutation_rounds = mutation_rounds;
|
||||
}
|
||||
}
|
||||
|
||||
pub fn validate(&self) -> Result<(), ConfigError> {
|
||||
@@ -132,6 +137,11 @@ impl Config {
|
||||
size: self.gene_size,
|
||||
});
|
||||
}
|
||||
if !(1..=shared::constants::MAX_MUTATION_ROUNDS).contains(&self.mutation_rounds) {
|
||||
return Err(ConfigError::InvalidMutationRounds {
|
||||
rounds: self.mutation_rounds,
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -214,6 +224,11 @@ impl Config {
|
||||
self.gene_size = val;
|
||||
}
|
||||
}
|
||||
if let Ok(value) = env::var("CHRONOSEAL_MUTATION_ROUNDS") {
|
||||
if let Ok(val) = value.parse() {
|
||||
self.mutation_rounds = val;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -234,6 +249,9 @@ pub enum ConfigError {
|
||||
InvalidGeneSize {
|
||||
size: usize,
|
||||
},
|
||||
InvalidMutationRounds {
|
||||
rounds: u8,
|
||||
},
|
||||
}
|
||||
|
||||
impl std::fmt::Display for ConfigError {
|
||||
@@ -253,6 +271,13 @@ impl std::fmt::Display for ConfigError {
|
||||
shared::constants::MAX_GENE_SIZE
|
||||
)
|
||||
}
|
||||
Self::InvalidMutationRounds { rounds } => {
|
||||
write!(
|
||||
f,
|
||||
"invalid mutation rounds {rounds}; expected 1..={}",
|
||||
shared::constants::MAX_MUTATION_ROUNDS
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -316,6 +341,7 @@ mod tests {
|
||||
db_path: None,
|
||||
frontend_dir: None,
|
||||
log_file: None,
|
||||
mutation_rounds: None,
|
||||
};
|
||||
cfg.apply_run_args(&args);
|
||||
assert_eq!(cfg.db_type, DbType::SqliteInDisk);
|
||||
|
||||
@@ -14,6 +14,9 @@ pub enum SessionError {
|
||||
#[error("Database error: {0}")]
|
||||
Database(#[from] rusqlite::Error),
|
||||
|
||||
#[error("Storage error: {0}")]
|
||||
Storage(String),
|
||||
|
||||
#[error("R2D2 pool error: {0}")]
|
||||
Pool(#[from] r2d2::Error),
|
||||
|
||||
@@ -48,6 +51,9 @@ pub enum VerificationError {
|
||||
#[error("Database error: {0}")]
|
||||
Database(#[from] rusqlite::Error),
|
||||
|
||||
#[error("Storage error: {0}")]
|
||||
Storage(String),
|
||||
|
||||
#[error("Hex decoding error: {0}")]
|
||||
Hex(#[from] hex::FromHexError),
|
||||
|
||||
|
||||
@@ -29,22 +29,7 @@ pub async fn handler(
|
||||
}
|
||||
|
||||
let config = state.get_config();
|
||||
let conn = match state.db_pool.get() {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
tracing::error!("Db pool error: {}", e);
|
||||
return (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(HeartbeatResponse {
|
||||
status: "error".into(),
|
||||
next_salt: None,
|
||||
next_mutation_step: None,
|
||||
next_mutation_order_b64: None,
|
||||
}),
|
||||
);
|
||||
}
|
||||
};
|
||||
match crate::session::verify_heartbeat(&conn, &config, &payload) {
|
||||
match crate::session::verify_heartbeat(&state.db_pool, &config, &payload) {
|
||||
Ok(result) => (
|
||||
StatusCode::OK,
|
||||
Json(HeartbeatResponse {
|
||||
@@ -145,7 +130,11 @@ mod tests {
|
||||
hardware_concurrency: 8,
|
||||
},
|
||||
mutation_step,
|
||||
gene_commitment: shared::gene::commitment_hex(&candidate),
|
||||
gene_commitment: shared::gene::commitment_hex_with_context(
|
||||
&candidate,
|
||||
&init.session_id,
|
||||
mutation_step,
|
||||
),
|
||||
signature: String::new(),
|
||||
};
|
||||
sign_request(sk, &mut req);
|
||||
@@ -157,7 +146,7 @@ mod tests {
|
||||
) -> (Arc<AppState>, InitResponse, SigningKey) {
|
||||
let pool = crate::storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let state = Arc::new(AppState {
|
||||
db_pool: pool,
|
||||
db_pool: pool.clone(),
|
||||
rate_limiter: tokio::sync::Mutex::new(crate::ratelimit::RateLimiter::new()),
|
||||
config: std::sync::RwLock::new(config.clone()),
|
||||
});
|
||||
@@ -165,8 +154,7 @@ mod tests {
|
||||
let mut rng = rand::thread_rng();
|
||||
let sk = SigningKey::generate(&mut rng);
|
||||
let pk_hex = hex::encode(sk.verifying_key().to_bytes());
|
||||
let conn = state.db_pool.get().unwrap();
|
||||
let init = crate::session::create_session(&conn, &config, &pk_hex).unwrap();
|
||||
let init = crate::session::create_session(&pool, &config, &pk_hex).unwrap();
|
||||
(state, init, sk)
|
||||
}
|
||||
|
||||
|
||||
@@ -9,7 +9,6 @@ pub async fn handler(
|
||||
Json(payload): Json<InitRequest>,
|
||||
) -> Result<Json<InitResponse>, SessionError> {
|
||||
let config = state.get_config();
|
||||
let conn = state.db_pool.get()?;
|
||||
let resp = crate::session::create_session(&conn, &config, &payload.public_key)?;
|
||||
let resp = crate::session::create_session(&state.db_pool, &config, &payload.public_key)?;
|
||||
Ok(Json(resp))
|
||||
}
|
||||
+18
-31
@@ -204,14 +204,7 @@ pub fn db_type_report() -> DbTypeReport {
|
||||
}
|
||||
|
||||
fn init_db_pool(config: &Config) -> Result<storage::DbPool, Box<dyn std::error::Error>> {
|
||||
match config.db_type {
|
||||
crate::config::DbType::SqliteInMemory => storage::init_pool(Path::new(":memory:")),
|
||||
crate::config::DbType::SqliteInDisk => storage::init_pool(&config.db_path),
|
||||
crate::config::DbType::Valkey => {
|
||||
warn!("db_type=valkey selected; using sqlite-in-memory compatibility mode in v0.6.0");
|
||||
storage::init_pool(Path::new(":memory:"))
|
||||
}
|
||||
}
|
||||
storage::DbPool::init(config)
|
||||
}
|
||||
|
||||
pub fn probe_health(config: &Config) -> HealthReport {
|
||||
@@ -278,11 +271,9 @@ async fn health_handler() -> impl IntoResponse {
|
||||
async fn stats_handler(
|
||||
axum::extract::State(state): axum::extract::State<Arc<session::AppState>>,
|
||||
) -> Result<Json<StoreStats>, (StatusCode, String)> {
|
||||
let db = state
|
||||
state
|
||||
.db_pool
|
||||
.get()
|
||||
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?;
|
||||
storage::stats(&db)
|
||||
.stats()
|
||||
.map(Json)
|
||||
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))
|
||||
}
|
||||
@@ -290,11 +281,9 @@ async fn stats_handler(
|
||||
async fn metrics_handler(
|
||||
axum::extract::State(state): axum::extract::State<Arc<session::AppState>>,
|
||||
) -> Result<String, (StatusCode, String)> {
|
||||
let db = state
|
||||
state
|
||||
.db_pool
|
||||
.get()
|
||||
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?;
|
||||
storage::stats(&db)
|
||||
.stats()
|
||||
.map(|stats| {
|
||||
format!(
|
||||
"# HELP chronoseal_sessions Active ChronoSeal sessions\n# TYPE chronoseal_sessions gauge\nchronoseal_sessions {}\n# HELP chronoseal_expired_sessions Expired sessions not yet removed\n# TYPE chronoseal_expired_sessions gauge\nchronoseal_expired_sessions {}\n# HELP chronoseal_max_chain_length Maximum heartbeat chain length\n# TYPE chronoseal_max_chain_length gauge\nchronoseal_max_chain_length {}\n",
|
||||
@@ -430,6 +419,7 @@ mod tests {
|
||||
min_pause_count: 0,
|
||||
require_mouse_activity: false,
|
||||
gene_size: shared::constants::DEFAULT_GENE_SIZE,
|
||||
mutation_rounds: shared::constants::DEFAULT_MUTATION_ROUNDS,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -445,11 +435,10 @@ mod tests {
|
||||
fn test_init_db_pool_sqlite_in_memory() {
|
||||
let config = base_config();
|
||||
let pool = init_db_pool(&config).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let count: u64 = conn
|
||||
.query_row("SELECT COUNT(*) FROM sessions", [], |row| row.get(0))
|
||||
.unwrap();
|
||||
assert_eq!(count, 0);
|
||||
let stats = pool.stats().unwrap();
|
||||
assert_eq!(stats.sessions, 0);
|
||||
assert_eq!(stats.expired_sessions, 0);
|
||||
assert_eq!(stats.max_chain_length, 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -459,11 +448,10 @@ mod tests {
|
||||
config.db_path = std::path::PathBuf::from("/tmp/chronoseal-db-type-disk.sqlite");
|
||||
let _ = std::fs::remove_file(&config.db_path);
|
||||
let pool = init_db_pool(&config).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let count: u64 = conn
|
||||
.query_row("SELECT COUNT(*) FROM sessions", [], |row| row.get(0))
|
||||
.unwrap();
|
||||
assert_eq!(count, 0);
|
||||
let stats = pool.stats().unwrap();
|
||||
assert_eq!(stats.sessions, 0);
|
||||
assert_eq!(stats.expired_sessions, 0);
|
||||
assert_eq!(stats.max_chain_length, 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -471,10 +459,9 @@ mod tests {
|
||||
let mut config = base_config();
|
||||
config.db_type = crate::config::DbType::Valkey;
|
||||
let pool = init_db_pool(&config).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let count: u64 = conn
|
||||
.query_row("SELECT COUNT(*) FROM sessions", [], |row| row.get(0))
|
||||
.unwrap();
|
||||
assert_eq!(count, 0);
|
||||
let stats = pool.stats().unwrap();
|
||||
assert_eq!(stats.sessions, 0);
|
||||
assert_eq!(stats.expired_sessions, 0);
|
||||
assert_eq!(stats.max_chain_length, 0);
|
||||
}
|
||||
}
|
||||
+114
-147
@@ -15,7 +15,6 @@ impl AppState {
|
||||
}
|
||||
|
||||
use crate::{crypto, fingerprint, storage, trust, vm};
|
||||
use rusqlite::params;
|
||||
use shared::{
|
||||
gene::{self, GeneState},
|
||||
protocol::{HeartbeatRequest, InitResponse},
|
||||
@@ -30,7 +29,7 @@ pub struct HeartbeatVerificationResult {
|
||||
}
|
||||
|
||||
pub fn create_session(
|
||||
conn: &rusqlite::Connection,
|
||||
db: &storage::DbPool,
|
||||
config: &crate::config::Config,
|
||||
pub_key_hex: &str,
|
||||
) -> Result<InitResponse, crate::errors::SessionError> {
|
||||
@@ -56,25 +55,23 @@ pub fn create_session(
|
||||
let initial_mutation = vm_extensions::generate_order(1, config.gene_size);
|
||||
let initial_mutation_b64 = vm_extensions::encode_order_b64(&initial_mutation);
|
||||
|
||||
conn.execute(
|
||||
"INSERT INTO sessions (
|
||||
session_id, public_key, salt, last_hash, chain_length, created_at, last_seen, expires_at,
|
||||
gene, environment, pending_mutation, pending_mutation_step
|
||||
) VALUES (?1, ?2, ?3, ?4, 1, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
|
||||
params![
|
||||
session_id,
|
||||
pub_key,
|
||||
salt.to_vec(),
|
||||
initial_hash,
|
||||
now,
|
||||
now,
|
||||
expires_at,
|
||||
gene_state.gene,
|
||||
environment_blob,
|
||||
initial_mutation.program,
|
||||
initial_mutation.step,
|
||||
],
|
||||
)?;
|
||||
let record = storage::SessionRecord {
|
||||
session_id: session_id.clone(),
|
||||
public_key: pub_key,
|
||||
salt: salt.to_vec(),
|
||||
last_hash: initial_hash.clone(),
|
||||
chain_length: 1,
|
||||
created_at: now,
|
||||
last_seen: now,
|
||||
expires_at,
|
||||
gene: gene_state.gene,
|
||||
environment: environment_blob,
|
||||
pending_mutation: initial_mutation.program,
|
||||
pending_mutation_step: initial_mutation.step,
|
||||
};
|
||||
|
||||
db.insert_session(&record)
|
||||
.map_err(|err| crate::errors::SessionError::Storage(err.to_string()))?;
|
||||
|
||||
Ok(InitResponse {
|
||||
session_id,
|
||||
@@ -91,85 +88,55 @@ pub fn create_session(
|
||||
}
|
||||
|
||||
pub fn verify_heartbeat(
|
||||
conn: &rusqlite::Connection,
|
||||
db: &storage::DbPool,
|
||||
config: &crate::config::Config,
|
||||
req: &HeartbeatRequest,
|
||||
) -> Result<HeartbeatVerificationResult, crate::errors::VerificationError> {
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT public_key, salt, last_hash, expires_at, gene, environment, pending_mutation, pending_mutation_step
|
||||
FROM sessions WHERE session_id = ?1",
|
||||
)?;
|
||||
let (
|
||||
pub_key,
|
||||
salt,
|
||||
stored_last_hash,
|
||||
expires_at,
|
||||
gene_blob,
|
||||
environment_blob,
|
||||
pending_mutation,
|
||||
pending_step,
|
||||
): (
|
||||
Vec<u8>,
|
||||
Vec<u8>,
|
||||
Vec<u8>,
|
||||
u64,
|
||||
Vec<u8>,
|
||||
Vec<u8>,
|
||||
Vec<u8>,
|
||||
u64,
|
||||
) = stmt
|
||||
.query_row(params![req.session_id], |row| {
|
||||
Ok((
|
||||
row.get(0)?,
|
||||
row.get(1)?,
|
||||
row.get(2)?,
|
||||
row.get(3)?,
|
||||
row.get(4)?,
|
||||
row.get(5)?,
|
||||
row.get(6)?,
|
||||
row.get(7)?,
|
||||
))
|
||||
})
|
||||
.map_err(|e| {
|
||||
if matches!(e, rusqlite::Error::QueryReturnedNoRows) {
|
||||
crate::errors::VerificationError::SessionNotFound
|
||||
} else {
|
||||
crate::errors::VerificationError::Database(e)
|
||||
}
|
||||
})?;
|
||||
let session = db
|
||||
.load_session(&req.session_id)
|
||||
.map_err(|e| crate::errors::VerificationError::Storage(e.to_string()))?;
|
||||
let session = session.ok_or(crate::errors::VerificationError::SessionNotFound)?;
|
||||
|
||||
let now = storage::current_time_ms();
|
||||
if now > expires_at {
|
||||
if now > session.expires_at {
|
||||
return Err(crate::errors::VerificationError::Expired);
|
||||
}
|
||||
|
||||
// 1. Verify signature
|
||||
crypto::verify_signature(&pub_key, req)
|
||||
crypto::verify_signature(&session.public_key, req)
|
||||
.map_err(|e| crate::errors::VerificationError::Signature(e.to_string()))?;
|
||||
|
||||
// 2. Check chain continuity
|
||||
let prev_hash_bytes = hex::decode(&req.prev_hash)?;
|
||||
if stored_last_hash != prev_hash_bytes {
|
||||
if session.last_hash != prev_hash_bytes {
|
||||
return Err(crate::errors::VerificationError::ChainBroken);
|
||||
}
|
||||
|
||||
// 3. Mutation step and deterministic mutation parity
|
||||
if req.mutation_step != pending_step {
|
||||
if req.mutation_step != session.pending_mutation_step {
|
||||
return Err(crate::errors::VerificationError::MutationStepMismatch {
|
||||
expected: pending_step,
|
||||
expected: session.pending_mutation_step,
|
||||
got: req.mutation_step,
|
||||
});
|
||||
}
|
||||
|
||||
let environment = gene::decode_environment(&environment_blob)
|
||||
let environment = gene::decode_environment(&session.environment)
|
||||
.map_err(|e| crate::errors::VerificationError::GeneState(e.to_string()))?;
|
||||
let server_state = GeneState {
|
||||
gene: gene_blob,
|
||||
gene: session.gene.clone(),
|
||||
environment,
|
||||
};
|
||||
let candidate_state = vm_extensions::apply_program_clone(&server_state, &pending_mutation)
|
||||
.map_err(|e| crate::errors::VerificationError::MutationProgram(e.to_string()))?;
|
||||
let expected_gene_commitment = gene::commitment_hex(&candidate_state);
|
||||
let candidate_state = vm_extensions::apply_program_clone_with_rounds(
|
||||
&server_state,
|
||||
&session.pending_mutation,
|
||||
config.mutation_rounds,
|
||||
)
|
||||
.map_err(|e| crate::errors::VerificationError::MutationProgram(e.to_string()))?;
|
||||
let expected_gene_commitment = gene::commitment_hex_with_context(
|
||||
&candidate_state,
|
||||
&req.session_id,
|
||||
req.mutation_step,
|
||||
);
|
||||
if req.gene_commitment != expected_gene_commitment {
|
||||
return Err(crate::errors::VerificationError::MutationCommitmentMismatch);
|
||||
}
|
||||
@@ -192,11 +159,11 @@ pub fn verify_heartbeat(
|
||||
req.timestamp,
|
||||
&req.entropy_data,
|
||||
&req.stack_state,
|
||||
&salt,
|
||||
&session.salt,
|
||||
);
|
||||
|
||||
// 7. Prepare next mutation order and salt
|
||||
let next_step = pending_step + 1;
|
||||
let next_step = session.pending_mutation_step + 1;
|
||||
let next_mutation = vm_extensions::generate_order(next_step, candidate_state.gene.len());
|
||||
let next_mutation_b64 = vm_extensions::encode_order_b64(&next_mutation);
|
||||
|
||||
@@ -205,28 +172,22 @@ pub fn verify_heartbeat(
|
||||
let next_environment_blob = gene::encode_environment(&candidate_state.environment)
|
||||
.map_err(|e| crate::errors::VerificationError::GeneState(e.to_string()))?;
|
||||
|
||||
conn.execute(
|
||||
"UPDATE sessions SET
|
||||
last_hash=?1,
|
||||
salt=?2,
|
||||
chain_length=chain_length+1,
|
||||
last_seen=?3,
|
||||
gene=?4,
|
||||
environment=?5,
|
||||
pending_mutation=?6,
|
||||
pending_mutation_step=?7
|
||||
WHERE session_id=?8",
|
||||
params![
|
||||
new_hash,
|
||||
next_salt.to_vec(),
|
||||
now,
|
||||
candidate_state.gene,
|
||||
next_environment_blob,
|
||||
next_mutation.program,
|
||||
next_step,
|
||||
req.session_id
|
||||
],
|
||||
)?;
|
||||
let update_record = storage::SessionRecord {
|
||||
session_id: req.session_id.clone(),
|
||||
public_key: session.public_key,
|
||||
salt: next_salt.to_vec(),
|
||||
last_hash: new_hash.clone(),
|
||||
chain_length: session.chain_length + 1,
|
||||
created_at: session.created_at,
|
||||
last_seen: now,
|
||||
expires_at: session.expires_at,
|
||||
gene: candidate_state.gene,
|
||||
environment: next_environment_blob,
|
||||
pending_mutation: next_mutation.program,
|
||||
pending_mutation_step: next_step,
|
||||
};
|
||||
db.update_session(&update_record)
|
||||
.map_err(|e| crate::errors::VerificationError::Storage(e.to_string()))?;
|
||||
|
||||
Ok(HeartbeatVerificationResult {
|
||||
next_salt_hex,
|
||||
@@ -239,6 +200,7 @@ pub fn verify_heartbeat(
|
||||
mod tests {
|
||||
use super::*;
|
||||
use ed25519_dalek::{Signer, SigningKey};
|
||||
use rusqlite::params;
|
||||
use shared::protocol::{EntropyData, Fingerprint, HeartbeatRequest, MouseEvent, StackState};
|
||||
use std::path::Path;
|
||||
|
||||
@@ -310,13 +272,13 @@ mod tests {
|
||||
}
|
||||
|
||||
fn create_test_session(
|
||||
conn: &rusqlite::Connection,
|
||||
db: &storage::DbPool,
|
||||
config: &crate::config::Config,
|
||||
) -> (InitResponse, SigningKey) {
|
||||
let mut rng = rand::thread_rng();
|
||||
let sk = SigningKey::generate(&mut rng);
|
||||
let pk_hex = hex::encode(sk.verifying_key().to_bytes());
|
||||
let init = create_session(conn, config, &pk_hex).unwrap();
|
||||
let init = create_session(db, config, &pk_hex).unwrap();
|
||||
(init, sk)
|
||||
}
|
||||
|
||||
@@ -355,7 +317,11 @@ mod tests {
|
||||
stack_state: stack.clone(),
|
||||
fingerprint: test_fingerprint(),
|
||||
mutation_step: client.pending_mutation_step,
|
||||
gene_commitment: gene::commitment_hex(&candidate_state),
|
||||
gene_commitment: gene::commitment_hex_with_context(
|
||||
&candidate_state,
|
||||
&client.session_id,
|
||||
client.pending_mutation_step,
|
||||
),
|
||||
signature: String::new(),
|
||||
};
|
||||
sign_request(&client.signing_key, &mut req);
|
||||
@@ -382,28 +348,22 @@ mod tests {
|
||||
client.committed_gene_state = candidate_state;
|
||||
}
|
||||
|
||||
fn load_server_gene_state(conn: &rusqlite::Connection, session_id: &str) -> GeneState {
|
||||
let (gene_blob, env_blob): (Vec<u8>, Vec<u8>) = conn
|
||||
.query_row(
|
||||
"SELECT gene, environment FROM sessions WHERE session_id=?1",
|
||||
[session_id],
|
||||
|row| Ok((row.get(0)?, row.get(1)?)),
|
||||
)
|
||||
.unwrap();
|
||||
fn load_server_gene_state(db: &storage::DbPool, session_id: &str) -> GeneState {
|
||||
let session = db.load_session(session_id).unwrap().unwrap();
|
||||
GeneState {
|
||||
gene: gene_blob,
|
||||
environment: gene::decode_environment(&env_blob).unwrap(),
|
||||
gene: session.gene,
|
||||
environment: gene::decode_environment(&session.environment).unwrap(),
|
||||
}
|
||||
}
|
||||
|
||||
fn run_successful_heartbeat(
|
||||
conn: &rusqlite::Connection,
|
||||
db: &storage::DbPool,
|
||||
config: &crate::config::Config,
|
||||
client: &mut SimulatedClient,
|
||||
) -> HeartbeatRequest {
|
||||
let timestamp = storage::current_time_ms();
|
||||
let (req, candidate_state, entropy, stack) = build_request(client, timestamp);
|
||||
let result = verify_heartbeat(conn, config, &req).unwrap();
|
||||
let result = verify_heartbeat(db, config, &req).unwrap();
|
||||
apply_successful_response(client, &req, candidate_state, &entropy, &stack, &result);
|
||||
req
|
||||
}
|
||||
@@ -411,20 +371,19 @@ mod tests {
|
||||
#[test]
|
||||
fn test_session_lifecycle_and_verification() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let config = test_config();
|
||||
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
assert_eq!(init.gene_size, config.gene_size as u32);
|
||||
assert!(!init.mutation_order_b64.is_empty());
|
||||
assert_eq!(init.mutation_step, 1);
|
||||
|
||||
let mut client = client_from_init(&init, signing_key);
|
||||
for _ in 0..5 {
|
||||
run_successful_heartbeat(&conn, &config, &mut client);
|
||||
run_successful_heartbeat(&pool, &config, &mut client);
|
||||
}
|
||||
|
||||
let stats = storage::stats(&conn).unwrap();
|
||||
let stats = pool.stats().unwrap();
|
||||
assert_eq!(stats.sessions, 1);
|
||||
assert_eq!(stats.max_chain_length, 6);
|
||||
}
|
||||
@@ -432,14 +391,13 @@ mod tests {
|
||||
#[test]
|
||||
fn test_deterministic_server_client_parity_across_many_heartbeats() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let config = test_config();
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
let mut client = client_from_init(&init, signing_key);
|
||||
|
||||
for _ in 0..12 {
|
||||
run_successful_heartbeat(&conn, &config, &mut client);
|
||||
let server_state = load_server_gene_state(&conn, &client.session_id);
|
||||
run_successful_heartbeat(&pool, &config, &mut client);
|
||||
let server_state = load_server_gene_state(&pool, &client.session_id);
|
||||
assert_eq!(server_state, client.committed_gene_state);
|
||||
}
|
||||
}
|
||||
@@ -447,14 +405,13 @@ mod tests {
|
||||
#[test]
|
||||
fn test_replay_attack_is_rejected() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let config = test_config();
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
let mut client = client_from_init(&init, signing_key);
|
||||
|
||||
let timestamp = storage::current_time_ms();
|
||||
let (req, candidate_state, entropy, stack) = build_request(&client, timestamp);
|
||||
let result = verify_heartbeat(&conn, &config, &req).unwrap();
|
||||
let result = verify_heartbeat(&pool, &config, &req).unwrap();
|
||||
apply_successful_response(
|
||||
&mut client,
|
||||
&req,
|
||||
@@ -464,7 +421,7 @@ mod tests {
|
||||
&result,
|
||||
);
|
||||
|
||||
let replay = verify_heartbeat(&conn, &config, &req);
|
||||
let replay = verify_heartbeat(&pool, &config, &req);
|
||||
assert!(matches!(
|
||||
replay.unwrap_err(),
|
||||
crate::errors::VerificationError::ChainBroken
|
||||
@@ -474,9 +431,8 @@ mod tests {
|
||||
#[test]
|
||||
fn test_mutation_step_mismatch_is_rejected() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let config = test_config();
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
let client = client_from_init(&init, signing_key);
|
||||
|
||||
let timestamp = storage::current_time_ms();
|
||||
@@ -484,7 +440,7 @@ mod tests {
|
||||
req.mutation_step += 1;
|
||||
sign_request(&client.signing_key, &mut req);
|
||||
|
||||
let err = verify_heartbeat(&conn, &config, &req).unwrap_err();
|
||||
let err = verify_heartbeat(&pool, &config, &req).unwrap_err();
|
||||
assert!(matches!(
|
||||
err,
|
||||
crate::errors::VerificationError::MutationStepMismatch { .. }
|
||||
@@ -494,9 +450,8 @@ mod tests {
|
||||
#[test]
|
||||
fn test_mutation_commitment_tamper_is_rejected() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let config = test_config();
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
let client = client_from_init(&init, signing_key);
|
||||
|
||||
let timestamp = storage::current_time_ms();
|
||||
@@ -504,7 +459,7 @@ mod tests {
|
||||
req.gene_commitment = "00".repeat(32);
|
||||
sign_request(&client.signing_key, &mut req);
|
||||
|
||||
let err = verify_heartbeat(&conn, &config, &req).unwrap_err();
|
||||
let err = verify_heartbeat(&pool, &config, &req).unwrap_err();
|
||||
assert!(matches!(
|
||||
err,
|
||||
crate::errors::VerificationError::MutationCommitmentMismatch
|
||||
@@ -514,9 +469,12 @@ mod tests {
|
||||
#[test]
|
||||
fn test_malformed_server_mutation_program_is_rejected() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let conn = match &pool {
|
||||
storage::DbPool::Sqlite(pool) => pool.get().unwrap(),
|
||||
_ => panic!("expected sqlite pool for test"),
|
||||
};
|
||||
let config = test_config();
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
let client = client_from_init(&init, signing_key);
|
||||
|
||||
conn.execute(
|
||||
@@ -525,9 +483,18 @@ mod tests {
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let updated: Vec<u8> = conn
|
||||
.query_row(
|
||||
"SELECT pending_mutation FROM sessions WHERE session_id=?1",
|
||||
params![client.session_id.clone()],
|
||||
|row| row.get(0),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(updated, vec![0xFFu8]);
|
||||
|
||||
let timestamp = storage::current_time_ms();
|
||||
let (req, _, _, _) = build_request(&client, timestamp);
|
||||
let err = verify_heartbeat(&conn, &config, &req).unwrap_err();
|
||||
let err = verify_heartbeat(&pool, &config, &req).unwrap_err();
|
||||
assert!(matches!(
|
||||
err,
|
||||
crate::errors::VerificationError::MutationProgram(_)
|
||||
@@ -537,9 +504,12 @@ mod tests {
|
||||
#[test]
|
||||
fn test_expired_session_is_rejected() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let conn = match &pool {
|
||||
storage::DbPool::Sqlite(pool) => pool.get().unwrap(),
|
||||
_ => panic!("expected sqlite pool for test"),
|
||||
};
|
||||
let config = test_config();
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
let client = client_from_init(&init, signing_key);
|
||||
|
||||
conn.execute(
|
||||
@@ -550,16 +520,15 @@ mod tests {
|
||||
|
||||
let timestamp = storage::current_time_ms();
|
||||
let (req, _, _, _) = build_request(&client, timestamp);
|
||||
let err = verify_heartbeat(&conn, &config, &req).unwrap_err();
|
||||
let err = verify_heartbeat(&pool, &config, &req).unwrap_err();
|
||||
assert!(matches!(err, crate::errors::VerificationError::Expired));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_create_session_rejects_invalid_public_key_length() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let config = test_config();
|
||||
let err = create_session(&conn, &config, "00ff").unwrap_err();
|
||||
let err = create_session(&pool, &config, "00ff").unwrap_err();
|
||||
assert!(matches!(
|
||||
err,
|
||||
crate::errors::SessionError::InvalidPublicKeyLength
|
||||
@@ -569,19 +538,18 @@ mod tests {
|
||||
#[test]
|
||||
fn test_stale_mutation_step_after_success_is_rejected() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let config = test_config();
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
let mut client = client_from_init(&init, signing_key);
|
||||
|
||||
run_successful_heartbeat(&conn, &config, &mut client);
|
||||
run_successful_heartbeat(&pool, &config, &mut client);
|
||||
|
||||
let timestamp = storage::current_time_ms();
|
||||
let (mut req, _, _, _) = build_request(&client, timestamp);
|
||||
req.mutation_step -= 1;
|
||||
sign_request(&client.signing_key, &mut req);
|
||||
|
||||
let err = verify_heartbeat(&conn, &config, &req).unwrap_err();
|
||||
let err = verify_heartbeat(&pool, &config, &req).unwrap_err();
|
||||
assert!(matches!(
|
||||
err,
|
||||
crate::errors::VerificationError::MutationStepMismatch { .. }
|
||||
@@ -591,15 +559,14 @@ mod tests {
|
||||
#[test]
|
||||
fn test_repeated_simulation_keeps_server_and_client_commitments_equal() {
|
||||
let pool = storage::init_pool(Path::new(":memory:")).unwrap();
|
||||
let conn = pool.get().unwrap();
|
||||
let mut config = test_config();
|
||||
config.gene_size = 128;
|
||||
let (init, signing_key) = create_test_session(&conn, &config);
|
||||
let (init, signing_key) = create_test_session(&pool, &config);
|
||||
let mut client = client_from_init(&init, signing_key);
|
||||
|
||||
for _ in 0..10 {
|
||||
run_successful_heartbeat(&conn, &config, &mut client);
|
||||
let server_state = load_server_gene_state(&conn, &client.session_id);
|
||||
run_successful_heartbeat(&pool, &config, &mut client);
|
||||
let server_state = load_server_gene_state(&pool, &client.session_id);
|
||||
assert_eq!(
|
||||
gene::commitment(&server_state),
|
||||
gene::commitment(&client.committed_gene_state)
|
||||
|
||||
+282
-21
@@ -1,7 +1,9 @@
|
||||
use rusqlite::Connection;
|
||||
use crate::config::Config;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::path::Path;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use valkey::Client as ValkeyClient;
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct StoreStats {
|
||||
@@ -10,9 +12,214 @@ pub struct StoreStats {
|
||||
pub max_chain_length: u64,
|
||||
}
|
||||
|
||||
pub type DbPool = r2d2::Pool<r2d2_sqlite::SqliteConnectionManager>;
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum DbPool {
|
||||
Sqlite(r2d2::Pool<r2d2_sqlite::SqliteConnectionManager>),
|
||||
Valkey(ValkeyStore),
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct ValkeyStore {
|
||||
client: Arc<Mutex<ValkeyClient>>,
|
||||
index_key: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct SessionRecord {
|
||||
pub session_id: String,
|
||||
pub public_key: Vec<u8>,
|
||||
pub salt: Vec<u8>,
|
||||
pub last_hash: Vec<u8>,
|
||||
pub chain_length: u64,
|
||||
pub created_at: u64,
|
||||
pub last_seen: u64,
|
||||
pub expires_at: u64,
|
||||
pub gene: Vec<u8>,
|
||||
pub environment: Vec<u8>,
|
||||
pub pending_mutation: Vec<u8>,
|
||||
pub pending_mutation_step: u64,
|
||||
}
|
||||
|
||||
impl DbPool {
|
||||
pub fn init(config: &Config) -> Result<Self, Box<dyn std::error::Error>> {
|
||||
match config.db_type {
|
||||
crate::config::DbType::SqliteInMemory => {
|
||||
let pool = init_sqlite_pool(Path::new(":memory:"))?;
|
||||
Ok(DbPool::Sqlite(pool))
|
||||
}
|
||||
crate::config::DbType::SqliteInDisk => {
|
||||
let pool = init_sqlite_pool(&config.db_path)?;
|
||||
Ok(DbPool::Sqlite(pool))
|
||||
}
|
||||
crate::config::DbType::Valkey => {
|
||||
let addr = std::env::var("CHRONOSEAL_VALKEY_ADDR").unwrap_or_else(|_| "127.0.0.1:6666".to_string());
|
||||
match ValkeyClient::connect(addr) {
|
||||
Ok(client) => Ok(DbPool::Valkey(ValkeyStore {
|
||||
client: Arc::new(Mutex::new(client)),
|
||||
index_key: "sessions:ids".to_string(),
|
||||
})),
|
||||
Err(err) => {
|
||||
tracing::warn!("valkey connection failed, falling back to sqlite-in-memory: {err}");
|
||||
let pool = init_sqlite_pool(Path::new(":memory:"))?;
|
||||
Ok(DbPool::Sqlite(pool))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn insert_session(&self, record: &SessionRecord) -> Result<(), Box<dyn std::error::Error>> {
|
||||
match self {
|
||||
DbPool::Sqlite(pool) => {
|
||||
let conn = pool.get()?;
|
||||
let mut stmt = conn.prepare(
|
||||
"INSERT INTO sessions (
|
||||
session_id, public_key, salt, last_hash, chain_length,
|
||||
created_at, last_seen, expires_at, gene, environment,
|
||||
pending_mutation, pending_mutation_step
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
|
||||
)?;
|
||||
stmt.execute(rusqlite::params![
|
||||
record.session_id,
|
||||
&record.public_key,
|
||||
&record.salt,
|
||||
&record.last_hash,
|
||||
record.chain_length,
|
||||
record.created_at,
|
||||
record.last_seen,
|
||||
record.expires_at,
|
||||
&record.gene,
|
||||
&record.environment,
|
||||
&record.pending_mutation,
|
||||
record.pending_mutation_step,
|
||||
])?;
|
||||
Ok(())
|
||||
}
|
||||
DbPool::Valkey(store) => store.insert_session(record),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn load_session(&self, session_id: &str) -> Result<Option<SessionRecord>, Box<dyn std::error::Error>> {
|
||||
match self {
|
||||
DbPool::Sqlite(pool) => {
|
||||
let conn = pool.get()?;
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT session_id, public_key, salt, last_hash, chain_length, created_at, last_seen, expires_at, gene, environment, pending_mutation, pending_mutation_step
|
||||
FROM sessions WHERE session_id = ?1",
|
||||
)?;
|
||||
let row = stmt.query_row([session_id], |row| {
|
||||
Ok(SessionRecord {
|
||||
session_id: row.get(0)?,
|
||||
public_key: row.get(1)?,
|
||||
salt: row.get(2)?,
|
||||
last_hash: row.get(3)?,
|
||||
chain_length: row.get(4)?,
|
||||
created_at: row.get(5)?,
|
||||
last_seen: row.get(6)?,
|
||||
expires_at: row.get(7)?,
|
||||
gene: row.get(8)?,
|
||||
environment: row.get(9)?,
|
||||
pending_mutation: row.get(10)?,
|
||||
pending_mutation_step: row.get(11)?,
|
||||
})
|
||||
});
|
||||
match row {
|
||||
Ok(rec) => Ok(Some(rec)),
|
||||
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
|
||||
Err(err) => Err(Box::new(err)),
|
||||
}
|
||||
}
|
||||
DbPool::Valkey(store) => store.load_session(session_id),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn update_session(&self, record: &SessionRecord) -> Result<(), Box<dyn std::error::Error>> {
|
||||
match self {
|
||||
DbPool::Sqlite(pool) => {
|
||||
let conn = pool.get()?;
|
||||
conn.execute(
|
||||
"UPDATE sessions SET
|
||||
public_key=?1,
|
||||
salt=?2,
|
||||
last_hash=?3,
|
||||
chain_length=?4,
|
||||
created_at=?5,
|
||||
last_seen=?6,
|
||||
expires_at=?7,
|
||||
gene=?8,
|
||||
environment=?9,
|
||||
pending_mutation=?10,
|
||||
pending_mutation_step=?11
|
||||
WHERE session_id=?12",
|
||||
rusqlite::params![
|
||||
&record.public_key,
|
||||
&record.salt,
|
||||
&record.last_hash,
|
||||
record.chain_length,
|
||||
record.created_at,
|
||||
record.last_seen,
|
||||
record.expires_at,
|
||||
&record.gene,
|
||||
&record.environment,
|
||||
&record.pending_mutation,
|
||||
record.pending_mutation_step,
|
||||
&record.session_id,
|
||||
],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
DbPool::Valkey(store) => store.insert_session(record),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn delete_expired_sessions(&self) -> Result<(), Box<dyn std::error::Error>> {
|
||||
match self {
|
||||
DbPool::Sqlite(pool) => {
|
||||
let conn = pool.get()?;
|
||||
conn.execute(
|
||||
"DELETE FROM sessions WHERE expires_at < ?1",
|
||||
rusqlite::params![current_time_ms()],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
DbPool::Valkey(store) => store.purge_expired_sessions(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn stats(&self) -> Result<StoreStats, Box<dyn std::error::Error>> {
|
||||
match self {
|
||||
DbPool::Sqlite(pool) => {
|
||||
let conn = pool.get()?;
|
||||
let now = current_time_ms();
|
||||
let sessions = conn.query_row("SELECT COUNT(*) FROM sessions", [], |row| row.get(0))?;
|
||||
let expired_sessions = conn.query_row(
|
||||
"SELECT COUNT(*) FROM sessions WHERE expires_at < ?1",
|
||||
[now],
|
||||
|row| row.get(0),
|
||||
)?;
|
||||
let max_chain_length = conn.query_row(
|
||||
"SELECT COALESCE(MAX(chain_length), 0) FROM sessions",
|
||||
[],
|
||||
|row| row.get(0),
|
||||
)?;
|
||||
Ok(StoreStats {
|
||||
sessions,
|
||||
expired_sessions,
|
||||
max_chain_length,
|
||||
})
|
||||
}
|
||||
DbPool::Valkey(store) => store.stats(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub fn init_pool(path: &Path) -> Result<DbPool, Box<dyn std::error::Error>> {
|
||||
let pool = init_sqlite_pool(path)?;
|
||||
Ok(DbPool::Sqlite(pool))
|
||||
}
|
||||
|
||||
fn init_sqlite_pool(path: &Path) -> Result<r2d2::Pool<r2d2_sqlite::SqliteConnectionManager>, Box<dyn std::error::Error>> {
|
||||
let manager = if path == Path::new(":memory:") {
|
||||
r2d2_sqlite::SqliteConnectionManager::memory()
|
||||
} else {
|
||||
@@ -21,7 +228,6 @@ pub fn init_pool(path: &Path) -> Result<DbPool, Box<dyn std::error::Error>> {
|
||||
}
|
||||
r2d2_sqlite::SqliteConnectionManager::file(path)
|
||||
};
|
||||
|
||||
let pool = r2d2::Pool::new(manager)?;
|
||||
let conn = pool.get()?;
|
||||
init_schema(&conn)?;
|
||||
@@ -89,24 +295,79 @@ fn ensure_column(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn stats(conn: &Connection) -> Result<StoreStats, rusqlite::Error> {
|
||||
let now = current_time_ms();
|
||||
let sessions = conn.query_row("SELECT COUNT(*) FROM sessions", [], |row| row.get(0))?;
|
||||
let expired_sessions = conn.query_row(
|
||||
"SELECT COUNT(*) FROM sessions WHERE expires_at < ?1",
|
||||
[now],
|
||||
|row| row.get(0),
|
||||
)?;
|
||||
let max_chain_length = conn.query_row(
|
||||
"SELECT COALESCE(MAX(chain_length), 0) FROM sessions",
|
||||
[],
|
||||
|row| row.get(0),
|
||||
)?;
|
||||
Ok(StoreStats {
|
||||
sessions,
|
||||
expired_sessions,
|
||||
max_chain_length,
|
||||
})
|
||||
impl ValkeyStore {
|
||||
fn session_key(&self, session_id: &str) -> String {
|
||||
format!("session:{}", session_id)
|
||||
}
|
||||
|
||||
fn load_session(&self, session_id: &str) -> Result<Option<SessionRecord>, Box<dyn std::error::Error>> {
|
||||
let mut client = self.client.lock().unwrap();
|
||||
if let Some(payload) = client.get(&self.session_key(session_id))? {
|
||||
let record = serde_json::from_str(&payload)?;
|
||||
Ok(Some(record))
|
||||
} else {
|
||||
Ok(None)
|
||||
}
|
||||
}
|
||||
|
||||
fn insert_session(&self, record: &SessionRecord) -> Result<(), Box<dyn std::error::Error>> {
|
||||
let mut client = self.client.lock().unwrap();
|
||||
let value = serde_json::to_string(record)?;
|
||||
client.set(&self.session_key(&record.session_id), &value)?;
|
||||
let existing = client.get(&self.index_key)?;
|
||||
let mut ids = existing.unwrap_or_default();
|
||||
if !ids.split('\n').any(|id| id == record.session_id) {
|
||||
if !ids.is_empty() {
|
||||
ids.push('\n');
|
||||
}
|
||||
ids.push_str(&record.session_id);
|
||||
client.set(&self.index_key, &ids)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn purge_expired_sessions(&self) -> Result<(), Box<dyn std::error::Error>> {
|
||||
let mut client = self.client.lock().unwrap();
|
||||
let ids = client.get(&self.index_key)?.unwrap_or_default();
|
||||
let now = current_time_ms();
|
||||
let mut remaining: Vec<String> = Vec::new();
|
||||
for id in ids.split('\n').filter(|id| !id.is_empty()) {
|
||||
if let Some(payload) = client.get(&self.session_key(id))? {
|
||||
if let Ok(record) = serde_json::from_str::<SessionRecord>(&payload) {
|
||||
if record.expires_at > now {
|
||||
remaining.push(id.to_string());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
client.set(&self.index_key, &remaining.join("\n"))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn stats(&self) -> Result<StoreStats, Box<dyn std::error::Error>> {
|
||||
let mut client = self.client.lock().unwrap();
|
||||
let ids = client.get(&self.index_key)?.unwrap_or_default();
|
||||
let now = current_time_ms();
|
||||
let mut sessions = 0;
|
||||
let mut expired_sessions = 0;
|
||||
let mut max_chain_length = 0;
|
||||
for id in ids.split('\n').filter(|id| !id.is_empty()) {
|
||||
if let Some(payload) = client.get(&self.session_key(id))? {
|
||||
if let Ok(record) = serde_json::from_str::<SessionRecord>(&payload) {
|
||||
sessions += 1;
|
||||
if record.expires_at < now {
|
||||
expired_sessions += 1;
|
||||
}
|
||||
max_chain_length = max_chain_length.max(record.chain_length);
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(StoreStats {
|
||||
sessions,
|
||||
expired_sessions,
|
||||
max_chain_length,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
pub fn current_time_ms() -> u64 {
|
||||
|
||||
Reference in new issue
Block a user