diff --git a/Cargo.lock b/Cargo.lock index 13be25a..8468ce6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -296,6 +296,7 @@ dependencies = [ "base64", "clap", "clap_complete", + "dashmap", "ed25519-dalek", "hex", "r2d2", @@ -509,6 +510,20 @@ dependencies = [ "syn", ] +[[package]] +name = "dashmap" +version = "6.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6361d5c062261c78a176addb82d4c821ae42bed6089de0e12603cd25de2059c" +dependencies = [ + "cfg-if", + "crossbeam-utils", + "hashbrown 0.14.5", + "lock_api", + "once_cell", + "parking_lot_core", +] + [[package]] name = "der" version = "0.7.10" diff --git a/Dockerfile b/Dockerfile index 94d49bd..d29b0da 100644 --- a/Dockerfile +++ b/Dockerfile @@ -24,4 +24,9 @@ ENV CHRONOSEAL_DB_PATH=/var/lib/chronoseal/chronoseal.sqlite ENV CHRONOSEAL_FRONTEND_DIR=/usr/share/chronoseal/frontend ENV CHRONOSEAL_PID_FILE=/run/chronoseal.pid +RUN useradd -r -s /bin/false chronoseal +USER chronoseal + +HEALTHCHECK --interval=30s --timeout=3s CMD chronoseal health || exit 1 + CMD ["chronoseal", "run"] diff --git a/chronoseal-replay/src/main.rs b/chronoseal-replay/src/main.rs index 053335d..7dcac8e 100644 --- a/chronoseal-replay/src/main.rs +++ b/chronoseal-replay/src/main.rs @@ -117,7 +117,11 @@ fn test_entropy() -> EntropyData { } } -fn do_handshake(client: &reqwest::blocking::Client, base_url: &str, sk: &SigningKey) -> Result { +fn do_handshake( + client: &reqwest::blocking::Client, + base_url: &str, + sk: &SigningKey, +) -> Result { let pk_hex = hex::encode(sk.verifying_key().to_bytes()); let init_req = InitRequest { public_key: pk_hex }; let resp = client @@ -126,7 +130,10 @@ fn do_handshake(client: &reqwest::blocking::Client, base_url: &str, sk: &Signing .send()?; if !resp.status().is_success() { - return Err(anyhow!("Handshake failed with HTTP status: {}", resp.status())); + return Err(anyhow!( + "Handshake failed with HTTP status: {}", + resp.status() + )); } let init_resp: InitResponse = resp.json()?; @@ -137,11 +144,17 @@ fn run_built_in_scenarios(client: &reqwest::blocking::Client, base_url: &str) -> let mut failures = 0; let scenarios = [ - ("valid_progression", run_valid_progression as fn(&reqwest::blocking::Client, &str) -> Result<()>), + ( + "valid_progression", + run_valid_progression as fn(&reqwest::blocking::Client, &str) -> Result<()>, + ), ("stale_replay", run_stale_replay), ("invalid_signature", run_invalid_signature), ("invalid_vm_stack", run_invalid_vm_stack), - ("invalid_mutation_commitment", run_invalid_mutation_commitment), + ( + "invalid_mutation_commitment", + run_invalid_mutation_commitment, + ), ("drifted_timestamp", run_drifted_timestamp), ("concurrent_heartbeat", run_concurrent_heartbeat), ("rate_limit_trigger", run_rate_limit_trigger), @@ -197,7 +210,8 @@ fn run_valid_progression(client: &reqwest::blocking::Client, base_url: &str) -> init.mutation_rounds, )?; - let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, mutation_step); + let commitment = + shared::gene::commitment_hex_with_context(&candidate, &init.session_id, mutation_step); let timestamp = current_time_ms(); let entropy = test_entropy(); @@ -215,13 +229,14 @@ fn run_valid_progression(client: &reqwest::blocking::Client, base_url: &str) -> sign_request(&sk, &mut req)?; - let resp = client - .post(format!("{}/hb", base_url)) - .json(&req) - .send()?; + let resp = client.post(format!("{}/hb", base_url)).json(&req).send()?; if !resp.status().is_success() { - return Err(anyhow!("Step {} /hb returned HTTP error: {}", step, resp.status())); + return Err(anyhow!( + "Step {} /hb returned HTTP error: {}", + step, + resp.status() + )); } let hb_resp: HeartbeatResponse = resp.json()?; @@ -230,9 +245,15 @@ fn run_valid_progression(client: &reqwest::blocking::Client, base_url: &str) -> } // Verify it was a successful validation (not a silent rejection) - let next_salt = hb_resp.next_salt.ok_or_else(|| anyhow!("Step {} was silently rejected", step))?; - let next_step = hb_resp.next_mutation_step.ok_or_else(|| anyhow!("Step {} missing next mutation step", step))?; - let next_order = hb_resp.next_mutation_order_b64.ok_or_else(|| anyhow!("Step {} missing next mutation order", step))?; + let next_salt = hb_resp + .next_salt + .ok_or_else(|| anyhow!("Step {} was silently rejected", step))?; + let next_step = hb_resp + .next_mutation_step + .ok_or_else(|| anyhow!("Step {} missing next mutation step", step))?; + let next_order = hb_resp + .next_mutation_order_b64 + .ok_or_else(|| anyhow!("Step {} missing next mutation order", step))?; println!("Step {} successful. Salt rotated: {}", step, next_salt); @@ -265,12 +286,21 @@ fn run_stale_replay(client: &reqwest::blocking::Client, base_url: &str) -> Resul let sk = SigningKey::generate(&mut csprng); let init = do_handshake(client, base_url, &sk)?; - let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?; + let opcodes = base64::Engine::decode( + &base64::engine::general_purpose::STANDARD, + &init.opcodes_b64, + )?; let stack_state = shared::vm::execute(&opcodes); - let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; + let order = + shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap(); - let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?; - let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); + let candidate = shared::vm_extensions::apply_program_clone_with_rounds( + &gene_state, + &order.program, + init.mutation_rounds, + )?; + let commitment = + shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); let mut req = HeartbeatRequest { session_id: init.session_id.clone(), @@ -296,7 +326,9 @@ fn run_stale_replay(client: &reqwest::blocking::Client, base_url: &str) -> Resul let resp2 = client.post(format!("{}/hb", base_url)).json(&req).send()?; let hb2: HeartbeatResponse = resp2.json()?; if hb2.next_salt.is_some() { - return Err(anyhow!("Replayed heartbeat was successfully accepted (broken replay protection)")); + return Err(anyhow!( + "Replayed heartbeat was successfully accepted (broken replay protection)" + )); } println!("Stale replay correctly rejected."); @@ -308,12 +340,21 @@ fn run_invalid_signature(client: &reqwest::blocking::Client, base_url: &str) -> let sk = SigningKey::generate(&mut csprng); let init = do_handshake(client, base_url, &sk)?; - let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?; + let opcodes = base64::Engine::decode( + &base64::engine::general_purpose::STANDARD, + &init.opcodes_b64, + )?; let stack_state = shared::vm::execute(&opcodes); - let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; + let order = + shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap(); - let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?; - let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); + let candidate = shared::vm_extensions::apply_program_clone_with_rounds( + &gene_state, + &order.program, + init.mutation_rounds, + )?; + let commitment = + shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); let mut req = HeartbeatRequest { session_id: init.session_id.clone(), @@ -344,10 +385,16 @@ fn run_invalid_vm_stack(client: &reqwest::blocking::Client, base_url: &str) -> R let sk = SigningKey::generate(&mut csprng); let init = do_handshake(client, base_url, &sk)?; - let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; + let order = + shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap(); - let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?; - let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); + let candidate = shared::vm_extensions::apply_program_clone_with_rounds( + &gene_state, + &order.program, + init.mutation_rounds, + )?; + let commitment = + shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); let mut req = HeartbeatRequest { session_id: init.session_id.clone(), @@ -375,12 +422,18 @@ fn run_invalid_vm_stack(client: &reqwest::blocking::Client, base_url: &str) -> R Ok(()) } -fn run_invalid_mutation_commitment(client: &reqwest::blocking::Client, base_url: &str) -> Result<()> { +fn run_invalid_mutation_commitment( + client: &reqwest::blocking::Client, + base_url: &str, +) -> Result<()> { let mut csprng = OsRng; let sk = SigningKey::generate(&mut csprng); let init = do_handshake(client, base_url, &sk)?; - let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?; + let opcodes = base64::Engine::decode( + &base64::engine::general_purpose::STANDARD, + &init.opcodes_b64, + )?; let stack_state = shared::vm::execute(&opcodes); let mut req = HeartbeatRequest { @@ -411,12 +464,21 @@ fn run_drifted_timestamp(client: &reqwest::blocking::Client, base_url: &str) -> let sk = SigningKey::generate(&mut csprng); let init = do_handshake(client, base_url, &sk)?; - let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?; + let opcodes = base64::Engine::decode( + &base64::engine::general_purpose::STANDARD, + &init.opcodes_b64, + )?; let stack_state = shared::vm::execute(&opcodes); - let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; + let order = + shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap(); - let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?; - let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); + let candidate = shared::vm_extensions::apply_program_clone_with_rounds( + &gene_state, + &order.program, + init.mutation_rounds, + )?; + let commitment = + shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); let mut req = HeartbeatRequest { session_id: init.session_id.clone(), @@ -446,12 +508,21 @@ fn run_concurrent_heartbeat(client: &reqwest::blocking::Client, base_url: &str) let sk = SigningKey::generate(&mut csprng); let init = do_handshake(client, base_url, &sk)?; - let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?; + let opcodes = base64::Engine::decode( + &base64::engine::general_purpose::STANDARD, + &init.opcodes_b64, + )?; let stack_state = shared::vm::execute(&opcodes); - let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; + let order = + shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap(); - let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?; - let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); + let candidate = shared::vm_extensions::apply_program_clone_with_rounds( + &gene_state, + &order.program, + init.mutation_rounds, + )?; + let commitment = + shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); let mut req = HeartbeatRequest { session_id: init.session_id.clone(), @@ -470,10 +541,8 @@ fn run_concurrent_heartbeat(client: &reqwest::blocking::Client, base_url: &str) let client_clone = client.clone(); let req_clone = req.clone(); let url_clone = format!("{}/hb", base_url); - - let handle = std::thread::spawn(move || { - client_clone.post(&url_clone).json(&req_clone).send() - }); + + let handle = std::thread::spawn(move || client_clone.post(&url_clone).json(&req_clone).send()); let resp2 = client.post(format!("{}/hb", base_url)).json(&req).send()?; let resp1_res = handle.join().map_err(|_| anyhow!("Thread panicked"))?; @@ -485,7 +554,10 @@ fn run_concurrent_heartbeat(client: &reqwest::blocking::Client, base_url: &str) // One must succeed and one must fail (silent rejection) because of CAS check let successes = (hb1.next_salt.is_some() as usize) + (hb2.next_salt.is_some() as usize); if successes != 1 { - return Err(anyhow!("Expected exactly one concurrent heartbeat to succeed. Got: {}", successes)); + return Err(anyhow!( + "Expected exactly one concurrent heartbeat to succeed. Got: {}", + successes + )); } println!("Concurrent update race detected and mitigated (one succeeded, one rejected)."); @@ -497,12 +569,21 @@ fn run_rate_limit_trigger(client: &reqwest::blocking::Client, base_url: &str) -> let sk = SigningKey::generate(&mut csprng); let init = do_handshake(client, base_url, &sk)?; - let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?; + let opcodes = base64::Engine::decode( + &base64::engine::general_purpose::STANDARD, + &init.opcodes_b64, + )?; let stack_state = shared::vm::execute(&opcodes); - let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; + let order = + shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?; let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap(); - let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?; - let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); + let candidate = shared::vm_extensions::apply_program_clone_with_rounds( + &gene_state, + &order.program, + init.mutation_rounds, + )?; + let commitment = + shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); let mut req = HeartbeatRequest { session_id: init.session_id.clone(), @@ -534,14 +615,20 @@ fn run_rate_limit_trigger(client: &reqwest::blocking::Client, base_url: &str) -> } if !rate_limited { - return Err(anyhow!("Rate limiter was not triggered after 35 rapid requests")); + return Err(anyhow!( + "Rate limiter was not triggered after 35 rapid requests" + )); } println!("Rate limiter correctly triggered."); Ok(()) } -fn run_file_scenario(_client: &reqwest::blocking::Client, _base_url: &str, file_path: &str) -> Result<()> { +fn run_file_scenario( + _client: &reqwest::blocking::Client, + _base_url: &str, + file_path: &str, +) -> Result<()> { let scenario_content = std::fs::read_to_string(file_path)?; let scenario: serde_json::Value = serde_json::from_str(&scenario_content)?; diff --git a/chronoseal.service b/chronoseal.service index 6dd7452..71ba079 100644 --- a/chronoseal.service +++ b/chronoseal.service @@ -27,6 +27,7 @@ RestrictSUIDSGID=yes LockPersonality=yes SystemCallArchitectures=native ReadWritePaths=/run/chronoseal.pid +ReadWritePaths=/var/lib/chronoseal # Logging StandardOutput=journal diff --git a/docs/API.md b/docs/API.md index 6339575..8eef399 100644 --- a/docs/API.md +++ b/docs/API.md @@ -161,7 +161,7 @@ Content-Type: application/json | `stack_state.ip` | number | yes | VM instruction pointer as an unsigned 16-bit value | | `fingerprint.aspectRatio` | string | yes | Screen aspect ratio; server accepts numeric strings in range `0.5..=3.0` | | `fingerprint.devicePixelRatio` | string | yes | Device pixel ratio; server accepts numeric strings in range `(0, 5]` | -| `fingerprint.hardwareConcurrency` | number | yes | Positive hardware concurrency value | +| `fingerprint.hardwareConcurrency` | number | yes | Hardware concurrency value; server accepts integers in range `1..=256` | | `mutation_step` | number | yes | Mutation step currently expected by the server | | `gene_commitment` | string | yes | Context-bound commitment produced by the WASM mutation preview | | `signature` | string | yes | Ed25519 signature over the canonical payload | diff --git a/frontend/favicon/apple-touch-icon.png b/frontend/favicon/apple-touch-icon.png new file mode 100644 index 0000000..6ae5322 Binary files /dev/null and b/frontend/favicon/apple-touch-icon.png differ diff --git a/frontend/favicon/favicon-96x96.png b/frontend/favicon/favicon-96x96.png new file mode 100644 index 0000000..91b532e Binary files /dev/null and b/frontend/favicon/favicon-96x96.png differ diff --git a/frontend/favicon/favicon.ico b/frontend/favicon/favicon.ico new file mode 100644 index 0000000..c9fa133 Binary files /dev/null and b/frontend/favicon/favicon.ico differ diff --git a/frontend/favicon/favicon.svg b/frontend/favicon/favicon.svg new file mode 100644 index 0000000..596088b --- /dev/null +++ b/frontend/favicon/favicon.svg @@ -0,0 +1 @@ +RealFaviconGeneratorhttps://realfavicongenerator.netchronoseal \ No newline at end of file diff --git a/frontend/favicon/site.webmanifest b/frontend/favicon/site.webmanifest new file mode 100644 index 0000000..ccf313a --- /dev/null +++ b/frontend/favicon/site.webmanifest @@ -0,0 +1,21 @@ +{ + "name": "MyWebSite", + "short_name": "MySite", + "icons": [ + { + "src": "/web-app-manifest-192x192.png", + "sizes": "192x192", + "type": "image/png", + "purpose": "maskable" + }, + { + "src": "/web-app-manifest-512x512.png", + "sizes": "512x512", + "type": "image/png", + "purpose": "maskable" + } + ], + "theme_color": "#ffffff", + "background_color": "#ffffff", + "display": "standalone" +} \ No newline at end of file diff --git a/frontend/favicon/web-app-manifest-192x192.png b/frontend/favicon/web-app-manifest-192x192.png new file mode 100644 index 0000000..c26b7ca Binary files /dev/null and b/frontend/favicon/web-app-manifest-192x192.png differ diff --git a/frontend/favicon/web-app-manifest-512x512.png b/frontend/favicon/web-app-manifest-512x512.png new file mode 100644 index 0000000..b3eda46 Binary files /dev/null and b/frontend/favicon/web-app-manifest-512x512.png differ diff --git a/frontend/heartbeat.js b/frontend/heartbeat.js index 219f125..f90ce56 100644 --- a/frontend/heartbeat.js +++ b/frontend/heartbeat.js @@ -77,7 +77,6 @@ async function sendHeartbeat() { const sig = sign_message(msg); if (!sig) { discard_gene_preview(); - console.error('Keypair not initialised — skipping heartbeat'); return; } const resp = await sendRequest('/hb', 'POST', { @@ -107,11 +106,9 @@ async function sendHeartbeat() { pendingMutationOrderB64 = resp.next_mutation_order_b64; } else { discard_gene_preview(); - console.warn('Heartbeat rejected'); } } catch (e) { discard_gene_preview(); - console.error(e); } finally { scheduleNext(); } diff --git a/frontend/index.html b/frontend/index.html index b3dc31e..4e841d4 100644 --- a/frontend/index.html +++ b/frontend/index.html @@ -2,6 +2,7 @@ + Anti-Scraper Demo diff --git a/frontend/main.js b/frontend/main.js index e44681f..a9900af 100644 --- a/frontend/main.js +++ b/frontend/main.js @@ -1,5 +1,3 @@ import { initHeartbeat } from './heartbeat.js'; -(async () => { - await initHeartbeat(); -})(); \ No newline at end of file +initHeartbeat().catch(() => {}); \ No newline at end of file diff --git a/scripts/dev.sh b/scripts/dev.sh index 9a6840d..0cb3c42 100644 --- a/scripts/dev.sh +++ b/scripts/dev.sh @@ -1,4 +1,5 @@ #!/bin/bash +set -euo pipefail echo "Starting server with static frontend serving..." cd ../server cargo run --release \ No newline at end of file diff --git a/server/Cargo.toml b/server/Cargo.toml index e39534e..0166fbc 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -31,4 +31,5 @@ base64 = "0.22" rand = "0.8" ed25519-dalek = "2" redis = { version = "0.29", features = ["r2d2"] } +dashmap = "6" diff --git a/server/src/cleanup.rs b/server/src/cleanup.rs index ffa020a..ed15562 100644 --- a/server/src/cleanup.rs +++ b/server/src/cleanup.rs @@ -20,12 +20,10 @@ pub async fn cleanup_loop(state: Arc) { } } - // Evict stale rate-limiter entries to prevent unbounded HashMap growth. + // Evict stale rate-limiter entries to prevent unbounded map growth. { let window_secs = state.get_config().rate_limit_window_secs; - let mut rl = state.rate_limiter.lock().await; - rl.evict_stale(window_secs); + state.rate_limiter.evict_stale(window_secs); } } } - diff --git a/server/src/config.rs b/server/src/config.rs index 730365b..a45fd6a 100644 --- a/server/src/config.rs +++ b/server/src/config.rs @@ -137,7 +137,9 @@ impl Config { size: self.gene_size, }); } - if !(1..=shared::constants::MAX_MUTATION_ROUNDS).contains(&self.mutation_rounds) { + if !(shared::constants::MIN_MUTATION_ROUNDS..=shared::constants::MAX_MUTATION_ROUNDS) + .contains(&self.mutation_rounds) + { return Err(ConfigError::InvalidMutationRounds { rounds: self.mutation_rounds, }); @@ -274,7 +276,8 @@ impl std::fmt::Display for ConfigError { Self::InvalidMutationRounds { rounds } => { write!( f, - "invalid mutation rounds {rounds}; expected 1..={}", + "invalid mutation rounds {rounds}; expected {}..={}", + shared::constants::MIN_MUTATION_ROUNDS, shared::constants::MAX_MUTATION_ROUNDS ) } diff --git a/server/src/crypto.rs b/server/src/crypto.rs index ead61c1..100c9ee 100644 --- a/server/src/crypto.rs +++ b/server/src/crypto.rs @@ -44,4 +44,3 @@ pub fn verify_signature( pk.verify_strict(message.as_bytes(), &sig)?; Ok(()) } - diff --git a/server/src/errors.rs b/server/src/errors.rs index aab48bf..c75df8a 100644 --- a/server/src/errors.rs +++ b/server/src/errors.rs @@ -25,12 +25,19 @@ pub enum SessionError { #[error("Invalid gene configuration: {0}")] InvalidGeneConfiguration(String), + + #[error("Rate limited")] + RateLimited, } impl IntoResponse for SessionError { fn into_response(self) -> Response { let (status, error_message) = match self { SessionError::InvalidPublicKeyLength => (StatusCode::BAD_REQUEST, self.to_string()), + SessionError::RateLimited => ( + StatusCode::TOO_MANY_REQUESTS, + "Too many requests".to_string(), + ), _ => ( StatusCode::INTERNAL_SERVER_ERROR, "Internal server error".to_string(), diff --git a/server/src/fingerprint.rs b/server/src/fingerprint.rs index 18982c1..3d4e8fa 100644 --- a/server/src/fingerprint.rs +++ b/server/src/fingerprint.rs @@ -1,5 +1,10 @@ use shared::protocol::Fingerprint; +const MIN_ASPECT_RATIO: f64 = 0.5; +const MAX_ASPECT_RATIO: f64 = 3.0; +const MAX_DEVICE_PIXEL_RATIO: f64 = 5.0; +const MAX_HARDWARE_CONCURRENCY: u32 = 256; + /// Validates the browser fingerprint fields submitted by the client. /// /// Checks basic screen aspect ratio thresholds, device pixel ratio limits, @@ -9,16 +14,66 @@ use shared::protocol::Fingerprint; /// * `fp` - The client's hardware and screen layout fingerprint. pub fn validate(fp: &Fingerprint) -> Result<(), Box> { let ar: f64 = fp.aspect_ratio.parse().map_err(|_| "ar")?; - if !(0.5..=3.0).contains(&ar) { + if !ar.is_finite() || !(MIN_ASPECT_RATIO..=MAX_ASPECT_RATIO).contains(&ar) { return Err("aspect ratio".into()); } + let dpr: f64 = fp.device_pixel_ratio.parse().map_err(|_| "dpr")?; - if dpr <= 0.0 || dpr > 5.0 { + if !dpr.is_finite() || dpr <= 0.0 || dpr > MAX_DEVICE_PIXEL_RATIO { return Err("dpr".into()); } - if fp.hardware_concurrency == 0 { + + if fp.hardware_concurrency == 0 || fp.hardware_concurrency > MAX_HARDWARE_CONCURRENCY { return Err("hw".into()); } + Ok(()) } +#[cfg(test)] +mod tests { + use super::*; + + fn fingerprint( + aspect_ratio: impl Into, + device_pixel_ratio: impl Into, + hardware_concurrency: u32, + ) -> Fingerprint { + Fingerprint { + aspect_ratio: aspect_ratio.into(), + device_pixel_ratio: device_pixel_ratio.into(), + hardware_concurrency, + } + } + + #[test] + fn accepts_valid_fingerprint() { + assert!(validate(&fingerprint("1.7777777778", "2", 8)).is_ok()); + } + + #[test] + fn accepts_boundary_values() { + assert!(validate(&fingerprint("0.5", "1", 1)).is_ok()); + assert!(validate(&fingerprint("3.0", "5.0", MAX_HARDWARE_CONCURRENCY)).is_ok()); + } + + #[test] + fn rejects_invalid_aspect_ratios() { + for aspect_ratio in ["not-a-number", "NaN", "inf", "0.49", "3.01"] { + assert!(validate(&fingerprint(aspect_ratio, "2", 8)).is_err()); + } + } + + #[test] + fn rejects_invalid_device_pixel_ratios() { + for device_pixel_ratio in ["not-a-number", "NaN", "inf", "0", "-1", "5.01"] { + assert!(validate(&fingerprint("1.77", device_pixel_ratio, 8)).is_err()); + } + } + + #[test] + fn rejects_invalid_hardware_concurrency() { + assert!(validate(&fingerprint("1.77", "2", 0)).is_err()); + assert!(validate(&fingerprint("1.77", "2", MAX_HARDWARE_CONCURRENCY + 1)).is_err()); + } +} diff --git a/server/src/main.rs b/server/src/main.rs index d404c3b..a127f65 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -30,9 +30,6 @@ async fn main() { async fn try_main() -> Result<(), Box> { let cli = Cli::parse(); - if let Some(config_path) = cli.globals.config.as_deref() { - std::env::set_var("CHRONOSEAL_CONFIG", config_path); - } let log_filter = cli.globals.log.as_deref().unwrap_or("info"); let log_file = log_file_for_command(&cli); let _log_guard = init_logging(log_filter, log_file)?; diff --git a/server/src/middleware.rs b/server/src/middleware.rs index b2b6971..e30d343 100644 --- a/server/src/middleware.rs +++ b/server/src/middleware.rs @@ -9,3 +9,25 @@ pub async fn log_request(req: Request, next: Next) -> Response { tracing::info!("{} {} -> {}", method, uri, response.status()); response } + +/// Injects defensive HTTP response headers on every response. +/// +/// These headers mitigate several classes of attacks: +/// - `X-Content-Type-Options: nosniff` — prevents MIME-type sniffing. +/// - `X-Frame-Options: DENY` — blocks clickjacking via framing. +/// - `Referrer-Policy: no-referrer` — suppresses referrer leakage. +/// - `X-XSS-Protection: 0` — disables legacy XSS auditors (can introduce bugs). +/// - `Permissions-Policy` — restricts powerful browser features. +pub async fn security_headers(req: Request, next: Next) -> Response { + let mut response = next.run(req).await; + let headers = response.headers_mut(); + headers.insert("x-content-type-options", "nosniff".parse().unwrap()); + headers.insert("x-frame-options", "DENY".parse().unwrap()); + headers.insert("referrer-policy", "no-referrer".parse().unwrap()); + headers.insert("x-xss-protection", "0".parse().unwrap()); + headers.insert( + "permissions-policy", + "camera=(), microphone=(), geolocation=()".parse().unwrap(), + ); + response +} diff --git a/server/src/ratelimit.rs b/server/src/ratelimit.rs index 45d174b..b286e85 100644 --- a/server/src/ratelimit.rs +++ b/server/src/ratelimit.rs @@ -1,17 +1,20 @@ -use std::collections::HashMap; +use dashmap::DashMap; use std::time::Instant; -/// A simple, in-memory sliding-window rate limiter for tracking client heartbeat frequency. +/// A lock-free, concurrent sliding-window rate limiter backed by `DashMap`. +/// +/// All public methods take `&self` (no `&mut self`), so the limiter can live in +/// an `Arc` without a `Mutex` wrapper. pub struct RateLimiter { - /// Maps session identifiers to request counts and window start timestamps. - buckets: HashMap, + /// Maps rate-limit keys to request counts and window start timestamps. + buckets: DashMap, } impl RateLimiter { /// Creates a new, empty `RateLimiter`. pub fn new() -> Self { Self { - buckets: HashMap::new(), + buckets: DashMap::new(), } } @@ -20,19 +23,21 @@ impl RateLimiter { /// Returns `true` if allowed, or `false` if the rate limit is exceeded. /// /// # Arguments - /// * `key` - The unique identifier to rate-limit (e.g., session ID). + /// * `key` - The unique identifier to rate-limit (e.g., client IP address). /// * `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 { + pub fn check(&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)); - if now.duration_since(entry.1).as_secs() >= window_secs { - *entry = (1, now); + let mut entry = self.buckets.entry(key.to_string()).or_insert((0, now)); + let (count, ts) = entry.value_mut(); + if now.duration_since(*ts).as_secs() >= window_secs { + *count = 1; + *ts = now; true - } else if entry.0 >= limit { + } else if *count >= limit { false } else { - entry.0 += 1; + *count += 1; true } } @@ -43,7 +48,7 @@ impl RateLimiter { /// /// # Arguments /// * `window_secs` - The active rate-limiting window duration in seconds. - pub fn evict_stale(&mut self, window_secs: u64) { + pub fn evict_stale(&self, window_secs: u64) { let now = Instant::now(); self.buckets .retain(|_, (_, ts)| now.duration_since(*ts).as_secs() < window_secs); @@ -58,7 +63,7 @@ mod tests { #[test] fn test_rate_limiter() { - let mut rl = RateLimiter::new(); + let rl = RateLimiter::new(); // Limit of 2 requests per 1 second window assert!(rl.check("user1", 2, 1)); assert!(rl.check("user1", 2, 1)); @@ -72,7 +77,7 @@ mod tests { #[test] fn test_rate_limiter_eviction() { - let mut rl = RateLimiter::new(); + let rl = RateLimiter::new(); assert!(rl.check("user1", 1, 1)); assert_eq!(rl.buckets.len(), 1); diff --git a/server/src/routes/heartbeat.rs b/server/src/routes/heartbeat.rs index c619348..54316e6 100644 --- a/server/src/routes/heartbeat.rs +++ b/server/src/routes/heartbeat.rs @@ -8,21 +8,51 @@ pub async fn handler( Json(payload): Json, ) -> (StatusCode, Json) { let start_http = std::time::Instant::now(); - state.heartbeats_total.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + state + .heartbeats_total + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); - // Rate limiting + // Cap entropy events to prevent oversized payloads from exhausting memory. + if payload.entropy_data.events.len() > 1000 { + 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 { + status: "ok".into(), + next_salt: None, + next_mutation_step: None, + next_mutation_order_b64: None, + }), + ); + } + + // Rate limiting (lock-free via DashMap) { let (limit, window_secs) = { let cfg = state.get_config(); (cfg.rate_limit_count, cfg.rate_limit_window_secs) }; - let mut rl = state.rate_limiter.lock().await; - if !rl.check(&payload.session_id, limit, window_secs) { + if !state + .rate_limiter + .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); + 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); + 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 { @@ -39,8 +69,12 @@ pub async fn handler( 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 + 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) => ( @@ -54,15 +88,21 @@ 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); + 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); + 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); + state + .mutation_failures_total + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); } _ => {} } @@ -79,8 +119,12 @@ 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); + 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 } @@ -138,7 +182,11 @@ mod tests { }, ], }; - let program_bytes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64).unwrap(); + 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 = @@ -176,7 +224,7 @@ mod tests { let pool = crate::storage::init_pool(Path::new(":memory:")).unwrap(); let state = Arc::new(AppState { db_pool: pool.clone(), - rate_limiter: tokio::sync::Mutex::new(crate::ratelimit::RateLimiter::new()), + rate_limiter: 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), diff --git a/server/src/routes/init.rs b/server/src/routes/init.rs index 338a42d..f21f106 100644 --- a/server/src/routes/init.rs +++ b/server/src/routes/init.rs @@ -10,18 +10,35 @@ pub async fn handler( ) -> Result, SessionError> { let start_http = std::time::Instant::now(); let config = state.get_config(); - + + // Rate limit session creation by public key to prevent storage exhaustion. + if !state.rate_limiter.check( + &payload.public_key, + config.rate_limit_count, + config.rate_limit_window_secs, + ) { + return Err(SessionError::RateLimited); + } + 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); + 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); + 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)) } diff --git a/server/src/runtime.rs b/server/src/runtime.rs index 812445f..a70fcfd 100644 --- a/server/src/runtime.rs +++ b/server/src/runtime.rs @@ -5,17 +5,19 @@ use crate::{ routes, session, storage::{self, StoreStats}, }; -use axum::{http::StatusCode, response::IntoResponse, routing::get, Json, Router}; +use axum::{ + extract::ConnectInfo, http::StatusCode, response::IntoResponse, routing::get, Json, Router, +}; use serde::Serialize; use std::{ fs, io::{Read, Write}, - net::{SocketAddr, TcpStream}, + net::{IpAddr, SocketAddr, TcpStream}, path::Path, sync::Arc, time::Duration, }; -use tokio::sync::{Mutex, Notify}; +use tokio::sync::Notify; use tracing::{error, info, warn}; #[derive(Debug, Serialize)] @@ -144,7 +146,7 @@ pub async fn run_daemon(config: Config) -> Result<(), Box let db_pool = init_db_pool(&config)?; let state = Arc::new(session::AppState { db_pool, - rate_limiter: Mutex::new(RateLimiter::new()), + rate_limiter: 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), @@ -170,7 +172,11 @@ pub async fn run_daemon(config: Config) -> Result<(), Box tower_http::services::ServeDir::new(&config.frontend_dir), ) .layer(tower_http::cors::CorsLayer::permissive()) + .layer(axum::middleware::from_fn( + crate::middleware::security_headers, + )) .layer(axum::middleware::from_fn(crate::middleware::log_request)) + .layer(axum::extract::DefaultBodyLimit::max(64 * 1024)) // 64 KiB .with_state(state.clone()); let addr: SocketAddr = config.bind.parse()?; @@ -178,9 +184,12 @@ pub async fn run_daemon(config: Config) -> Result<(), Box info!(bind = %config.bind, "chronoseal daemon started"); let shutdown = signal_task(state.clone()); - let result = axum::serve(listener, app) - .with_graceful_shutdown(shutdown) - .await; + let result = axum::serve( + listener, + app.into_make_service_with_connect_info::(), + ) + .with_graceful_shutdown(shutdown) + .await; remove_pid_file(&config.pid_file); result?; @@ -277,8 +286,12 @@ async fn health_handler() -> impl IntoResponse { } async fn stats_handler( + ConnectInfo(addr): ConnectInfo, axum::extract::State(state): axum::extract::State>, ) -> Result, (StatusCode, String)> { + if !is_loopback(addr.ip()) { + return Err((StatusCode::FORBIDDEN, "Forbidden".to_string())); + } state .db_pool .stats() @@ -287,25 +300,45 @@ async fn stats_handler( } async fn metrics_handler( + ConnectInfo(addr): ConnectInfo, axum::extract::State(state): axum::extract::State>, ) -> Result { + if !is_loopback(addr.ip()) { + return Err((StatusCode::FORBIDDEN, "Forbidden".to_string())); + } let stats = state .db_pool .stats() .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 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_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 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_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); + let http_count = state + .http_ops_count + .load(std::sync::atomic::Ordering::Relaxed); Ok(format!( "# HELP chronoseal_active_sessions Active ChronoSeal sessions\n\ @@ -355,6 +388,13 @@ async fn metrics_handler( )) } +fn is_loopback(ip: IpAddr) -> bool { + match ip { + IpAddr::V4(v4) => v4.is_loopback(), + IpAddr::V6(v6) => v6.is_loopback(), + } +} + async fn signal_task(state: Arc) { let shutdown = Arc::new(Notify::new()); diff --git a/server/src/session.rs b/server/src/session.rs index a8467b0..bedc000 100644 --- a/server/src/session.rs +++ b/server/src/session.rs @@ -1,6 +1,6 @@ pub struct AppState { pub db_pool: crate::storage::DbPool, - pub rate_limiter: tokio::sync::Mutex, + pub rate_limiter: crate::ratelimit::RateLimiter, pub config: std::sync::RwLock, pub heartbeats_total: std::sync::atomic::AtomicU64, pub verification_failures_total: std::sync::atomic::AtomicU64, @@ -271,7 +271,6 @@ mod tests { } } - fn test_fingerprint() -> Fingerprint { Fingerprint { aspect_ratio: "1.77".to_string(), @@ -323,8 +322,12 @@ mod tests { vm_extensions::apply_program_clone(&client.committed_gene_state, &order.program) .unwrap(); let entropy = test_entropy(); - - let program_bytes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &client.opcodes_b64).unwrap(); + + 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 { diff --git a/server/src/storage.rs b/server/src/storage.rs index 709714e..8ffc76c 100644 --- a/server/src/storage.rs +++ b/server/src/storage.rs @@ -1,8 +1,8 @@ use crate::config::Config; +use redis::Commands; use serde::{Deserialize, Serialize}; use std::path::Path; use std::time::{SystemTime, UNIX_EPOCH}; -use redis::Commands; #[derive(Debug, Clone, Serialize, Deserialize)] pub struct StoreStats { @@ -54,11 +54,12 @@ impl DbPool { crate::config::DbType::Valkey => { let addr = std::env::var("CHRONOSEAL_VALKEY_ADDR") .unwrap_or_else(|_| "127.0.0.1:6666".to_string()); - let connection_string = if addr.starts_with("redis://") || addr.starts_with("rediss://") { - addr.clone() - } else { - format!("redis://{}", addr) - }; + 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 { @@ -354,9 +355,19 @@ impl ValkeyStore { 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) + .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(()) } @@ -400,9 +411,19 @@ impl ValkeyStore { 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) + .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 { @@ -422,8 +443,12 @@ impl ValkeyStore { 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) + .cmd("ZREM") + .arg(&self.index_key) + .arg(&expired_ids) + .cmd("ZREM") + .arg("sessions:chain_lengths") + .arg(&expired_ids) .query::<()>(&mut *conn)?; } Ok(()) @@ -435,8 +460,12 @@ impl ValkeyStore { 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); + 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, @@ -449,7 +478,7 @@ impl ValkeyStore { pub fn current_time_ms() -> u64 { SystemTime::now() .duration_since(UNIX_EPOCH) - .unwrap() + .expect("system clock is before UNIX epoch; check system time") .as_millis() as u64 } @@ -459,7 +488,8 @@ mod valkey_tests { #[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 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, @@ -522,10 +552,11 @@ mod valkey_tests { #[test] fn test_valkey_pool_concurrency() { - use std::thread; use std::sync::Arc; + use std::thread; - let addr = std::env::var("CHRONOSEAL_VALKEY_ADDR").unwrap_or_else(|_| "127.0.0.1:6379".to_string()); + 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, @@ -589,7 +620,8 @@ mod valkey_tests { // 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.session_key(&format!("valkey_concurrent_{}", t))); } let _: Result<(), _> = conn.del(&store_arc.index_key); } @@ -598,8 +630,8 @@ mod valkey_tests { #[cfg(test)] mod sqlite_tests { use super::*; - use std::thread; use std::sync::Arc; + use std::thread; #[test] fn test_sqlite_pool_concurrency() { @@ -655,5 +687,3 @@ mod sqlite_tests { let _ = std::fs::remove_file(db_path); } } - - diff --git a/server/src/trust.rs b/server/src/trust.rs index df05bb5..6c3f498 100644 --- a/server/src/trust.rs +++ b/server/src/trust.rs @@ -49,7 +49,6 @@ pub fn validate_mouse( Ok(()) } - #[cfg(test)] mod tests { use super::*; diff --git a/server/src/vm.rs b/server/src/vm.rs index af6775b..5ff0be6 100644 --- a/server/src/vm.rs +++ b/server/src/vm.rs @@ -80,7 +80,6 @@ pub fn execute_mutation_order( vm_extensions::execute_program(state, &order.program) } - #[cfg(test)] mod tests { use super::*; diff --git a/shared/src/constants.rs b/shared/src/constants.rs index 7abfb32..36dd3ea 100644 --- a/shared/src/constants.rs +++ b/shared/src/constants.rs @@ -10,3 +10,4 @@ pub const MAX_MUTATION_ROUNDS: u8 = 10; pub const MAX_MUTATION_INSTRUCTION_BUDGET: usize = 2048; pub const HASH_OPCODE_INSTRUCTION_COST: usize = 16; pub const SOFT_CAP_DURATION_MS: u128 = 50; +pub const MAX_STACK_DEPTH: usize = 64; diff --git a/shared/src/hashing.rs b/shared/src/hashing.rs index b3e00bf..65f1c65 100644 --- a/shared/src/hashing.rs +++ b/shared/src/hashing.rs @@ -37,8 +37,8 @@ pub fn next_chain_hash( stack: &StackState, salt: &[u8], ) -> Vec { - let entropy_bytes = serde_json::to_vec(entropy).unwrap(); - let stack_bytes = serde_json::to_vec(stack).unwrap(); + let entropy_bytes = serde_json::to_vec(entropy).unwrap_or_default(); + let stack_bytes = serde_json::to_vec(stack).unwrap_or_default(); let entropy_hash = blake3::hash(&entropy_bytes); let stack_hash = blake3::hash(&stack_bytes); @@ -63,4 +63,3 @@ pub fn hash_stack(stack: &[u32]) -> u32 { let hash = blake3::hash(&data); u32::from_le_bytes(hash.as_bytes()[..4].try_into().unwrap()) } - diff --git a/shared/src/protocol.rs b/shared/src/protocol.rs index f79ccdc..52af8fb 100644 --- a/shared/src/protocol.rs +++ b/shared/src/protocol.rs @@ -116,4 +116,3 @@ pub struct StackState { /// The final instruction pointer location at program completion or termination. pub ip: u16, } - diff --git a/shared/src/vm.rs b/shared/src/vm.rs index e26ce6f..e857f15 100644 --- a/shared/src/vm.rs +++ b/shared/src/vm.rs @@ -25,6 +25,9 @@ pub fn execute(program: &[u8]) -> StackState { program[ip + 3], ]); ip += 4; + if stack.len() >= crate::constants::MAX_STACK_DEPTH { + break; + } stack.push(val); } 0x01..=0x07 => { @@ -43,6 +46,9 @@ pub fn execute(program: &[u8]) -> StackState { 0x07 => a.rotate_left(b % 32), _ => unreachable!(), }; + if stack.len() >= crate::constants::MAX_STACK_DEPTH { + break; + } stack.push(r); } 0x08 => { @@ -50,11 +56,17 @@ pub fn execute(program: &[u8]) -> StackState { break; } let a = stack.pop().unwrap(); + if stack.len() >= crate::constants::MAX_STACK_DEPTH { + break; + } stack.push(!a); } 0x09 => { let r = crate::hashing::hash_stack(&stack); stack.clear(); + if stack.len() >= crate::constants::MAX_STACK_DEPTH { + break; + } stack.push(r); } _ => break, diff --git a/shared/src/vm_extensions.rs b/shared/src/vm_extensions.rs index 34644ed..24aac7f 100644 --- a/shared/src/vm_extensions.rs +++ b/shared/src/vm_extensions.rs @@ -269,7 +269,6 @@ pub fn apply_program(state: &mut GeneState, program: &[u8]) -> Result<(), Mutati Ok(()) } - /// Executes the mutation program on the mutable `GeneState` reference for multiple rounds. /// /// Implements a soft instruction cost-budget cap check to prevent hostile/inefficient diff --git a/wasm/src/crypto.rs b/wasm/src/crypto.rs index 83eec88..e45da2e 100644 --- a/wasm/src/crypto.rs +++ b/wasm/src/crypto.rs @@ -51,8 +51,14 @@ pub fn compute_next_hash( ) -> String { let prev = hex::decode(prev_hash_hex).unwrap_or_default(); let salt = hex::decode(salt_hex).unwrap_or_default(); - let entropy = serde_json::from_str::(entropy_data_json).unwrap(); - let stack = serde_json::from_str::(stack_state_json).unwrap(); + let entropy = match serde_json::from_str::(entropy_data_json) { + Ok(v) => v, + Err(_) => return String::new(), + }; + let stack = match serde_json::from_str::(stack_state_json) { + Ok(v) => v, + Err(_) => return String::new(), + }; let new = shared::hashing::next_chain_hash(&prev, timestamp, &entropy, &stack, &salt); hex::encode(new) } diff --git a/wasm/src/vm.rs b/wasm/src/vm.rs index 1f6a070..e40e425 100644 --- a/wasm/src/vm.rs +++ b/wasm/src/vm.rs @@ -4,11 +4,15 @@ use wasm_bindgen::prelude::*; #[wasm_bindgen] pub fn run_program(program_b64: &str) -> JsValue { use base64::Engine; - let bytes = base64::engine::general_purpose::STANDARD - .decode(program_b64) - .unwrap(); + let bytes = match base64::engine::general_purpose::STANDARD.decode(program_b64) { + Ok(v) => v, + Err(_) => return JsValue::NULL, + }; let state = execute(&bytes); - serde_wasm_bindgen::to_value(&state).unwrap() + match serde_wasm_bindgen::to_value(&state) { + Ok(v) => v, + Err(_) => JsValue::NULL, + } } fn execute(program: &[u8]) -> StackState {