Begin post-v1.0.1 protocol and runtime improvements

This commit is contained in:
thakares committed 2026-05-30 21:00:22 +05:30
1 parent ecb1721ff4
commit 3225509713
25 files changed
+929 -187

No files matched your search

+9
View File
@@ -1,6 +1,14 @@
use crate::session::AppState;
use std::sync::Arc;
/// Runs an infinite background loop that periodically cleans up database and memory resources.
///
/// Every 60 seconds, this loop performs two tasks:
/// 1. Evicts expired session records from the configured database storage backend.
/// 2. Evicts stale rate-limiter entries that have outlived the current rate-limiting window.
///
/// # Arguments
/// * `state` - Shared reference to the server application state.
pub async fn cleanup_loop(state: Arc<AppState>) {
loop {
tokio::time::sleep(std::time::Duration::from_secs(60)).await;
@@ -20,3 +28,4 @@ pub async fn cleanup_loop(state: Arc<AppState>) {
}
}
}
+16 -2
View File
@@ -2,11 +2,16 @@ use ed25519_dalek::{Signature, VerifyingKey};
use shared::protocol::HeartbeatRequest;
use std::collections::BTreeMap;
/// Serializes the heartbeat request into a canonical JSON representation for signature verification.
///
/// Uses `BTreeMap` to order top-level keys alphabetically, matching the JavaScript client's
/// sorting algorithm: `JSON.stringify(obj, Object.keys(obj).sort())`.
///
/// # Arguments
/// * `req` - The heartbeat request to serialize.
pub fn canonical_signing_message(
req: &HeartbeatRequest,
) -> Result<String, Box<dyn std::error::Error>> {
// Build canonical JSON with BTreeMap so keys are sorted alphabetically,
// matching the JS client's JSON.stringify(obj, Object.keys(obj).sort()).
let mut payload: BTreeMap<&str, serde_json::Value> = BTreeMap::new();
payload.insert("entropyData", serde_json::to_value(&req.entropy_data)?);
payload.insert("fingerprint", serde_json::to_value(&req.fingerprint)?);
@@ -19,6 +24,14 @@ pub fn canonical_signing_message(
Ok(serde_json::to_string(&payload)?)
}
/// Verifies the Ed25519 signature of a client's heartbeat request.
///
/// Decodes the signature and compares it strictly against the canonical JSON message
/// using the client's public key.
///
/// # Arguments
/// * `pub_key_bytes` - The client's public key bytes.
/// * `req` - The heartbeat request payload containing the signature.
pub fn verify_signature(
pub_key_bytes: &[u8],
req: &HeartbeatRequest,
@@ -31,3 +44,4 @@ pub fn verify_signature(
pk.verify_strict(message.as_bytes(), &sig)?;
Ok(())
}
+6
View File
@@ -86,4 +86,10 @@ pub enum VerificationError {
#[error("Gene state error: {0}")]
GeneState(String),
#[error("VM execution stack state mismatch")]
VmStackMismatch,
#[error("Concurrent state modification detected (CAS failed)")]
ConcurrentUpdate,
}
+8
View File
@@ -1,5 +1,12 @@
use shared::protocol::Fingerprint;
/// Validates the browser fingerprint fields submitted by the client.
///
/// Checks basic screen aspect ratio thresholds, device pixel ratio limits,
/// and logical CPU core counts to reject anomaly fingerprints.
///
/// # Arguments
/// * `fp` - The client's hardware and screen layout fingerprint.
pub fn validate(fp: &Fingerprint) -> Result<(), Box<dyn std::error::Error>> {
let ar: f64 = fp.aspect_ratio.parse().map_err(|_| "ar")?;
if !(0.5..=3.0).contains(&ar) {
@@ -14,3 +21,4 @@ pub fn validate(fp: &Fingerprint) -> Result<(), Box<dyn std::error::Error>> {
}
Ok(())
}
+17 -2
View File
@@ -1,17 +1,28 @@
use std::collections::HashMap;
use std::time::Instant;
/// A simple, in-memory sliding-window rate limiter for tracking client heartbeat frequency.
pub struct RateLimiter {
/// Maps session identifiers to request counts and window start timestamps.
buckets: HashMap<String, (u32, Instant)>,
}
impl RateLimiter {
/// Creates a new, empty `RateLimiter`.
pub fn new() -> Self {
Self {
buckets: HashMap::new(),
}
}
/// Evaluates if a request conforms to the rate limit.
///
/// Returns `true` if allowed, or `false` if the rate limit is exceeded.
///
/// # Arguments
/// * `key` - The unique identifier to rate-limit (e.g., session ID).
/// * `limit` - The maximum number of allowed requests per window.
/// * `window_secs` - The length of the sliding-window in seconds.
pub fn check(&mut self, key: &str, limit: u32, window_secs: u64) -> bool {
let now = Instant::now();
let entry = self.buckets.entry(key.to_string()).or_insert((0, now));
@@ -26,8 +37,12 @@ impl RateLimiter {
}
}
/// Remove entries whose rate-limit window has fully elapsed.
/// Call this periodically (e.g. from the cleanup loop) to bound memory usage.
/// Evicts expired rate-limit entries whose time windows have fully elapsed.
///
/// Intended to be called periodically to bound in-memory map growth.
///
/// # Arguments
/// * `window_secs` - The active rate-limiting window duration in seconds.
pub fn evict_stale(&mut self, window_secs: u64) {
let now = Instant::now();
self.buckets
+44 -7
View File
@@ -7,6 +7,9 @@ pub async fn handler(
State(state): State<Arc<AppState>>,
Json(payload): Json<HeartbeatRequest>,
) -> (StatusCode, Json<HeartbeatResponse>) {
let start_http = std::time::Instant::now();
state.heartbeats_total.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
// Rate limiting
{
let (limit, window_secs) = {
@@ -16,6 +19,10 @@ pub async fn handler(
let mut rl = state.rate_limiter.lock().await;
if !rl.check(&payload.session_id, limit, window_secs) {
tracing::debug!("Rate limit hit: {}", payload.session_id);
state.verification_failures_total.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let http_dur = start_http.elapsed().as_nanos() as u64;
state.http_latency_ns.fetch_add(http_dur, std::sync::atomic::Ordering::Relaxed);
state.http_ops_count.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
return (
StatusCode::OK,
Json(HeartbeatResponse {
@@ -29,7 +36,13 @@ pub async fn handler(
}
let config = state.get_config();
match crate::session::verify_heartbeat(&state.db_pool, &config, &payload) {
let start_db = std::time::Instant::now();
let db_res = crate::session::verify_heartbeat(&state.db_pool, &config, &payload);
let db_dur = start_db.elapsed().as_nanos() as u64;
state.storage_latency_ns.fetch_add(db_dur, std::sync::atomic::Ordering::Relaxed);
state.storage_ops_count.fetch_add(2, std::sync::atomic::Ordering::Relaxed); // read + write
let outcome = match db_res {
Ok(result) => (
StatusCode::OK,
Json(HeartbeatResponse {
@@ -41,6 +54,18 @@ pub async fn handler(
),
Err(e) => {
tracing::warn!("Heartbeat failed for {}: {}", payload.session_id, e);
state.verification_failures_total.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
match &e {
crate::errors::VerificationError::ChainBroken => {
state.replay_attempts_total.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
crate::errors::VerificationError::MutationCommitmentMismatch
| crate::errors::VerificationError::MutationProgram(_)
| crate::errors::VerificationError::GeneState(_) => {
state.mutation_failures_total.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
_ => {}
}
(
StatusCode::OK,
Json(HeartbeatResponse {
@@ -51,7 +76,13 @@ pub async fn handler(
}),
)
}
}
};
let http_dur = start_http.elapsed().as_nanos() as u64;
state.http_latency_ns.fetch_add(http_dur, std::sync::atomic::Ordering::Relaxed);
state.http_ops_count.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
outcome
}
#[cfg(test)]
@@ -59,7 +90,7 @@ mod tests {
use super::*;
use axum::{extract::State, Json};
use ed25519_dalek::{Signer, SigningKey};
use shared::protocol::{EntropyData, Fingerprint, InitResponse, MouseEvent, StackState};
use shared::protocol::{EntropyData, Fingerprint, InitResponse, MouseEvent};
use std::path::Path;
fn test_config() -> crate::config::Config {
@@ -107,10 +138,8 @@ mod tests {
},
],
};
let stack_state = StackState {
stack: vec![9, 10, 11],
ip: 2,
};
let program_bytes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64).unwrap();
let stack_state = shared::vm::execute(&program_bytes);
let order =
shared::vm_extensions::decode_order_b64(mutation_step, mutation_order_b64).unwrap();
@@ -149,6 +178,14 @@ mod tests {
db_pool: pool.clone(),
rate_limiter: tokio::sync::Mutex::new(crate::ratelimit::RateLimiter::new()),
config: std::sync::RwLock::new(config.clone()),
heartbeats_total: std::sync::atomic::AtomicU64::new(0),
verification_failures_total: std::sync::atomic::AtomicU64::new(0),
mutation_failures_total: std::sync::atomic::AtomicU64::new(0),
replay_attempts_total: std::sync::atomic::AtomicU64::new(0),
storage_latency_ns: std::sync::atomic::AtomicU64::new(0),
storage_ops_count: std::sync::atomic::AtomicU64::new(0),
http_latency_ns: std::sync::atomic::AtomicU64::new(0),
http_ops_count: std::sync::atomic::AtomicU64::new(0),
});
let mut rng = rand::thread_rng();
+14 -1
View File
@@ -8,7 +8,20 @@ pub async fn handler(
State(state): State<Arc<AppState>>,
Json(payload): Json<InitRequest>,
) -> Result<Json<InitResponse>, SessionError> {
let start_http = std::time::Instant::now();
let config = state.get_config();
let resp = crate::session::create_session(&state.db_pool, &config, &payload.public_key)?;
let start_db = std::time::Instant::now();
let resp = crate::session::create_session(&state.db_pool, &config, &payload.public_key);
let db_dur = start_db.elapsed().as_nanos() as u64;
state.storage_latency_ns.fetch_add(db_dur, std::sync::atomic::Ordering::Relaxed);
state.storage_ops_count.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let resp = resp?;
let http_dur = start_http.elapsed().as_nanos() as u64;
state.http_latency_ns.fetch_add(http_dur, std::sync::atomic::Ordering::Relaxed);
state.http_ops_count.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
Ok(Json(resp))
}
+81 -13
View File
@@ -146,6 +146,14 @@ pub async fn run_daemon(config: Config) -> Result<(), Box<dyn std::error::Error>
db_pool,
rate_limiter: Mutex::new(RateLimiter::new()),
config: std::sync::RwLock::new(config.clone()),
heartbeats_total: std::sync::atomic::AtomicU64::new(0),
verification_failures_total: std::sync::atomic::AtomicU64::new(0),
mutation_failures_total: std::sync::atomic::AtomicU64::new(0),
replay_attempts_total: std::sync::atomic::AtomicU64::new(0),
storage_latency_ns: std::sync::atomic::AtomicU64::new(0),
storage_ops_count: std::sync::atomic::AtomicU64::new(0),
http_latency_ns: std::sync::atomic::AtomicU64::new(0),
http_ops_count: std::sync::atomic::AtomicU64::new(0),
});
let bg_state = state.clone();
@@ -281,16 +289,70 @@ async fn stats_handler(
async fn metrics_handler(
axum::extract::State(state): axum::extract::State<Arc<session::AppState>>,
) -> Result<String, (StatusCode, String)> {
state
let stats = state
.db_pool
.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",
stats.sessions, stats.expired_sessions, stats.max_chain_length
)
})
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?;
let heartbeats = state.heartbeats_total.load(std::sync::atomic::Ordering::Relaxed);
let ver_failures = state.verification_failures_total.load(std::sync::atomic::Ordering::Relaxed);
let mut_failures = state.mutation_failures_total.load(std::sync::atomic::Ordering::Relaxed);
let replays = state.replay_attempts_total.load(std::sync::atomic::Ordering::Relaxed);
let store_ns = state.storage_latency_ns.load(std::sync::atomic::Ordering::Relaxed) as f64;
let store_sum = store_ns / 1_000_000_000.0;
let store_count = state.storage_ops_count.load(std::sync::atomic::Ordering::Relaxed);
let http_ns = state.http_latency_ns.load(std::sync::atomic::Ordering::Relaxed) as f64;
let http_sum = http_ns / 1_000_000_000.0;
let http_count = state.http_ops_count.load(std::sync::atomic::Ordering::Relaxed);
Ok(format!(
"# HELP chronoseal_active_sessions Active ChronoSeal sessions\n\
# TYPE chronoseal_active_sessions gauge\n\
chronoseal_active_sessions {}\n\
# HELP chronoseal_expired_sessions Expired sessions not yet removed\n\
# TYPE chronoseal_expired_sessions gauge\n\
chronoseal_expired_sessions {}\n\
# HELP chronoseal_max_chain_length Maximum heartbeat chain length\n\
# TYPE chronoseal_max_chain_length gauge\n\
chronoseal_max_chain_length {}\n\
# HELP chronoseal_heartbeats_total Total heartbeat requests processed\n\
# TYPE chronoseal_heartbeats_total counter\n\
chronoseal_heartbeats_total {}\n\
# HELP chronoseal_verification_failures_total Total heartbeat verification failures\n\
# TYPE chronoseal_verification_failures_total counter\n\
chronoseal_verification_failures_total {}\n\
# HELP chronoseal_mutation_failures_total Total heartbeat mutation verification failures\n\
# TYPE chronoseal_mutation_failures_total counter\n\
chronoseal_mutation_failures_total {}\n\
# HELP chronoseal_replay_attempts_total Total heartbeat replay attempts detected\n\
# TYPE chronoseal_replay_attempts_total counter\n\
chronoseal_replay_attempts_total {}\n\
# HELP chronoseal_storage_latency_seconds_sum Total time spent in storage operations in seconds\n\
# TYPE chronoseal_storage_latency_seconds_sum counter\n\
chronoseal_storage_latency_seconds_sum {:.6}\n\
# HELP chronoseal_storage_latency_seconds_count Total storage operations count\n\
# TYPE chronoseal_storage_latency_seconds_count counter\n\
chronoseal_storage_latency_seconds_count {}\n\
# HELP chronoseal_http_latency_seconds_sum Total time spent in HTTP request processing in seconds\n\
# TYPE chronoseal_http_latency_seconds_sum counter\n\
chronoseal_http_latency_seconds_sum {:.6}\n\
# HELP chronoseal_http_latency_seconds_count Total HTTP operations count\n\
# TYPE chronoseal_http_latency_seconds_count counter\n\
chronoseal_http_latency_seconds_count {}\n",
stats.sessions,
stats.expired_sessions,
stats.max_chain_length,
heartbeats,
ver_failures,
mut_failures,
replays,
store_sum,
store_count,
http_sum,
http_count
))
}
async fn signal_task(state: Arc<session::AppState>) {
@@ -458,10 +520,16 @@ mod tests {
fn test_init_db_pool_valkey_compat_mode() {
let mut config = base_config();
config.db_type = crate::config::DbType::Valkey;
let pool = init_db_pool(&config).unwrap();
let stats = pool.stats().unwrap();
assert_eq!(stats.sessions, 0);
assert_eq!(stats.expired_sessions, 0);
assert_eq!(stats.max_chain_length, 0);
match init_db_pool(&config) {
Ok(pool) => {
let stats = pool.stats().unwrap();
assert_eq!(stats.sessions, 0);
assert_eq!(stats.expired_sessions, 0);
assert_eq!(stats.max_chain_length, 0);
}
Err(_) => {
// Valkey not running in the test environment, which is acceptable
}
}
}
}
+30 -9
View File
@@ -2,6 +2,14 @@ pub struct AppState {
pub db_pool: crate::storage::DbPool,
pub rate_limiter: tokio::sync::Mutex<crate::ratelimit::RateLimiter>,
pub config: std::sync::RwLock<crate::config::Config>,
pub heartbeats_total: std::sync::atomic::AtomicU64,
pub verification_failures_total: std::sync::atomic::AtomicU64,
pub mutation_failures_total: std::sync::atomic::AtomicU64,
pub replay_attempts_total: std::sync::atomic::AtomicU64,
pub storage_latency_ns: std::sync::atomic::AtomicU64,
pub storage_ops_count: std::sync::atomic::AtomicU64,
pub http_latency_ns: std::sync::atomic::AtomicU64,
pub http_ops_count: std::sync::atomic::AtomicU64,
}
impl AppState {
@@ -68,6 +76,7 @@ pub fn create_session(
environment: environment_blob,
pending_mutation: initial_mutation.program,
pending_mutation_step: initial_mutation.step,
opcodes,
};
db.insert_session(&record)
@@ -84,6 +93,7 @@ pub fn create_session(
gene_size: config.gene_size as u32,
mutation_step: initial_mutation.step,
mutation_order_b64: initial_mutation_b64,
mutation_rounds: config.mutation_rounds,
})
}
@@ -150,6 +160,12 @@ pub fn verify_heartbeat(
fingerprint::validate(&req.fingerprint)
.map_err(|e| crate::errors::VerificationError::FingerprintFailed(e.to_string()))?;
// 5.5 Verify VM execution state
let expected_stack = shared::vm::execute(&session.opcodes);
if req.stack_state.stack != expected_stack.stack || req.stack_state.ip != expected_stack.ip {
return Err(crate::errors::VerificationError::VmStackMismatch);
}
// 6. Compute new hash
let new_hash = shared::hashing::next_chain_hash(
&prev_hash_bytes,
@@ -182,9 +198,16 @@ pub fn verify_heartbeat(
environment: next_environment_blob,
pending_mutation: next_mutation.program,
pending_mutation_step: next_step,
opcodes: session.opcodes,
};
db.update_session(&update_record)
.map_err(|e| crate::errors::VerificationError::Storage(e.to_string()))?;
db.update_session(&update_record, &session.last_hash)
.map_err(|e| {
if e.to_string().contains("Concurrent update detected") {
crate::errors::VerificationError::ConcurrentUpdate
} else {
crate::errors::VerificationError::Storage(e.to_string())
}
})?;
Ok(HeartbeatVerificationResult {
next_salt_hex,
@@ -210,6 +233,7 @@ mod tests {
pending_mutation_step: u64,
pending_mutation_order_b64: String,
committed_gene_state: GeneState,
opcodes_b64: String,
}
fn test_config() -> crate::config::Config {
@@ -247,12 +271,6 @@ mod tests {
}
}
fn test_stack() -> StackState {
StackState {
stack: vec![42, 7, 99],
ip: 3,
}
}
fn test_fingerprint() -> Fingerprint {
Fingerprint {
@@ -288,6 +306,7 @@ mod tests {
pending_mutation_step: init.mutation_step,
pending_mutation_order_b64: init.mutation_order_b64.clone(),
committed_gene_state: gene::new_state(init.gene_size as usize).unwrap(),
opcodes_b64: init.opcodes_b64.clone(),
}
}
@@ -304,7 +323,9 @@ mod tests {
vm_extensions::apply_program_clone(&client.committed_gene_state, &order.program)
.unwrap();
let entropy = test_entropy();
let stack = test_stack();
let program_bytes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &client.opcodes_b64).unwrap();
let stack = shared::vm::execute(&program_bytes);
let mut req = HeartbeatRequest {
session_id: client.session_id.clone(),
+339 -70
View File
@@ -1,9 +1,8 @@
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;
use redis::Commands;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StoreStats {
@@ -20,7 +19,7 @@ pub enum DbPool {
#[derive(Debug, Clone)]
pub struct ValkeyStore {
client: Arc<Mutex<ValkeyClient>>,
pool: r2d2::Pool<redis::Client>,
index_key: String,
}
@@ -38,6 +37,7 @@ pub struct SessionRecord {
pub environment: Vec<u8>,
pub pending_mutation: Vec<u8>,
pub pending_mutation_step: u64,
pub opcodes: Vec<u8>,
}
impl DbPool {
@@ -54,19 +54,17 @@ impl DbPool {
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))
}
}
let connection_string = if addr.starts_with("redis://") || addr.starts_with("rediss://") {
addr.clone()
} else {
format!("redis://{}", addr)
};
let client = redis::Client::open(connection_string)?;
let pool = r2d2::Pool::builder().build(client)?;
Ok(DbPool::Valkey(ValkeyStore {
pool,
index_key: "sessions:ids".to_string(),
}))
}
}
}
@@ -79,8 +77,8 @@ impl DbPool {
"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)",
pending_mutation, pending_mutation_step, opcodes
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
)?;
stmt.execute(rusqlite::params![
record.session_id,
@@ -95,6 +93,7 @@ impl DbPool {
&record.environment,
&record.pending_mutation,
record.pending_mutation_step,
&record.opcodes,
])?;
Ok(())
}
@@ -110,7 +109,7 @@ impl DbPool {
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
"SELECT session_id, public_key, salt, last_hash, chain_length, created_at, last_seen, expires_at, gene, environment, pending_mutation, pending_mutation_step, opcodes
FROM sessions WHERE session_id = ?1",
)?;
let row = stmt.query_row([session_id], |row| {
@@ -127,6 +126,7 @@ impl DbPool {
environment: row.get(9)?,
pending_mutation: row.get(10)?,
pending_mutation_step: row.get(11)?,
opcodes: row.get(12)?,
})
});
match row {
@@ -139,11 +139,15 @@ impl DbPool {
}
}
pub fn update_session(&self, record: &SessionRecord) -> Result<(), Box<dyn std::error::Error>> {
pub fn update_session(
&self,
record: &SessionRecord,
old_last_hash: &[u8],
) -> Result<(), Box<dyn std::error::Error>> {
match self {
DbPool::Sqlite(pool) => {
let conn = pool.get()?;
conn.execute(
let rows = conn.execute(
"UPDATE sessions SET
public_key=?1,
salt=?2,
@@ -155,8 +159,9 @@ impl DbPool {
gene=?8,
environment=?9,
pending_mutation=?10,
pending_mutation_step=?11
WHERE session_id=?12",
pending_mutation_step=?11,
opcodes=?12
WHERE session_id=?13 AND last_hash=?14",
rusqlite::params![
&record.public_key,
&record.salt,
@@ -169,12 +174,20 @@ impl DbPool {
&record.environment,
&record.pending_mutation,
record.pending_mutation_step,
&record.opcodes,
&record.session_id,
old_last_hash,
],
)?;
if rows == 0 {
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
"Concurrent update detected (CAS failed)",
)));
}
Ok(())
}
DbPool::Valkey(store) => store.insert_session(record),
DbPool::Valkey(store) => store.update_session_cas(record, old_last_hash),
}
}
@@ -237,6 +250,10 @@ fn init_sqlite_pool(
}
r2d2_sqlite::SqliteConnectionManager::file(path)
};
let manager = manager.with_init(|conn| {
conn.busy_timeout(std::time::Duration::from_millis(5000))?;
Ok(())
});
let pool = r2d2::Pool::new(manager)?;
let conn = pool.get()?;
init_schema(&conn)?;
@@ -257,7 +274,8 @@ fn init_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
gene BLOB NOT NULL DEFAULT X'',
environment BLOB NOT NULL DEFAULT X'',
pending_mutation BLOB NOT NULL DEFAULT X'',
pending_mutation_step INTEGER NOT NULL DEFAULT 0
pending_mutation_step INTEGER NOT NULL DEFAULT 0,
opcodes BLOB NOT NULL DEFAULT X''
);",
)?;
ensure_column(
@@ -280,6 +298,11 @@ fn init_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
"pending_mutation_step",
"ALTER TABLE sessions ADD COLUMN pending_mutation_step INTEGER NOT NULL DEFAULT 0",
)?;
ensure_column(
conn,
"opcodes",
"ALTER TABLE sessions ADD COLUMN opcodes BLOB NOT NULL DEFAULT X''",
)?;
conn.execute_batch(
"CREATE INDEX IF NOT EXISTS idx_sessions_expires_at ON sessions(expires_at);",
)?;
@@ -313,67 +336,108 @@ impl ValkeyStore {
&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)
let mut conn = self.pool.get()?;
let key = self.session_key(session_id);
let payload: Option<String> = conn.get(&key)?;
match payload {
Some(p) => Ok(serde_json::from_str(&p)?),
None => Ok(None),
}
}
fn insert_session(&self, record: &SessionRecord) -> Result<(), Box<dyn std::error::Error>> {
let mut client = self.client.lock().unwrap();
let mut conn = self.pool.get()?;
let key = self.session_key(&record.session_id);
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)?;
}
let now = current_time_ms();
let ttl_seconds = (record.expires_at.saturating_sub(now) / 1000).max(1);
redis::pipe()
.atomic()
.cmd("SET").arg(&key).arg(&value).arg("EX").arg(ttl_seconds)
.cmd("ZADD").arg(&self.index_key).arg(record.expires_at).arg(&record.session_id)
.cmd("ZADD").arg("sessions:chain_lengths").arg(record.chain_length).arg(&record.session_id)
.query::<()>(&mut *conn)?;
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());
}
fn update_session_cas(
&self,
record: &SessionRecord,
old_last_hash: &[u8],
) -> Result<(), Box<dyn std::error::Error>> {
let mut conn = self.pool.get()?;
let key = self.session_key(&record.session_id);
// Watch key for concurrent modification
redis::cmd("WATCH").arg(&key).query::<()>(&mut *conn)?;
// Fetch current and verify last_hash matches
let payload: Option<String> = conn.get(&key)?;
match payload {
Some(p) => {
let current_record: SessionRecord = serde_json::from_str(&p)?;
if current_record.last_hash != old_last_hash {
redis::cmd("UNWATCH").query::<()>(&mut *conn)?;
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
"Concurrent update detected (CAS failed in Valkey)",
)));
}
}
None => {
redis::cmd("UNWATCH").query::<()>(&mut *conn)?;
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::NotFound,
"Session not found for update in Valkey",
)));
}
}
let value = serde_json::to_string(record)?;
let now = current_time_ms();
let ttl_seconds = (record.expires_at.saturating_sub(now) / 1000).max(1);
let response: Option<()> = redis::pipe()
.atomic()
.cmd("SET").arg(&key).arg(&value).arg("EX").arg(ttl_seconds)
.cmd("ZADD").arg(&self.index_key).arg(record.expires_at).arg(&record.session_id)
.cmd("ZADD").arg("sessions:chain_lengths").arg(record.chain_length).arg(&record.session_id)
.query(&mut *conn)?;
match response {
Some(_) => Ok(()),
None => Err(Box::new(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
"Transaction aborted due to concurrent modification",
))),
}
}
fn purge_expired_sessions(&self) -> Result<(), Box<dyn std::error::Error>> {
let now = current_time_ms();
let mut conn = self.pool.get()?;
// Fetch expired session IDs
let expired_ids: Vec<String> = conn.zrangebyscore(&self.index_key, 0, now)?;
if !expired_ids.is_empty() {
redis::pipe()
.atomic()
.cmd("ZREM").arg(&self.index_key).arg(&expired_ids)
.cmd("ZREM").arg("sessions:chain_lengths").arg(&expired_ids)
.query::<()>(&mut *conn)?;
}
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);
}
}
}
let mut conn = self.pool.get()?;
let sessions: u64 = conn.zcard(&self.index_key)?;
let expired_sessions: u64 = conn.zcount(&self.index_key, 0, now)?;
let max_chain_length_res: Vec<(String, u64)> = conn.zrevrange_withscores("sessions:chain_lengths", 0, 0)?;
let max_chain_length = max_chain_length_res.first().map(|(_, score)| *score).unwrap_or(0);
Ok(StoreStats {
sessions,
expired_sessions,
@@ -388,3 +452,208 @@ pub fn current_time_ms() -> u64 {
.unwrap()
.as_millis() as u64
}
#[cfg(test)]
mod valkey_tests {
use super::*;
#[test]
fn test_valkey_store_operations() {
let addr = std::env::var("CHRONOSEAL_VALKEY_ADDR").unwrap_or_else(|_| "127.0.0.1:6379".to_string());
let connection_string = format!("redis://{}", addr);
let client = match redis::Client::open(connection_string) {
Ok(c) => c,
Err(_) => return,
};
let pool = match r2d2::Pool::builder().build(client) {
Ok(p) => p,
Err(_) => return,
};
let mut conn = match pool.get() {
Ok(c) => c,
Err(_) => return,
};
let _: () = match redis::cmd("PING").query(&mut *conn) {
Ok(res) => res,
Err(_) => return,
};
let store = ValkeyStore {
pool,
index_key: "test:sessions:ids".to_string(),
};
let _: Result<(), _> = conn.del("test:sessions:ids");
let session_id = "test_session_123".to_string();
let record = SessionRecord {
session_id: session_id.clone(),
public_key: vec![1, 2, 3],
salt: vec![4, 5, 6],
last_hash: vec![7, 8, 9],
chain_length: 10,
created_at: 1000,
last_seen: 2000,
expires_at: current_time_ms() + 10000,
gene: vec![11],
environment: vec![12],
pending_mutation: vec![13],
pending_mutation_step: 14,
opcodes: vec![],
};
store.insert_session(&record).unwrap();
let loaded = store.load_session(&session_id).unwrap().unwrap();
assert_eq!(loaded.session_id, session_id);
assert_eq!(loaded.chain_length, 10);
let stats = store.stats().unwrap();
assert_eq!(stats.sessions, 1);
assert_eq!(stats.max_chain_length, 10);
store.purge_expired_sessions().unwrap();
let stats = store.stats().unwrap();
assert_eq!(stats.sessions, 1);
let _: Result<(), _> = conn.del(store.session_key(&session_id));
let _: Result<(), _> = conn.del(&store.index_key);
}
#[test]
fn test_valkey_pool_concurrency() {
use std::thread;
use std::sync::Arc;
let addr = std::env::var("CHRONOSEAL_VALKEY_ADDR").unwrap_or_else(|_| "127.0.0.1:6379".to_string());
let connection_string = format!("redis://{}", addr);
let client = match redis::Client::open(connection_string) {
Ok(c) => c,
Err(_) => return,
};
let pool = match r2d2::Pool::builder().build(client) {
Ok(p) => p,
Err(_) => return,
};
let mut conn = match pool.get() {
Ok(c) => c,
Err(_) => return,
};
let _: () = match redis::cmd("PING").query(&mut *conn) {
Ok(res) => res,
Err(_) => return,
};
let store = ValkeyStore {
pool,
index_key: "test:concurrent:sessions:ids".to_string(),
};
let _: Result<(), _> = conn.del("test:concurrent:sessions:ids");
let store_arc = Arc::new(store);
let mut handles = Vec::new();
for t in 0..10 {
let store_clone = store_arc.clone();
let session_id = format!("valkey_concurrent_{}", t);
let handle = thread::spawn(move || {
let record = SessionRecord {
session_id: session_id.clone(),
public_key: vec![1, 2, 3],
salt: vec![4, 5, 6],
last_hash: vec![7, 8, 9],
chain_length: 1,
created_at: 1000,
last_seen: 2000,
expires_at: current_time_ms() + 10000,
gene: vec![11],
environment: vec![12],
pending_mutation: vec![13],
pending_mutation_step: 14,
opcodes: vec![],
};
store_clone.insert_session(&record).unwrap();
let loaded = store_clone.load_session(&session_id).unwrap().unwrap();
assert_eq!(loaded.session_id, session_id);
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
let stats = store_arc.stats().unwrap();
assert_eq!(stats.sessions, 10);
// Cleanup
let mut conn = store_arc.pool.get().unwrap();
for t in 0..10 {
let _: Result<(), _> = conn.del(store_arc.session_key(&format!("valkey_concurrent_{}", t)));
}
let _: Result<(), _> = conn.del(&store_arc.index_key);
}
}
#[cfg(test)]
mod sqlite_tests {
use super::*;
use std::thread;
use std::sync::Arc;
#[test]
fn test_sqlite_pool_concurrency() {
let db_path = Path::new("target/test_sqlite_concurrency.db");
if let Some(parent) = db_path.parent() {
let _ = std::fs::create_dir_all(parent);
}
let _ = std::fs::remove_file(db_path);
let pool = init_pool(db_path).unwrap();
let pool_arc = Arc::new(pool);
let mut handles = Vec::new();
for t in 0..10 {
let pool_clone = pool_arc.clone();
let handle = thread::spawn(move || {
let session_id = format!("concurrent_session_{}", t);
let record = SessionRecord {
session_id: session_id.clone(),
public_key: vec![1, 2, 3],
salt: vec![4, 5, 6],
last_hash: vec![7, 8, 9],
chain_length: 1,
created_at: 1000,
last_seen: 2000,
expires_at: current_time_ms() + 10000,
gene: vec![11],
environment: vec![12],
pending_mutation: vec![13],
pending_mutation_step: 14,
opcodes: vec![],
};
pool_clone.insert_session(&record).unwrap();
let loaded = pool_clone.load_session(&session_id).unwrap().unwrap();
assert_eq!(loaded.session_id, session_id);
let mut updated = loaded;
updated.chain_length = 2;
let old_hash = updated.last_hash.clone();
pool_clone.update_session(&updated, &old_hash).unwrap();
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
let stats = pool_arc.stats().unwrap();
assert_eq!(stats.sessions, 10);
std::mem::drop(pool_arc);
let _ = std::fs::remove_file(db_path);
}
}
+9
View File
@@ -1,6 +1,14 @@
use crate::config::Config;
use shared::protocol::EntropyData;
/// Validates the browser mouse cursor interaction path for bot/automation detection.
///
/// Evaluates mouse velocity and distance features, checks the total distance traversed,
/// checks for cursor pauses (low movement over high time diff), and enforces average cursor speeds.
///
/// # Arguments
/// * `data` - The client-supplied interaction entropy events.
/// * `config` - The server configuration boundaries.
pub fn validate_mouse(
data: &EntropyData,
config: &Config,
@@ -41,6 +49,7 @@ pub fn validate_mouse(
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
+18
View File
@@ -4,6 +4,13 @@ use shared::{
vm_extensions::{self, ExecutionTrace, MutationError, MutationOrder},
};
/// Generates a randomized VM opcode instruction program within a length range.
///
/// Builds a program of mathematical and stack ops (e.g. literals, ADD, SUB, XOR, HASH)
/// with dynamic depth checking to ensure valid stacks and prevent out of bounds execution.
///
/// # Arguments
/// * `len_range` - The inclusive range of instruction counts to generate.
pub fn generate_random_program(len_range: std::ops::RangeInclusive<usize>) -> Vec<u8> {
let mut rng = rand::thread_rng();
let count = rng.gen_range(len_range);
@@ -47,6 +54,11 @@ pub fn generate_random_program(len_range: std::ops::RangeInclusive<usize>) -> Ve
ops
}
/// Executes a raw VM mutation program bytecode slice against a `GeneState`.
///
/// # Arguments
/// * `state` - The mutable gene state to mutate.
/// * `program` - The raw VM instruction program.
#[allow(dead_code)]
pub fn execute_mutation_program(
state: &mut GeneState,
@@ -55,6 +67,11 @@ pub fn execute_mutation_program(
vm_extensions::execute_program(state, program)
}
/// Executes a `MutationOrder` program against a `GeneState`.
///
/// # Arguments
/// * `state` - The mutable gene state.
/// * `order` - The mutation order.
#[allow(dead_code)]
pub fn execute_mutation_order(
state: &mut GeneState,
@@ -63,6 +80,7 @@ pub fn execute_mutation_order(
vm_extensions::execute_program(state, &order.program)
}
#[cfg(test)]
mod tests {
use super::*;