6 Commits
50 changed files with 1564 additions and 402 deletions

No files matched your search

-2
View File
@@ -5,5 +5,3 @@ dist/
*.log *.log
.env .env
.idea/ .idea/
.antigravitycli/
Generated
+19 -4
View File
@@ -275,7 +275,7 @@ dependencies = [
[[package]] [[package]]
name = "chronoseal-replay" name = "chronoseal-replay"
version = "0.1.0" version = "1.0.1"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"base64", "base64",
@@ -290,12 +290,13 @@ dependencies = [
[[package]] [[package]]
name = "chronoseal-server" name = "chronoseal-server"
version = "0.6.0" version = "1.0.1"
dependencies = [ dependencies = [
"axum", "axum",
"base64", "base64",
"clap", "clap",
"clap_complete", "clap_complete",
"dashmap",
"ed25519-dalek", "ed25519-dalek",
"hex", "hex",
"r2d2", "r2d2",
@@ -319,7 +320,7 @@ dependencies = [
[[package]] [[package]]
name = "chronoseal-wasm" name = "chronoseal-wasm"
version = "0.6.0" version = "1.0.1"
dependencies = [ dependencies = [
"base64", "base64",
"blake3", "blake3",
@@ -509,6 +510,20 @@ dependencies = [
"syn", "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]] [[package]]
name = "der" name = "der"
version = "0.7.10" version = "0.7.10"
@@ -1946,7 +1961,7 @@ dependencies = [
[[package]] [[package]]
name = "shared" name = "shared"
version = "0.6.0" version = "1.0.1"
dependencies = [ dependencies = [
"base64", "base64",
"blake3", "blake3",
+2 -1
View File
@@ -4,7 +4,8 @@ resolver = "2"
members = [ members = [
"shared", "shared",
"server", "server",
"wasm" "wasm",
"chronoseal-replay"
] ]
[workspace.package] [workspace.package]
+5
View File
@@ -24,4 +24,9 @@ ENV CHRONOSEAL_DB_PATH=/var/lib/chronoseal/chronoseal.sqlite
ENV CHRONOSEAL_FRONTEND_DIR=/usr/share/chronoseal/frontend ENV CHRONOSEAL_FRONTEND_DIR=/usr/share/chronoseal/frontend
ENV CHRONOSEAL_PID_FILE=/run/chronoseal.pid 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"] CMD ["chronoseal", "run"]
+27 -2
View File
@@ -20,7 +20,7 @@
<img src="https://img.shields.io/badge/rust-stable%20%E2%89%A5%201.87-orange.svg" alt="Rust stable >= 1.87"> <img src="https://img.shields.io/badge/rust-stable%20%E2%89%A5%201.87-orange.svg" alt="Rust stable >= 1.87">
</a> </a>
<a href="https://github.com/thakares/chronoseal-rs/blob/main/docs/REFRACTORING-v0.6.0.md"> <a href="https://github.com/thakares/chronoseal-rs/blob/main/docs/REFRACTORING-v0.6.0.md">
<img src="https://img.shields.io/badge/version-v0.6.0-green.svg" alt="v0.6.0"> <img src="https://img.shields.io/badge/version-v1.0.1-green.svg" alt="v0.6.0">
</a> </a>
<img src="https://img.shields.io/badge/wasm-rust--compiled-blueviolet.svg" alt="WASM"> <img src="https://img.shields.io/badge/wasm-rust--compiled-blueviolet.svg" alt="WASM">
</p> </p>
@@ -536,7 +536,32 @@ Set storage mode with `db_type` or `CHRONOSEAL_DB_TYPE`.
| `sqlite-in-disk` | SQLite database persisted at `db_path`. | | `sqlite-in-disk` | SQLite database persisted at `db_path`. |
| `valkey` | Valkey-compatible backend mode. | | `valkey` | Valkey-compatible backend mode. |
For Valkey mode, the server reads `CHRONOSEAL_VALKEY_ADDR` and defaults to `127.0.0.1:6666` when it is not set. If the Valkey connection fails, the current implementation falls back to in-memory SQLite and logs a warning. For Valkey mode, the server reads `CHRONOSEAL_VALKEY_ADDR` (defaulting to `127.0.0.1:6666`) and establishes a thread-safe connection pool using `r2d2` and the `redis` client crate. It leverages native Valkey sets for session ID indexing and native key expiration for automatic session cleanup. If the Valkey connection fails, the server falls back to in-memory SQLite and logs a warning.
### Valkey / Redis Server Setup
To quickly run a local Valkey/Redis instance for testing or production:
```bash
# Option A: Start a local Valkey/Redis server on port 6666
valkey-server --port 6666 --bind 127.0.0.1
# Or
redis-server --port 6666 --bind 127.0.0.1
# Option B: Spin up via Docker
docker run -d --name chronoseal-valkey -p 6666:6379 valkey/valkey:latest
```
Configure ChronoSeal to use it:
```bash
export CHRONOSEAL_DB_TYPE=valkey
export CHRONOSEAL_VALKEY_ADDR=127.0.0.1:6666
```
If your Valkey or Redis server requires credentials or secure TLS:
* **Password Only**: `redis://:your_password@127.0.0.1:6666`
* **Username & Password**: `redis://your_username:your_password@127.0.0.1:6666`
* **Secure Connection (SSL/TLS)**: `rediss://your_username:your_password@secure-host.example.com:6379`
## Operations ## Operations
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "chronoseal-replay" name = "chronoseal-replay"
version = "0.1.0" version = "1.0.1"
edition = "2021" edition = "2021"
[dependencies] [dependencies]
+133 -46
View File
@@ -117,7 +117,11 @@ fn test_entropy() -> EntropyData {
} }
} }
fn do_handshake(client: &reqwest::blocking::Client, base_url: &str, sk: &SigningKey) -> Result<InitResponse> { fn do_handshake(
client: &reqwest::blocking::Client,
base_url: &str,
sk: &SigningKey,
) -> Result<InitResponse> {
let pk_hex = hex::encode(sk.verifying_key().to_bytes()); let pk_hex = hex::encode(sk.verifying_key().to_bytes());
let init_req = InitRequest { public_key: pk_hex }; let init_req = InitRequest { public_key: pk_hex };
let resp = client let resp = client
@@ -126,7 +130,10 @@ fn do_handshake(client: &reqwest::blocking::Client, base_url: &str, sk: &Signing
.send()?; .send()?;
if !resp.status().is_success() { 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()?; 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 mut failures = 0;
let scenarios = [ 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), ("stale_replay", run_stale_replay),
("invalid_signature", run_invalid_signature), ("invalid_signature", run_invalid_signature),
("invalid_vm_stack", run_invalid_vm_stack), ("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), ("drifted_timestamp", run_drifted_timestamp),
("concurrent_heartbeat", run_concurrent_heartbeat), ("concurrent_heartbeat", run_concurrent_heartbeat),
("rate_limit_trigger", run_rate_limit_trigger), ("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, 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 timestamp = current_time_ms();
let entropy = test_entropy(); let entropy = test_entropy();
@@ -215,13 +229,14 @@ fn run_valid_progression(client: &reqwest::blocking::Client, base_url: &str) ->
sign_request(&sk, &mut req)?; sign_request(&sk, &mut req)?;
let resp = client let resp = client.post(format!("{}/hb", base_url)).json(&req).send()?;
.post(format!("{}/hb", base_url))
.json(&req)
.send()?;
if !resp.status().is_success() { 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()?; 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) // 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_salt = hb_resp
let next_step = hb_resp.next_mutation_step.ok_or_else(|| anyhow!("Step {} missing next mutation step", step))?; .next_salt
let next_order = hb_resp.next_mutation_order_b64.ok_or_else(|| anyhow!("Step {} missing next mutation order", step))?; .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); 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 sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?; 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 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 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 candidate = shared::vm_extensions::apply_program_clone_with_rounds(
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); &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 { let mut req = HeartbeatRequest {
session_id: init.session_id.clone(), 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 resp2 = client.post(format!("{}/hb", base_url)).json(&req).send()?;
let hb2: HeartbeatResponse = resp2.json()?; let hb2: HeartbeatResponse = resp2.json()?;
if hb2.next_salt.is_some() { 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."); 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 sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?; 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 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 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 candidate = shared::vm_extensions::apply_program_clone_with_rounds(
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); &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 { let mut req = HeartbeatRequest {
session_id: init.session_id.clone(), 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 sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?; 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 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 candidate = shared::vm_extensions::apply_program_clone_with_rounds(
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); &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 { let mut req = HeartbeatRequest {
session_id: init.session_id.clone(), session_id: init.session_id.clone(),
@@ -375,12 +422,18 @@ fn run_invalid_vm_stack(client: &reqwest::blocking::Client, base_url: &str) -> R
Ok(()) 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 mut csprng = OsRng;
let sk = SigningKey::generate(&mut csprng); let sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?; 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 stack_state = shared::vm::execute(&opcodes);
let mut req = HeartbeatRequest { 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 sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?; 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 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 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 candidate = shared::vm_extensions::apply_program_clone_with_rounds(
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); &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 { let mut req = HeartbeatRequest {
session_id: init.session_id.clone(), 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 sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?; 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 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 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 candidate = shared::vm_extensions::apply_program_clone_with_rounds(
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); &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 { let mut req = HeartbeatRequest {
session_id: init.session_id.clone(), 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 client_clone = client.clone();
let req_clone = req.clone(); let req_clone = req.clone();
let url_clone = format!("{}/hb", base_url); let url_clone = format!("{}/hb", base_url);
let handle = std::thread::spawn(move || { let handle = std::thread::spawn(move || client_clone.post(&url_clone).json(&req_clone).send());
client_clone.post(&url_clone).json(&req_clone).send()
});
let resp2 = client.post(format!("{}/hb", base_url)).json(&req).send()?; let resp2 = client.post(format!("{}/hb", base_url)).json(&req).send()?;
let resp1_res = handle.join().map_err(|_| anyhow!("Thread panicked"))?; 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 // 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); let successes = (hb1.next_salt.is_some() as usize) + (hb2.next_salt.is_some() as usize);
if successes != 1 { 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)."); 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 sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?; 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 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 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 candidate = shared::vm_extensions::apply_program_clone_with_rounds(
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step); &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 { let mut req = HeartbeatRequest {
session_id: init.session_id.clone(), 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 { 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."); println!("Rate limiter correctly triggered.");
Ok(()) 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_content = std::fs::read_to_string(file_path)?;
let scenario: serde_json::Value = serde_json::from_str(&scenario_content)?; let scenario: serde_json::Value = serde_json::from_str(&scenario_content)?;
+1
View File
@@ -27,6 +27,7 @@ RestrictSUIDSGID=yes
LockPersonality=yes LockPersonality=yes
SystemCallArchitectures=native SystemCallArchitectures=native
ReadWritePaths=/run/chronoseal.pid ReadWritePaths=/run/chronoseal.pid
ReadWritePaths=/var/lib/chronoseal
# Logging # Logging
StandardOutput=journal StandardOutput=journal
+1 -1
View File
@@ -161,7 +161,7 @@ Content-Type: application/json
| `stack_state.ip` | number | yes | VM instruction pointer as an unsigned 16-bit value | | `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.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.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 | | `mutation_step` | number | yes | Mutation step currently expected by the server |
| `gene_commitment` | string | yes | Context-bound commitment produced by the WASM mutation preview | | `gene_commitment` | string | yes | Context-bound commitment produced by the WASM mutation preview |
| `signature` | string | yes | Ed25519 signature over the canonical payload | | `signature` | string | yes | Ed25519 signature over the canonical payload |
+2 -2
View File
@@ -379,7 +379,7 @@ Storage is abstracted by `DbPool`.
|---|---|---| |---|---|---|
| SQLite memory | `sqlite-in-memory` | default, process-local, ephemeral | | SQLite memory | `sqlite-in-memory` | default, process-local, ephemeral |
| SQLite disk | `sqlite-in-disk` | persisted SQLite file at `db_path` | | SQLite disk | `sqlite-in-disk` | persisted SQLite file at `db_path` |
| Valkey | `valkey` | Valkey-compatible session store | | Valkey | `valkey` | Valkey-compatible session store utilizing thread-safe connection pooling |
The storage layer must support: The storage layer must support:
@@ -389,7 +389,7 @@ The storage layer must support:
- delete expired sessions - delete expired sessions
- report statistics - report statistics
`valkey` mode reads `CHRONOSEAL_VALKEY_ADDR`, defaulting to `127.0.0.1:6666`. If connection setup fails, the current implementation logs a warning and falls back to in-memory SQLite. `valkey` mode reads `CHRONOSEAL_VALKEY_ADDR`, defaulting to `127.0.0.1:6666`. It establishes a connection pool using `r2d2` and the `redis` client crate. Session IDs are indexed using native Valkey sets (`sessions:ids`) to minimize overhead and avoid lock contention, while individual sessions are persisted with a native TTL (`SET ... EX`) matching their expiration times. If connection setup fails, it logs a warning and falls back to in-memory SQLite.
## Metrics and Observability ## Metrics and Observability
+51 -1
View File
@@ -178,13 +178,63 @@ sudo mkdir -p /var/lib/chronoseal
sudo chown -R chronoseal:chronoseal /var/lib/chronoseal sudo chown -R chronoseal:chronoseal /var/lib/chronoseal
``` ```
For Valkey: For Valkey / Redis:
ChronoSeal expects a running Valkey or Redis instance when `db_type` is set to `valkey`.
### 1. Installing Valkey or Redis
To install Valkey (the recommended open-source option) or Redis on Linux:
* **Valkey (Debian/Ubuntu)**:
```bash
sudo apt-get install -y valkey-server
```
* **Redis (Debian/Ubuntu)**:
```bash
sudo apt-get install -y redis-server
```
### 2. Local Setup and Startup
By default, ChronoSeal searches for Valkey/Redis on `127.0.0.1:6666`.
You can start a local instance manually:
```bash
# Start Valkey on port 6666
valkey-server --port 6666 --bind 127.0.0.1
# Or start Redis on port 6666
redis-server --port 6666 --bind 127.0.0.1
```
Or run it via Docker:
```bash
# Run Valkey container mapping host port 6666 to container port 6379
docker run -d --name chronoseal-valkey -p 6666:6379 valkey/valkey:latest
```
### 3. Service Configuration
Configure the environment variables to point ChronoSeal to your instance:
```bash ```bash
export CHRONOSEAL_DB_TYPE=valkey export CHRONOSEAL_DB_TYPE=valkey
export CHRONOSEAL_VALKEY_ADDR=127.0.0.1:6666 export CHRONOSEAL_VALKEY_ADDR=127.0.0.1:6666
``` ```
#### Providing Credentials & SSL/TLS
If your Valkey/Redis server requires authentication or secure TLS/SSL, include them directly in the `CHRONOSEAL_VALKEY_ADDR` connection URL:
* **Password Only**:
```bash
export CHRONOSEAL_VALKEY_ADDR=redis://:your_password@127.0.0.1:6666
```
* **Username & Password**:
```bash
export CHRONOSEAL_VALKEY_ADDR=redis://your_username:your_password@127.0.0.1:6666
```
* **Secure Connection (SSL/TLS)**: Use the `rediss://` scheme prefix:
```bash
export CHRONOSEAL_VALKEY_ADDR=rediss://your_username:your_password@secure-valkey-host.example.com:6379
```
If Valkey connection setup fails, the current implementation logs a warning and falls back to in-memory SQLite. If Valkey connection setup fails, the current implementation logs a warning and falls back to in-memory SQLite.
## systemd ## systemd
+181 -120
View File
@@ -2,14 +2,15 @@
ChronoSeal maintains a rigorous, security-first test suite focused on cryptographic correctness, deterministic server ↔ WASM parity, mutation engine integrity, replay resistance, tampering detection, behavioral validation, and storage reliability. ChronoSeal maintains a rigorous, security-first test suite focused on cryptographic correctness, deterministic server ↔ WASM parity, mutation engine integrity, replay resistance, tampering detection, behavioral validation, and storage reliability.
As of **v0.6.1**, the project contains **89 passing tests** across the server, WASM, and shared protocol crates. As of **v1.0.1**, the project contains **95 passing tests** across the server, WASM, and shared protocol crates.
| Crate | Tests | | Crate | Tests |
| ------------------- | -----: | | ---------------------------- | ------ |
| `chronoseal-server` | 30 | | `chronoseal-server` | 33 |
| `chronoseal-wasm` | 24 | | `chronoseal-wasm` | 24 |
| `shared` | 35 | | `shared` (unit) | 36 |
| **Total** | **89** | | `shared` (property-based) | 2 |
| **Total** | **95** |
--- ---
@@ -17,12 +18,12 @@ As of **v0.6.1**, the project contains **89 passing tests** across the server, W
ChronoSeal testing prioritizes: ChronoSeal testing prioritizes:
* **Security invariants** over raw coverage metrics - **Security invariants** over raw coverage metrics
* **Deterministic parity** between server and browser WASM runtimes - **Deterministic parity** between server and browser WASM runtimes
* **Negative-path testing** (tampering, replay, malformed input, edge cases) - **Negative-path testing** (tampering, replay, malformed input, edge cases)
* **Fuzz-style and randomized testing** for mutation logic - **Property-based and fuzz-style testing** for mutation logic and VM robustness
* **Performance regression detection** - **Performance regression detection**
* **Long-term protocol stability** - **Long-term protocol stability**
Particular emphasis is placed on ensuring that browser-side WASM execution produces identical results to server-side validation. Particular emphasis is placed on ensuring that browser-side WASM execution produces identical results to server-side validation.
@@ -34,27 +35,28 @@ Particular emphasis is placed on ensuring that browser-side WASM execution produ
Configuration tests verify: Configuration tests verify:
* Database backend selection - Database backend selection
* TOML configuration parsing - TOML configuration parsing
* Command-line override behavior - Command-line override behavior
* Default configuration values - Default configuration values
* Runtime initialization logic - Runtime initialization logic
Supported backends include: Supported backends include:
* `sqlite-in-memory` - `sqlite-in-memory`
* `sqlite-in-disk` - `sqlite-in-disk`
* `valkey` - `valkey` (redis-compatible via r2d2 connection pool)
Example tests: Example tests:
```text ```
test_apply_run_args_overrides_db_type test_apply_run_args_overrides_db_type
test_default_db_type_is_sqlite_in_memory test_default_db_type_is_sqlite_in_memory
test_toml_parses_db_type_kebab_case test_toml_parses_db_type_kebab_case
test_init_db_pool_sqlite_in_memory test_init_db_pool_sqlite_in_memory
test_init_db_pool_sqlite_in_disk test_init_db_pool_sqlite_in_disk
test_init_db_pool_valkey_compat_mode test_init_db_pool_valkey_compat_mode
test_valkey_store_operations
``` ```
--- ---
@@ -63,17 +65,17 @@ test_init_db_pool_valkey_compat_mode
Session tests validate: Session tests validate:
* Session creation - Session creation
* Public key validation - Public key validation
* Expiration handling - Expiration handling
* Replay attack prevention - Replay attack prevention
* Mutation step enforcement - Mutation step enforcement
* Commitment verification - Commitment verification
* Long-running deterministic parity - Long-running deterministic parity
Example tests: Example tests:
```text ```
test_create_session_rejects_invalid_public_key_length test_create_session_rejects_invalid_public_key_length
test_expired_session_is_rejected test_expired_session_is_rejected
test_replay_attack_is_rejected test_replay_attack_is_rejected
@@ -92,17 +94,17 @@ The Synthetic Gene Mutation Engine is one of the most security-critical componen
Testing focuses on: Testing focuses on:
* Deterministic server/client parity - Deterministic server/client parity
* Mutation order execution - Mutation order execution
* Gene state integrity - Gene state integrity
* Preview → Commit → Discard lifecycle - Preview → Commit → Discard lifecycle
* Randomized mutation programs - Randomized mutation programs
* Edge-case validation - Edge-case validation
* Performance regression detection - Performance regression detection
Example tests: Example tests:
```text ```
test_server_client_parity_across_random_orders test_server_client_parity_across_random_orders
test_generate_order_is_deterministic_for_seeded_rng test_generate_order_is_deterministic_for_seeded_rng
test_invalid_positions_wrap_deterministically test_invalid_positions_wrap_deterministically
@@ -117,15 +119,15 @@ test_performance_smoke_mutation_execution
Heartbeat validation tests verify: Heartbeat validation tests verify:
* Successful state advancement - Successful state advancement
* Silent rejection behavior - Silent rejection behavior
* Commitment validation - Commitment validation
* Rate limiting - Rate limiting
* Next-state mutation generation - Next-state mutation generation
Example tests: Example tests:
```text ```
test_handler_success_returns_next_mutation_fields test_handler_success_returns_next_mutation_fields
test_handler_tampered_commitment_is_silent_failure test_handler_tampered_commitment_is_silent_failure
test_handler_rate_limit_returns_no_mutation_data test_handler_rate_limit_returns_no_mutation_data
@@ -137,16 +139,16 @@ test_handler_rate_limit_returns_no_mutation_data
Behavioral validation tests verify: Behavioral validation tests verify:
* Minimum mouse activity - Minimum mouse activity
* Minimum movement distance - Minimum movement distance
* Pause detection - Pause detection
* Speed thresholds - Speed thresholds
* Optional activity requirements - Optional activity requirements
* Fingerprint-related validation paths - Fingerprint-related validation paths
Example tests: Example tests:
```text ```
test_validate_mouse_success test_validate_mouse_success
test_validate_mouse_insufficient_events test_validate_mouse_insufficient_events
test_validate_mouse_insufficient_distance test_validate_mouse_insufficient_distance
@@ -161,15 +163,28 @@ test_validate_mouse_require_activity_toggle
Storage tests verify: Storage tests verify:
* SQLite in-memory operation - SQLite in-memory operation
* SQLite disk-backed operation - SQLite disk-backed operation and pool concurrency
* Valkey compatibility mode - Valkey compatibility mode and r2d2 pool concurrency
* Session CRUD behavior - Session CRUD behavior including the `opcodes` field
* Expiration cleanup - Expiration cleanup
* Runtime statistics reporting - Runtime statistics reporting
These tests ensure storage implementations remain interchangeable without affecting protocol behavior. These tests ensure storage implementations remain interchangeable without affecting protocol behavior.
Example tests:
```
test_sqlite_pool_concurrency
test_valkey_pool_concurrency
test_valkey_store_operations
```
> **Note:** The concurrent write collision path in `update_session` (the `old_last_hash`
> optimistic concurrency guard) is not yet covered by an automated test. Two goroutines
> advancing the same chain simultaneously is a security-relevant race condition.
> A dedicated test is planned for v1.1.0 (see Future Improvements).
--- ---
## 7. VM Core ## 7. VM Core
@@ -178,28 +193,19 @@ The VM core is tested extensively across both WASM and shared crates.
Coverage includes: Coverage includes:
* ADD - ADD, SUB, MUL, XOR, AND, OR, NOT, HASH, ROT, PUSH
* SUB
* MUL
* XOR
* AND
* OR
* NOT
* HASH
* ROT
* PUSH
Edge cases include: Edge cases include:
* Stack underflow - Stack underflow
* Truncated instructions - Truncated instructions
* Unknown opcodes - Unknown opcodes
* Wrapping arithmetic - Wrapping arithmetic
* Invalid instruction streams - Invalid instruction streams
Example tests: Example tests:
```text ```
test_add test_add
test_add_wrapping test_add_wrapping
test_sub test_sub
@@ -215,16 +221,38 @@ test_rejects_truncated_instruction
--- ---
## 8. Property-Based Tests
ChronoSeal uses [`proptest`](https://github.com/proptest-rs/proptest) for property-based testing of core protocol invariants against arbitrary random input.
Tests live in `shared/tests/proptests.rs` and run as part of `cargo test --workspace`.
```
test_vm_execute_never_panics
test_gene_environment_roundtrip_never_panics
```
`test_vm_execute_never_panics` feeds arbitrary `Vec<u8>` byte sequences into the stack machine
and asserts that execution never panics and that the instruction pointer never exceeds the
program length. This guards against any future VM opcode handler introducing undefined
behaviour on malformed input.
`test_gene_environment_roundtrip_never_panics` feeds arbitrary bytes into
`gene::decode_environment` and asserts graceful failure rather than a panic, covering the full
space of malformed environment payloads a client could send.
---
# Server Test Coverage (`chronoseal-server`) # Server Test Coverage (`chronoseal-server`)
The server crate currently contains **30 tests** covering: The server crate currently contains **33 tests** covering:
* Configuration - Configuration
* Runtime initialization - Runtime initialization
* Session management - Session management
* Heartbeat validation - Heartbeat validation
* Rate limiting - Rate limiting
* Trust validation - Trust validation
The server tests focus heavily on protocol enforcement and security validation. The server tests focus heavily on protocol enforcement and security validation.
@@ -234,16 +262,16 @@ The server tests focus heavily on protocol enforcement and security validation.
The WASM crate currently contains **24 tests** covering: The WASM crate currently contains **24 tests** covering:
* VM execution - VM execution
* Browser-side mutation lifecycle - Browser-side mutation lifecycle
* Gene initialization - Gene initialization
* Mutation preview - Mutation preview
* Mutation commit/discard behavior - Mutation commit/discard behavior
* Deterministic parity with shared logic - Deterministic parity with shared logic
Example tests: Example tests:
```text ```
test_preview_commitment_matches_shared_engine test_preview_commitment_matches_shared_engine
test_commit_applies_preview test_commit_applies_preview
test_discard_preview_keeps_committed_state test_discard_preview_keeps_committed_state
@@ -256,13 +284,14 @@ These tests ensure browser-generated commitments remain consistent with server e
# Shared Crate Coverage (`shared`) # Shared Crate Coverage (`shared`)
The shared crate currently contains **35 tests** and represents the core protocol implementation used by both server and browser runtimes. The shared crate currently contains **36 unit tests** and **2 property-based tests**, representing
the core protocol implementation used by both server and browser runtimes.
Coverage includes: Coverage includes:
### Synthetic Gene Engine ### Synthetic Gene Engine
```text ```
test_new_state_with_default_size test_new_state_with_default_size
test_new_state_rejects_invalid_sizes test_new_state_rejects_invalid_sizes
test_commitment_changes_when_gene_or_environment_changes test_commitment_changes_when_gene_or_environment_changes
@@ -272,7 +301,7 @@ test_table_driven_randomized_environment_roundtrip
### Mutation Engine ### Mutation Engine
```text ```
test_opcode_insert test_opcode_insert
test_opcode_delete test_opcode_delete
test_opcode_mutate_point test_opcode_mutate_point
@@ -283,7 +312,7 @@ test_mutation_chain
### Validation & Hardening ### Validation & Hardening
```text ```
test_rejects_stack_underflow test_rejects_stack_underflow
test_rejects_truncated_instruction test_rejects_truncated_instruction
test_rejects_unknown_opcode test_rejects_unknown_opcode
@@ -292,7 +321,7 @@ test_zero_length_gene_is_rejected
### Deterministic Parity ### Deterministic Parity
```text ```
test_server_client_parity_across_random_orders test_server_client_parity_across_random_orders
test_generate_order_is_deterministic_for_seeded_rng test_generate_order_is_deterministic_for_seeded_rng
test_invalid_positions_wrap_deterministically test_invalid_positions_wrap_deterministically
@@ -300,24 +329,47 @@ test_invalid_positions_wrap_deterministically
### Fuzz & Regression Testing ### Fuzz & Regression Testing
```text ```
test_fuzz_style_random_program_bytes_do_not_diverge test_fuzz_style_random_program_bytes_do_not_diverge
test_performance_smoke_mutation_execution test_performance_smoke_mutation_execution
test_vm_instruction_budget_soft_cap
``` ```
### Property-Based Tests (`shared/tests/proptests.rs`)
```
test_vm_execute_never_panics
test_gene_environment_roundtrip_never_panics
```
---
# Tooling Crates
## `chronoseal-replay`
A standalone replay and audit tool for offline verification of recorded ChronoSeal session
chains. It is a developer and forensic utility, not a library, and currently carries no
automated tests. Integration tests against captured session fixtures are planned.
## `fuzz/`
Contains libFuzzer targets for deeper coverage of the VM and gene codec. Run separately
via `cargo +nightly fuzz run <target>` — not part of the standard `cargo test` suite.
--- ---
# Running the Test Suite # Running the Test Suite
Run the full workspace: Run the full workspace:
```bash ```
cargo test --workspace cargo test --workspace
``` ```
Run individual crates: Run individual crates:
```bash ```
cargo test -p shared cargo test -p shared
cargo test -p chronoseal-wasm cargo test -p chronoseal-wasm
cargo test -p chronoseal-server cargo test -p chronoseal-server
@@ -325,7 +377,7 @@ cargo test -p chronoseal-server
Display test output: Display test output:
```bash ```
cargo test -- --nocapture cargo test -- --nocapture
``` ```
@@ -333,18 +385,22 @@ cargo test -- --nocapture
# Critical Security Tests # Critical Security Tests
The following tests protect ChronoSeal's core protocol guarantees and should be treated as **release-blocking** if they fail: The following tests protect ChronoSeal's core protocol guarantees and should be treated as
**release-blocking** if they fail:
```text ```
test_mutation_commitment_tamper_is_rejected test_mutation_commitment_tamper_is_rejected
test_replay_attack_is_rejected test_replay_attack_is_rejected
test_handler_tampered_commitment_is_silent_failure test_handler_tampered_commitment_is_silent_failure
test_server_client_parity_across_random_orders test_server_client_parity_across_random_orders
test_deterministic_server_client_parity_across_many_heartbeats test_deterministic_server_client_parity_across_many_heartbeats
test_fuzz_style_random_program_bytes_do_not_diverge test_fuzz_style_random_program_bytes_do_not_diverge
test_vm_execute_never_panics
test_gene_environment_roundtrip_never_panics
``` ```
These tests directly validate resistance to replay attacks, protocol divergence, mutation tampering, and commitment forgery. These tests directly validate resistance to replay attacks, protocol divergence, mutation
tampering, commitment forgery, and VM panic on adversarial input.
--- ---
@@ -355,40 +411,45 @@ When adding new functionality:
1. Prefer placing protocol logic tests in `shared/` 1. Prefer placing protocol logic tests in `shared/`
2. Ensure server ↔ WASM parity is validated 2. Ensure server ↔ WASM parity is validated
3. Include negative-path test cases 3. Include negative-path test cases
4. Add randomized testing where appropriate 4. Add randomized or property-based testing where appropriate
5. Update this document when introducing major new categories 5. Update this document when introducing major new categories
--- ---
# Future Improvements # Future Improvements
Planned enhancements include: - **Concurrent chain write collision test** — verify that two simultaneous heartbeats for the
same session are handled correctly by the `old_last_hash` optimistic concurrency guard in
* Property-based testing using `proptest` `update_session` (security-critical, planned for v1.1.0)
* Browser-driven end-to-end integration tests - **Browser-driven end-to-end integration tests** — full Playwright or wasm-bindgen-test
* Valkey concurrency and failover testing harness exercising the complete init → heartbeat loop in a real browser environment
* Automated benchmark execution in CI - **Valkey failover testing** — verify graceful degradation and reconnection under r2d2 pool
* Expanded mutation-engine fuzzing exhaustion and server-side connection drops
* CI-enforced performance regression thresholds - **Automated benchmark execution in CI** — enforce performance regression thresholds for
mutation engine and hash chain operations
- **Expanded mutation-engine fuzzing** — additional libFuzzer targets for `vm_extensions`
opcodes introduced in v0.7.0
- **CI-enforced performance regression thresholds** — gate releases on measured latency bounds
--- ---
# Conclusion # Conclusion
ChronoSeal's testing strategy is centered on preserving deterministic behavior, cryptographic correctness, and protocol integrity. ChronoSeal's testing strategy is centered on preserving deterministic behavior, cryptographic
correctness, and protocol integrity.
The current suite of **89 tests** provides broad coverage across: The current suite of **95 tests** provides broad coverage across:
* Session security - Session security
* Heartbeat validation - Heartbeat validation
* Mutation engine correctness - Mutation engine correctness
* Deterministic server/WASM parity - Deterministic server/WASM parity
* Trust validation - Trust and behavioral validation
* Storage abstraction - Storage abstraction
* Replay resistance - Replay resistance
* Protocol hardening - Protocol hardening
- Property-based VM and gene codec robustness
Maintaining and expanding this test suite remains a core project priority as ChronoSeal evolves. Maintaining and expanding this test suite remains a core project priority as ChronoSeal evolves.
**Last Updated:** May 2026 (v0.6.1) **Last Updated:** May 2026 (v1.0.1)
+17
View File
@@ -106,6 +106,23 @@ Expected result:
- ChronoSeal does not claim complete prevention - ChronoSeal does not claim complete prevention
- additional application-level controls are required - additional application-level controls are required
## Attacker Classification Boundaries
### Protected
* **Commodity Scrapers:** Simple HTTP clients (`curl`, Python `requests`, Go HTTP clients) that cannot execute JavaScript or WebAssembly.
* **Simple Replay Attackers:** Intercepted heartbeat payloads cannot be reused because of the strict hash-chain sequencing and salt rotation.
* **Signature Forgers:** Heartbeats without the session's private key will fail Ed25519 verification.
### Partially Protected
* **Headless Automation (Puppeteer, Playwright):** Attackers must load the WASM runtime, execute the VM instructions, calculate gene mutations, and simulate realistic human mouse interactions. This significantly increases CPU and system memory overhead, reducing the scale of bot operations.
* **Stealth Automation Frameworks:** Advanced frameworks must maintain state sync across multiple heartbeat cycles, exposing them to timing detection.
### Unprotected
* **WASM Key Extraction:** A reverse engineer with full browser process control can extract the private key from WASM memory.
* **Malware Operators:** Keyloggers, screen scrapers, or memory dumpers operating at the OS level are outside the application trust boundary.
* **MITM Interceptors (without TLS):** Plaintext traffic can be intercepted. (TLS termination is assumed).
* **Insiders / Storage Tampering:** Attackers with direct write access to the SQLite database or Valkey instance can forge or hijack active session states.
## Attack Vectors and Mitigations ## Attack Vectors and Mitigations
### Replay ### Replay
Binary file not shown.

After

Width:  |  Height:  |  Size: 1.5 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 594 B

Binary file not shown.

After

Width:  |  Height:  |  Size: 15 KiB

+1
View File
@@ -0,0 +1 @@
<svg xmlns="http://www.w3.org/2000/svg" version="1.1" id="Layer_1" x="0px" y="0px" width="296.99997mm" viewBox="0 0 1122.5196 793.7008" enable-background="new 0 0 1254 1254" xml:space="preserve" height="210mm" sodipodi:docname="logo1.svg" inkscape:export-filename="logo1.png" inkscape:export-xdpi="96" inkscape:export-ydpi="96" inkscape:version="1.4.4 (dcaf3e7d9e, 2026-05-05)" xmlns:inkscape="http://www.inkscape.org/namespaces/inkscape" xmlns:sodipodi="http://sodipodi.sourceforge.net/DTD/sodipodi-0.dtd" xmlns:svg="http://www.w3.org/2000/svg"><metadata><rdf:RDF xmlns:rdf="http://www.w3.org/1999/02/22-rdf-syntax-ns#" xmlns:dc="http://purl.org/dc/elements/1.1/"><rdf:Description><dc:creator>RealFaviconGenerator</dc:creator><dc:source>https://realfavicongenerator.net</dc:source></rdf:Description></rdf:RDF></metadata><sodipodi:namedview id="namedview1" pagecolor="#ffffff" bordercolor="#000000" borderopacity="0.25" inkscape:showpageshadow="2" inkscape:pageopacity="0.0" inkscape:pagecheckerboard="0" inkscape:deskcolor="#d1d1d1" inkscape:document-units="mm" inkscape:zoom="1.1654266" inkscape:cx="561.16791" inkscape:cy="396.85039" inkscape:window-width="2048" inkscape:window-height="1205" inkscape:window-x="0" inkscape:window-y="0" inkscape:window-maximized="1" inkscape:current-layer="Layer_1"></sodipodi:namedview><defs id="defs44"></defs><path fill="none" opacity="1" stroke="none" d="m 791.89377,632.01366 c 1.04068,-1.03091 1.90903,-2.38956 3.14896,-3.04166 11.48425,-6.03955 22.98721,-12.04745 34.57742,-17.88034 8.52002,-4.28775 14.37558,-10.49713 14.23218,-20.46771 -0.14127,-9.82202 -6.25309,-15.51495 -14.66742,-19.5827 -61.76651,-29.85993 -123.47361,-59.84277 -185.21183,-89.76129 -10.48242,-5.07981 -20.81335,-10.5555 -31.58477,-14.93764 -19.71186,-8.01936 -39.87799,-8.40881 -59.18305,0.84148 -71.77527,34.39209 -143.33667,69.23047 -215.00567,103.84463 -8.40815,4.06091 -14.32166,9.84101 -14.39978,19.74008 -0.079,10.00498 5.52282,16.22797 14.08978,20.54743 12.3277,6.21567 24.51239,12.71484 37.34228,19.41369 7.51429,3.93381 14.42273,7.5896 21.37854,11.15274 31.68775,16.23208 63.40558,32.40548 95.07782,48.66778 14.943,7.6726 29.74964,15.6123 44.73053,23.2091 8.44464,4.2823 16.79547,9.1534 25.75525,11.9324 18.81989,5.8372 37.10712,3.9335 54.83185,-5.5578 25.51599,-13.6634 51.41541,-26.6141 77.2135,-39.7465 30.82071,-15.6894 61.7009,-31.26196 92.58359,-46.82907 1.56006,-0.78641 3.38776,-1.0419 5.09082,-1.54462 z" id="path2"></path><path fill="#4b4d51" opacity="1" stroke="none" d="m 791.89376,550.60089 c 1.04068,-1.03091 1.90903,-2.38956 3.14896,-3.04166 11.48425,-6.03955 22.98721,-12.04745 34.57742,-17.88034 8.52002,-4.28775 14.37558,-10.49713 14.23218,-20.46771 -0.14127,-9.82202 -6.25309,-15.51495 -14.66742,-19.5827 -61.76651,-29.85993 -123.47361,-59.84277 -185.21183,-89.76129 -10.48242,-5.07981 -20.81335,-10.5555 -31.58477,-14.93764 -19.71186,-8.01936 -39.87799,-8.40881 -59.18305,0.84148 -71.77522,34.39209 -143.33662,69.23047 -215.00562,103.84463 -8.4082,4.06091 -14.3217,9.84101 -14.3998,19.74008 -0.079,10.00498 5.5228,16.22797 14.0898,20.54743 12.3277,6.21567 24.5123,12.71484 37.3422,19.41369 7.5143,3.93381 14.4228,7.5896 21.3786,11.15274 31.6877,16.23212 63.4056,32.40552 95.0778,48.66782 14.943,7.67261 29.7496,15.61227 44.73052,23.20911 8.44464,4.28232 16.79547,9.15341 25.75525,11.93237 18.81989,5.83722 37.10712,3.93353 54.83185,-5.55777 25.51599,-13.66336 51.41541,-26.61407 77.2135,-39.74655 30.82071,-15.68933 61.7009,-31.26196 92.58359,-46.82907 1.56006,-0.78641 3.38776,-1.0419 5.09082,-1.54462 z" id="path45" style="fill:#e1e1e4;fill-opacity:1"></path><path fill="#4b4d51" opacity="1" stroke="none" d="m 791.89377,469.18814 c 1.04068,-1.03091 1.90903,-2.38956 3.14896,-3.04166 11.48425,-6.03955 22.98721,-12.04745 34.57742,-17.88034 8.52002,-4.28775 14.37558,-10.49713 14.23218,-20.46771 -0.14127,-9.82202 -6.25309,-15.51495 -14.66742,-19.5827 C 767.4184,378.3558 705.7113,348.37296 643.97308,318.45444 c -10.48242,-5.07981 -20.81335,-10.5555 -31.58477,-14.93764 -19.71186,-8.01936 -39.87799,-8.40881 -59.18305,0.84148 -71.77527,34.39209 -143.33667,69.23047 -215.00567,103.84463 -8.40815,4.06091 -14.32166,9.84101 -14.39978,19.74008 -0.079,10.00498 5.52282,16.22797 14.08978,20.54743 12.3277,6.21567 24.51239,12.71484 37.34228,19.41369 7.51429,3.93381 14.42273,7.5896 21.37854,11.15274 31.68775,16.23212 63.40558,32.40552 95.07782,48.66782 14.943,7.67261 29.74964,15.61227 44.73053,23.20911 8.44464,4.28232 16.79547,9.15341 25.75525,11.93237 18.81989,5.83722 37.10712,3.93353 54.83185,-5.55777 25.51599,-13.66336 51.41541,-26.61407 77.2135,-39.74655 30.82071,-15.68933 61.7009,-31.26196 92.58359,-46.82907 1.56006,-0.78641 3.38776,-1.0419 5.09082,-1.54462 z" id="path46" style="fill:#b0b2b8;fill-opacity:1"></path><path fill="#4b4d51" opacity="1" stroke="none" d="m 791.89377,387.77539 c 1.04068,-1.03091 1.90903,-2.38956 3.14896,-3.04166 11.48425,-6.03955 22.98721,-12.04745 34.57742,-17.88034 8.52002,-4.28775 14.37558,-10.49713 14.23218,-20Line truncated

After

Width:  |  Height:  |  Size: 8.1 KiB

+21
View File
@@ -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"
}
Binary file not shown.

After

Width:  |  Height:  |  Size: 1.6 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 7.0 KiB

+3 -4
View File
@@ -16,6 +16,7 @@ let minInterval = 12000;
let maxInterval = 25000; let maxInterval = 25000;
let pendingMutationStep = 0; let pendingMutationStep = 0;
let pendingMutationOrderB64 = ''; let pendingMutationOrderB64 = '';
let mutationRounds = 4;
export async function initHeartbeat() { export async function initHeartbeat() {
await init(); await init();
@@ -32,6 +33,7 @@ export async function initHeartbeat() {
} }
pendingMutationStep = initResp.mutation_step; pendingMutationStep = initResp.mutation_step;
pendingMutationOrderB64 = initResp.mutation_order_b64; pendingMutationOrderB64 = initResp.mutation_order_b64;
mutationRounds = initResp.mutation_rounds || 4;
lastTime = performance.now(); lastTime = performance.now();
scheduleNext(); scheduleNext();
} }
@@ -56,7 +58,7 @@ async function sendHeartbeat() {
const timestamp = Date.now(); const timestamp = Date.now();
const entropyData = { events: events.map(e => ({ x: e.x, y: e.y, t: e.t })) }; const entropyData = { events: events.map(e => ({ x: e.x, y: e.y, t: e.t })) };
const entropyJson = JSON.stringify(entropyData); const entropyJson = JSON.stringify(entropyData);
const geneCommitment = preview_gene_commitment(pendingMutationOrderB64); const geneCommitment = preview_gene_commitment(pendingMutationOrderB64, session, pendingMutationStep, mutationRounds);
if (!geneCommitment) { if (!geneCommitment) {
throw new Error('Unable to compute mutation commitment'); throw new Error('Unable to compute mutation commitment');
} }
@@ -75,7 +77,6 @@ async function sendHeartbeat() {
const sig = sign_message(msg); const sig = sign_message(msg);
if (!sig) { if (!sig) {
discard_gene_preview(); discard_gene_preview();
console.error('Keypair not initialised — skipping heartbeat');
return; return;
} }
const resp = await sendRequest('/hb', 'POST', { const resp = await sendRequest('/hb', 'POST', {
@@ -105,11 +106,9 @@ async function sendHeartbeat() {
pendingMutationOrderB64 = resp.next_mutation_order_b64; pendingMutationOrderB64 = resp.next_mutation_order_b64;
} else { } else {
discard_gene_preview(); discard_gene_preview();
console.warn('Heartbeat rejected');
} }
} catch (e) { } catch (e) {
discard_gene_preview(); discard_gene_preview();
console.error(e);
} finally { } finally {
scheduleNext(); scheduleNext();
} }
+1
View File
@@ -2,6 +2,7 @@
<html lang="en"> <html lang="en">
<head> <head>
<meta charset="UTF-8"> <meta charset="UTF-8">
<meta http-equiv="Content-Security-Policy" content="default-src 'self'; script-src 'self' 'wasm-unsafe-eval'; connect-src 'self'; style-src 'self' 'unsafe-inline'">
<title>Anti-Scraper Demo</title> <title>Anti-Scraper Demo</title>
</head> </head>
<body> <body>
+1 -3
View File
@@ -1,5 +1,3 @@
import { initHeartbeat } from './heartbeat.js'; import { initHeartbeat } from './heartbeat.js';
(async () => { initHeartbeat().catch(() => {});
await initHeartbeat();
})();
+1
View File
@@ -1,4 +1,5 @@
#!/bin/bash #!/bin/bash
set -euo pipefail
echo "Starting server with static frontend serving..." echo "Starting server with static frontend serving..."
cd ../server cd ../server
cargo run --release cargo run --release
+4 -2
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "chronoseal-server" name = "chronoseal-server"
version = "0.6.0" version = "1.0.1"
edition = "2021" edition = "2021"
[[bin]] [[bin]]
@@ -30,4 +30,6 @@ hex = "0.4"
base64 = "0.22" base64 = "0.22"
rand = "0.8" rand = "0.8"
ed25519-dalek = "2" ed25519-dalek = "2"
valkey = "0.0.0-alpha5" redis = { version = "0.29", features = ["r2d2"] }
dashmap = "6"
+10 -3
View File
@@ -1,6 +1,14 @@
use crate::session::AppState; use crate::session::AppState;
use std::sync::Arc; 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>) { pub async fn cleanup_loop(state: Arc<AppState>) {
loop { loop {
tokio::time::sleep(std::time::Duration::from_secs(60)).await; tokio::time::sleep(std::time::Duration::from_secs(60)).await;
@@ -12,11 +20,10 @@ pub async fn cleanup_loop(state: Arc<AppState>) {
} }
} }
// 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 window_secs = state.get_config().rate_limit_window_secs;
let mut rl = state.rate_limiter.lock().await; state.rate_limiter.evict_stale(window_secs);
rl.evict_stale(window_secs);
} }
} }
} }
+5 -2
View File
@@ -137,7 +137,9 @@ impl Config {
size: self.gene_size, 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 { return Err(ConfigError::InvalidMutationRounds {
rounds: self.mutation_rounds, rounds: self.mutation_rounds,
}); });
@@ -274,7 +276,8 @@ impl std::fmt::Display for ConfigError {
Self::InvalidMutationRounds { rounds } => { Self::InvalidMutationRounds { rounds } => {
write!( write!(
f, f,
"invalid mutation rounds {rounds}; expected 1..={}", "invalid mutation rounds {rounds}; expected {}..={}",
shared::constants::MIN_MUTATION_ROUNDS,
shared::constants::MAX_MUTATION_ROUNDS shared::constants::MAX_MUTATION_ROUNDS
) )
} }
+15 -2
View File
@@ -2,11 +2,16 @@ use ed25519_dalek::{Signature, VerifyingKey};
use shared::protocol::HeartbeatRequest; use shared::protocol::HeartbeatRequest;
use std::collections::BTreeMap; 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( pub fn canonical_signing_message(
req: &HeartbeatRequest, req: &HeartbeatRequest,
) -> Result<String, Box<dyn std::error::Error>> { ) -> 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(); let mut payload: BTreeMap<&str, serde_json::Value> = BTreeMap::new();
payload.insert("entropyData", serde_json::to_value(&req.entropy_data)?); payload.insert("entropyData", serde_json::to_value(&req.entropy_data)?);
payload.insert("fingerprint", serde_json::to_value(&req.fingerprint)?); payload.insert("fingerprint", serde_json::to_value(&req.fingerprint)?);
@@ -19,6 +24,14 @@ pub fn canonical_signing_message(
Ok(serde_json::to_string(&payload)?) 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 fn verify_signature(
pub_key_bytes: &[u8], pub_key_bytes: &[u8],
req: &HeartbeatRequest, req: &HeartbeatRequest,
+13
View File
@@ -25,12 +25,19 @@ pub enum SessionError {
#[error("Invalid gene configuration: {0}")] #[error("Invalid gene configuration: {0}")]
InvalidGeneConfiguration(String), InvalidGeneConfiguration(String),
#[error("Rate limited")]
RateLimited,
} }
impl IntoResponse for SessionError { impl IntoResponse for SessionError {
fn into_response(self) -> Response { fn into_response(self) -> Response {
let (status, error_message) = match self { let (status, error_message) = match self {
SessionError::InvalidPublicKeyLength => (StatusCode::BAD_REQUEST, self.to_string()), SessionError::InvalidPublicKeyLength => (StatusCode::BAD_REQUEST, self.to_string()),
SessionError::RateLimited => (
StatusCode::TOO_MANY_REQUESTS,
"Too many requests".to_string(),
),
_ => ( _ => (
StatusCode::INTERNAL_SERVER_ERROR, StatusCode::INTERNAL_SERVER_ERROR,
"Internal server error".to_string(), "Internal server error".to_string(),
@@ -86,4 +93,10 @@ pub enum VerificationError {
#[error("Gene state error: {0}")] #[error("Gene state error: {0}")]
GeneState(String), GeneState(String),
#[error("VM execution stack state mismatch")]
VmStackMismatch,
#[error("Concurrent state modification detected (CAS failed)")]
ConcurrentUpdate,
} }
+66 -3
View File
@@ -1,16 +1,79 @@
use shared::protocol::Fingerprint; 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,
/// 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>> { pub fn validate(fp: &Fingerprint) -> Result<(), Box<dyn std::error::Error>> {
let ar: f64 = fp.aspect_ratio.parse().map_err(|_| "ar")?; 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()); return Err("aspect ratio".into());
} }
let dpr: f64 = fp.device_pixel_ratio.parse().map_err(|_| "dpr")?; 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()); return Err("dpr".into());
} }
if fp.hardware_concurrency == 0 {
if fp.hardware_concurrency == 0 || fp.hardware_concurrency > MAX_HARDWARE_CONCURRENCY {
return Err("hw".into()); return Err("hw".into());
} }
Ok(()) Ok(())
} }
#[cfg(test)]
mod tests {
use super::*;
fn fingerprint(
aspect_ratio: impl Into<String>,
device_pixel_ratio: impl Into<String>,
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());
}
}
-3
View File
@@ -30,9 +30,6 @@ async fn main() {
async fn try_main() -> Result<(), Box<dyn std::error::Error>> { async fn try_main() -> Result<(), Box<dyn std::error::Error>> {
let cli = Cli::parse(); 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_filter = cli.globals.log.as_deref().unwrap_or("info");
let log_file = log_file_for_command(&cli); let log_file = log_file_for_command(&cli);
let _log_guard = init_logging(log_filter, log_file)?; let _log_guard = init_logging(log_filter, log_file)?;
+22
View File
@@ -9,3 +9,25 @@ pub async fn log_request(req: Request, next: Next) -> Response {
tracing::info!("{} {} -> {}", method, uri, response.status()); tracing::info!("{} {} -> {}", method, uri, response.status());
response 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
}
+34 -14
View File
@@ -1,34 +1,54 @@
use std::collections::HashMap; use dashmap::DashMap;
use std::time::Instant; use std::time::Instant;
/// 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<AppState>` without a `Mutex` wrapper.
pub struct RateLimiter { pub struct RateLimiter {
buckets: HashMap<String, (u32, Instant)>, /// Maps rate-limit keys to request counts and window start timestamps.
buckets: DashMap<String, (u32, Instant)>,
} }
impl RateLimiter { impl RateLimiter {
/// Creates a new, empty `RateLimiter`.
pub fn new() -> Self { pub fn new() -> Self {
Self { Self {
buckets: HashMap::new(), buckets: DashMap::new(),
} }
} }
pub fn check(&mut self, key: &str, limit: u32, window_secs: u64) -> bool { /// 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., 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(&self, key: &str, limit: u32, window_secs: u64) -> bool {
let now = Instant::now(); let now = Instant::now();
let entry = self.buckets.entry(key.to_string()).or_insert((0, now)); let mut entry = self.buckets.entry(key.to_string()).or_insert((0, now));
if now.duration_since(entry.1).as_secs() >= window_secs { let (count, ts) = entry.value_mut();
*entry = (1, now); if now.duration_since(*ts).as_secs() >= window_secs {
*count = 1;
*ts = now;
true true
} else if entry.0 >= limit { } else if *count >= limit {
false false
} else { } else {
entry.0 += 1; *count += 1;
true true
} }
} }
/// Remove entries whose rate-limit window has fully elapsed. /// Evicts expired rate-limit entries whose time windows have fully elapsed.
/// Call this periodically (e.g. from the cleanup loop) to bound memory usage. ///
pub fn evict_stale(&mut self, window_secs: u64) { /// 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(&self, window_secs: u64) {
let now = Instant::now(); let now = Instant::now();
self.buckets self.buckets
.retain(|_, (_, ts)| now.duration_since(*ts).as_secs() < window_secs); .retain(|_, (_, ts)| now.duration_since(*ts).as_secs() < window_secs);
@@ -43,7 +63,7 @@ mod tests {
#[test] #[test]
fn test_rate_limiter() { fn test_rate_limiter() {
let mut rl = RateLimiter::new(); let rl = RateLimiter::new();
// Limit of 2 requests per 1 second window // Limit of 2 requests per 1 second window
assert!(rl.check("user1", 2, 1)); assert!(rl.check("user1", 2, 1));
assert!(rl.check("user1", 2, 1)); assert!(rl.check("user1", 2, 1));
@@ -57,7 +77,7 @@ mod tests {
#[test] #[test]
fn test_rate_limiter_eviction() { fn test_rate_limiter_eviction() {
let mut rl = RateLimiter::new(); let rl = RateLimiter::new();
assert!(rl.check("user1", 1, 1)); assert!(rl.check("user1", 1, 1));
assert_eq!(rl.buckets.len(), 1); assert_eq!(rl.buckets.len(), 1);
+96 -11
View File
@@ -7,15 +7,52 @@ pub async fn handler(
State(state): State<Arc<AppState>>, State(state): State<Arc<AppState>>,
Json(payload): Json<HeartbeatRequest>, Json(payload): Json<HeartbeatRequest>,
) -> (StatusCode, Json<HeartbeatResponse>) { ) -> (StatusCode, Json<HeartbeatResponse>) {
// Rate limiting let start_http = std::time::Instant::now();
state
.heartbeats_total
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
// 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 (limit, window_secs) = {
let cfg = state.get_config(); let cfg = state.get_config();
(cfg.rate_limit_count, cfg.rate_limit_window_secs) (cfg.rate_limit_count, cfg.rate_limit_window_secs)
}; };
let mut rl = state.rate_limiter.lock().await; if !state
if !rl.check(&payload.session_id, limit, window_secs) { .rate_limiter
.check(&payload.session_id, limit, window_secs)
{
tracing::debug!("Rate limit hit: {}", payload.session_id); 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 ( return (
StatusCode::OK, StatusCode::OK,
Json(HeartbeatResponse { Json(HeartbeatResponse {
@@ -29,7 +66,17 @@ pub async fn handler(
} }
let config = state.get_config(); 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) => ( Ok(result) => (
StatusCode::OK, StatusCode::OK,
Json(HeartbeatResponse { Json(HeartbeatResponse {
@@ -41,6 +88,24 @@ pub async fn handler(
), ),
Err(e) => { Err(e) => {
tracing::warn!("Heartbeat failed for {}: {}", payload.session_id, 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, StatusCode::OK,
Json(HeartbeatResponse { Json(HeartbeatResponse {
@@ -51,7 +116,17 @@ 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)] #[cfg(test)]
@@ -59,7 +134,7 @@ mod tests {
use super::*; use super::*;
use axum::{extract::State, Json}; use axum::{extract::State, Json};
use ed25519_dalek::{Signer, SigningKey}; use ed25519_dalek::{Signer, SigningKey};
use shared::protocol::{EntropyData, Fingerprint, InitResponse, MouseEvent, StackState}; use shared::protocol::{EntropyData, Fingerprint, InitResponse, MouseEvent};
use std::path::Path; use std::path::Path;
fn test_config() -> crate::config::Config { fn test_config() -> crate::config::Config {
@@ -107,10 +182,12 @@ mod tests {
}, },
], ],
}; };
let stack_state = StackState { let program_bytes = base64::Engine::decode(
stack: vec![9, 10, 11], &base64::engine::general_purpose::STANDARD,
ip: 2, &init.opcodes_b64,
}; )
.unwrap();
let stack_state = shared::vm::execute(&program_bytes);
let order = let order =
shared::vm_extensions::decode_order_b64(mutation_step, mutation_order_b64).unwrap(); shared::vm_extensions::decode_order_b64(mutation_step, mutation_order_b64).unwrap();
@@ -147,8 +224,16 @@ mod tests {
let pool = crate::storage::init_pool(Path::new(":memory:")).unwrap(); let pool = crate::storage::init_pool(Path::new(":memory:")).unwrap();
let state = Arc::new(AppState { let state = Arc::new(AppState {
db_pool: pool.clone(), 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()), 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(); let mut rng = rand::thread_rng();
+31 -1
View File
@@ -8,7 +8,37 @@ pub async fn handler(
State(state): State<Arc<AppState>>, State(state): State<Arc<AppState>>,
Json(payload): Json<InitRequest>, Json(payload): Json<InitRequest>,
) -> Result<Json<InitResponse>, SessionError> { ) -> Result<Json<InitResponse>, SessionError> {
let start_http = std::time::Instant::now();
let config = state.get_config(); let config = state.get_config();
let resp = crate::session::create_session(&state.db_pool, &config, &payload.public_key)?;
// 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);
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)) Ok(Json(resp))
} }
+128 -20
View File
@@ -5,17 +5,19 @@ use crate::{
routes, session, routes, session,
storage::{self, StoreStats}, 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 serde::Serialize;
use std::{ use std::{
fs, fs,
io::{Read, Write}, io::{Read, Write},
net::{SocketAddr, TcpStream}, net::{IpAddr, SocketAddr, TcpStream},
path::Path, path::Path,
sync::Arc, sync::Arc,
time::Duration, time::Duration,
}; };
use tokio::sync::{Mutex, Notify}; use tokio::sync::Notify;
use tracing::{error, info, warn}; use tracing::{error, info, warn};
#[derive(Debug, Serialize)] #[derive(Debug, Serialize)]
@@ -144,8 +146,16 @@ pub async fn run_daemon(config: Config) -> Result<(), Box<dyn std::error::Error>
let db_pool = init_db_pool(&config)?; let db_pool = init_db_pool(&config)?;
let state = Arc::new(session::AppState { let state = Arc::new(session::AppState {
db_pool, db_pool,
rate_limiter: Mutex::new(RateLimiter::new()), rate_limiter: RateLimiter::new(),
config: std::sync::RwLock::new(config.clone()), 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(); let bg_state = state.clone();
@@ -162,7 +172,11 @@ pub async fn run_daemon(config: Config) -> Result<(), Box<dyn std::error::Error>
tower_http::services::ServeDir::new(&config.frontend_dir), tower_http::services::ServeDir::new(&config.frontend_dir),
) )
.layer(tower_http::cors::CorsLayer::permissive()) .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::middleware::from_fn(crate::middleware::log_request))
.layer(axum::extract::DefaultBodyLimit::max(64 * 1024)) // 64 KiB
.with_state(state.clone()); .with_state(state.clone());
let addr: SocketAddr = config.bind.parse()?; let addr: SocketAddr = config.bind.parse()?;
@@ -170,9 +184,12 @@ pub async fn run_daemon(config: Config) -> Result<(), Box<dyn std::error::Error>
info!(bind = %config.bind, "chronoseal daemon started"); info!(bind = %config.bind, "chronoseal daemon started");
let shutdown = signal_task(state.clone()); let shutdown = signal_task(state.clone());
let result = axum::serve(listener, app) let result = axum::serve(
.with_graceful_shutdown(shutdown) listener,
.await; app.into_make_service_with_connect_info::<SocketAddr>(),
)
.with_graceful_shutdown(shutdown)
.await;
remove_pid_file(&config.pid_file); remove_pid_file(&config.pid_file);
result?; result?;
@@ -269,8 +286,12 @@ async fn health_handler() -> impl IntoResponse {
} }
async fn stats_handler( async fn stats_handler(
ConnectInfo(addr): ConnectInfo<SocketAddr>,
axum::extract::State(state): axum::extract::State<Arc<session::AppState>>, axum::extract::State(state): axum::extract::State<Arc<session::AppState>>,
) -> Result<Json<StoreStats>, (StatusCode, String)> { ) -> Result<Json<StoreStats>, (StatusCode, String)> {
if !is_loopback(addr.ip()) {
return Err((StatusCode::FORBIDDEN, "Forbidden".to_string()));
}
state state
.db_pool .db_pool
.stats() .stats()
@@ -279,18 +300,99 @@ async fn stats_handler(
} }
async fn metrics_handler( async fn metrics_handler(
ConnectInfo(addr): ConnectInfo<SocketAddr>,
axum::extract::State(state): axum::extract::State<Arc<session::AppState>>, axum::extract::State(state): axum::extract::State<Arc<session::AppState>>,
) -> Result<String, (StatusCode, String)> { ) -> Result<String, (StatusCode, String)> {
state if !is_loopback(addr.ip()) {
return Err((StatusCode::FORBIDDEN, "Forbidden".to_string()));
}
let stats = state
.db_pool .db_pool
.stats() .stats()
.map(|stats| { .map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?;
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", let heartbeats = state
stats.sessions, stats.expired_sessions, stats.max_chain_length .heartbeats_total
) .load(std::sync::atomic::Ordering::Relaxed);
}) let ver_failures = state
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string())) .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
))
}
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<session::AppState>) { async fn signal_task(state: Arc<session::AppState>) {
@@ -458,10 +560,16 @@ mod tests {
fn test_init_db_pool_valkey_compat_mode() { fn test_init_db_pool_valkey_compat_mode() {
let mut config = base_config(); let mut config = base_config();
config.db_type = crate::config::DbType::Valkey; config.db_type = crate::config::DbType::Valkey;
let pool = init_db_pool(&config).unwrap(); match init_db_pool(&config) {
let stats = pool.stats().unwrap(); Ok(pool) => {
assert_eq!(stats.sessions, 0); let stats = pool.stats().unwrap();
assert_eq!(stats.expired_sessions, 0); assert_eq!(stats.sessions, 0);
assert_eq!(stats.max_chain_length, 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
}
}
} }
} }
+35 -11
View File
@@ -1,7 +1,15 @@
pub struct AppState { pub struct AppState {
pub db_pool: crate::storage::DbPool, pub db_pool: crate::storage::DbPool,
pub rate_limiter: tokio::sync::Mutex<crate::ratelimit::RateLimiter>, pub rate_limiter: crate::ratelimit::RateLimiter,
pub config: std::sync::RwLock<crate::config::Config>, 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 { impl AppState {
@@ -68,6 +76,7 @@ pub fn create_session(
environment: environment_blob, environment: environment_blob,
pending_mutation: initial_mutation.program, pending_mutation: initial_mutation.program,
pending_mutation_step: initial_mutation.step, pending_mutation_step: initial_mutation.step,
opcodes,
}; };
db.insert_session(&record) db.insert_session(&record)
@@ -84,6 +93,7 @@ pub fn create_session(
gene_size: config.gene_size as u32, gene_size: config.gene_size as u32,
mutation_step: initial_mutation.step, mutation_step: initial_mutation.step,
mutation_order_b64: initial_mutation_b64, mutation_order_b64: initial_mutation_b64,
mutation_rounds: config.mutation_rounds,
}) })
} }
@@ -150,6 +160,12 @@ pub fn verify_heartbeat(
fingerprint::validate(&req.fingerprint) fingerprint::validate(&req.fingerprint)
.map_err(|e| crate::errors::VerificationError::FingerprintFailed(e.to_string()))?; .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 // 6. Compute new hash
let new_hash = shared::hashing::next_chain_hash( let new_hash = shared::hashing::next_chain_hash(
&prev_hash_bytes, &prev_hash_bytes,
@@ -182,9 +198,16 @@ pub fn verify_heartbeat(
environment: next_environment_blob, environment: next_environment_blob,
pending_mutation: next_mutation.program, pending_mutation: next_mutation.program,
pending_mutation_step: next_step, pending_mutation_step: next_step,
opcodes: session.opcodes,
}; };
db.update_session(&update_record) db.update_session(&update_record, &session.last_hash)
.map_err(|e| crate::errors::VerificationError::Storage(e.to_string()))?; .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 { Ok(HeartbeatVerificationResult {
next_salt_hex, next_salt_hex,
@@ -210,6 +233,7 @@ mod tests {
pending_mutation_step: u64, pending_mutation_step: u64,
pending_mutation_order_b64: String, pending_mutation_order_b64: String,
committed_gene_state: GeneState, committed_gene_state: GeneState,
opcodes_b64: String,
} }
fn test_config() -> crate::config::Config { fn test_config() -> crate::config::Config {
@@ -247,13 +271,6 @@ mod tests {
} }
} }
fn test_stack() -> StackState {
StackState {
stack: vec![42, 7, 99],
ip: 3,
}
}
fn test_fingerprint() -> Fingerprint { fn test_fingerprint() -> Fingerprint {
Fingerprint { Fingerprint {
aspect_ratio: "1.77".to_string(), aspect_ratio: "1.77".to_string(),
@@ -288,6 +305,7 @@ mod tests {
pending_mutation_step: init.mutation_step, pending_mutation_step: init.mutation_step,
pending_mutation_order_b64: init.mutation_order_b64.clone(), pending_mutation_order_b64: init.mutation_order_b64.clone(),
committed_gene_state: gene::new_state(init.gene_size as usize).unwrap(), committed_gene_state: gene::new_state(init.gene_size as usize).unwrap(),
opcodes_b64: init.opcodes_b64.clone(),
} }
} }
@@ -304,7 +322,13 @@ mod tests {
vm_extensions::apply_program_clone(&client.committed_gene_state, &order.program) vm_extensions::apply_program_clone(&client.committed_gene_state, &order.program)
.unwrap(); .unwrap();
let entropy = test_entropy(); 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 { let mut req = HeartbeatRequest {
session_id: client.session_id.clone(), session_id: client.session_id.clone(),
+370 -71
View File
@@ -1,9 +1,8 @@
use crate::config::Config; use crate::config::Config;
use redis::Commands;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::path::Path; use std::path::Path;
use std::sync::{Arc, Mutex};
use std::time::{SystemTime, UNIX_EPOCH}; use std::time::{SystemTime, UNIX_EPOCH};
use valkey::Client as ValkeyClient;
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StoreStats { pub struct StoreStats {
@@ -20,7 +19,7 @@ pub enum DbPool {
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct ValkeyStore { pub struct ValkeyStore {
client: Arc<Mutex<ValkeyClient>>, pool: r2d2::Pool<redis::Client>,
index_key: String, index_key: String,
} }
@@ -38,6 +37,7 @@ pub struct SessionRecord {
pub environment: Vec<u8>, pub environment: Vec<u8>,
pub pending_mutation: Vec<u8>, pub pending_mutation: Vec<u8>,
pub pending_mutation_step: u64, pub pending_mutation_step: u64,
pub opcodes: Vec<u8>,
} }
impl DbPool { impl DbPool {
@@ -54,19 +54,18 @@ impl DbPool {
crate::config::DbType::Valkey => { crate::config::DbType::Valkey => {
let addr = std::env::var("CHRONOSEAL_VALKEY_ADDR") let addr = std::env::var("CHRONOSEAL_VALKEY_ADDR")
.unwrap_or_else(|_| "127.0.0.1:6666".to_string()); .unwrap_or_else(|_| "127.0.0.1:6666".to_string());
match ValkeyClient::connect(addr) { let connection_string =
Ok(client) => Ok(DbPool::Valkey(ValkeyStore { if addr.starts_with("redis://") || addr.starts_with("rediss://") {
client: Arc::new(Mutex::new(client)), addr.clone()
index_key: "sessions:ids".to_string(), } else {
})), format!("redis://{}", addr)
Err(err) => { };
tracing::warn!( let client = redis::Client::open(connection_string)?;
"valkey connection failed, falling back to sqlite-in-memory: {err}" let pool = r2d2::Pool::builder().build(client)?;
); Ok(DbPool::Valkey(ValkeyStore {
let pool = init_sqlite_pool(Path::new(":memory:"))?; pool,
Ok(DbPool::Sqlite(pool)) index_key: "sessions:ids".to_string(),
} }))
}
} }
} }
} }
@@ -79,8 +78,8 @@ impl DbPool {
"INSERT INTO sessions ( "INSERT INTO sessions (
session_id, public_key, salt, last_hash, chain_length, session_id, public_key, salt, last_hash, chain_length,
created_at, last_seen, expires_at, gene, environment, created_at, last_seen, expires_at, gene, environment,
pending_mutation, pending_mutation_step pending_mutation, pending_mutation_step, opcodes
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)", ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
)?; )?;
stmt.execute(rusqlite::params![ stmt.execute(rusqlite::params![
record.session_id, record.session_id,
@@ -95,6 +94,7 @@ impl DbPool {
&record.environment, &record.environment,
&record.pending_mutation, &record.pending_mutation,
record.pending_mutation_step, record.pending_mutation_step,
&record.opcodes,
])?; ])?;
Ok(()) Ok(())
} }
@@ -110,7 +110,7 @@ impl DbPool {
DbPool::Sqlite(pool) => { DbPool::Sqlite(pool) => {
let conn = pool.get()?; let conn = pool.get()?;
let mut stmt = conn.prepare( 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", FROM sessions WHERE session_id = ?1",
)?; )?;
let row = stmt.query_row([session_id], |row| { let row = stmt.query_row([session_id], |row| {
@@ -127,6 +127,7 @@ impl DbPool {
environment: row.get(9)?, environment: row.get(9)?,
pending_mutation: row.get(10)?, pending_mutation: row.get(10)?,
pending_mutation_step: row.get(11)?, pending_mutation_step: row.get(11)?,
opcodes: row.get(12)?,
}) })
}); });
match row { match row {
@@ -139,11 +140,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 { match self {
DbPool::Sqlite(pool) => { DbPool::Sqlite(pool) => {
let conn = pool.get()?; let conn = pool.get()?;
conn.execute( let rows = conn.execute(
"UPDATE sessions SET "UPDATE sessions SET
public_key=?1, public_key=?1,
salt=?2, salt=?2,
@@ -155,8 +160,9 @@ impl DbPool {
gene=?8, gene=?8,
environment=?9, environment=?9,
pending_mutation=?10, pending_mutation=?10,
pending_mutation_step=?11 pending_mutation_step=?11,
WHERE session_id=?12", opcodes=?12
WHERE session_id=?13 AND last_hash=?14",
rusqlite::params![ rusqlite::params![
&record.public_key, &record.public_key,
&record.salt, &record.salt,
@@ -169,12 +175,20 @@ impl DbPool {
&record.environment, &record.environment,
&record.pending_mutation, &record.pending_mutation,
record.pending_mutation_step, record.pending_mutation_step,
&record.opcodes,
&record.session_id, &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(()) Ok(())
} }
DbPool::Valkey(store) => store.insert_session(record), DbPool::Valkey(store) => store.update_session_cas(record, old_last_hash),
} }
} }
@@ -237,6 +251,10 @@ fn init_sqlite_pool(
} }
r2d2_sqlite::SqliteConnectionManager::file(path) 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 pool = r2d2::Pool::new(manager)?;
let conn = pool.get()?; let conn = pool.get()?;
init_schema(&conn)?; init_schema(&conn)?;
@@ -257,7 +275,8 @@ fn init_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
gene BLOB NOT NULL DEFAULT X'', gene BLOB NOT NULL DEFAULT X'',
environment BLOB NOT NULL DEFAULT X'', environment BLOB NOT NULL DEFAULT X'',
pending_mutation 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( ensure_column(
@@ -280,6 +299,11 @@ fn init_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
"pending_mutation_step", "pending_mutation_step",
"ALTER TABLE sessions ADD COLUMN pending_mutation_step INTEGER NOT NULL DEFAULT 0", "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( conn.execute_batch(
"CREATE INDEX IF NOT EXISTS idx_sessions_expires_at ON sessions(expires_at);", "CREATE INDEX IF NOT EXISTS idx_sessions_expires_at ON sessions(expires_at);",
)?; )?;
@@ -313,67 +337,136 @@ impl ValkeyStore {
&self, &self,
session_id: &str, session_id: &str,
) -> Result<Option<SessionRecord>, Box<dyn std::error::Error>> { ) -> Result<Option<SessionRecord>, Box<dyn std::error::Error>> {
let mut client = self.client.lock().unwrap(); let mut conn = self.pool.get()?;
if let Some(payload) = client.get(&self.session_key(session_id))? { let key = self.session_key(session_id);
let record = serde_json::from_str(&payload)?; let payload: Option<String> = conn.get(&key)?;
Ok(Some(record)) match payload {
} else { Some(p) => Ok(serde_json::from_str(&p)?),
Ok(None) None => Ok(None),
} }
} }
fn insert_session(&self, record: &SessionRecord) -> Result<(), Box<dyn std::error::Error>> { 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)?; let value = serde_json::to_string(record)?;
client.set(&self.session_key(&record.session_id), &value)?; let now = current_time_ms();
let existing = client.get(&self.index_key)?; let ttl_seconds = (record.expires_at.saturating_sub(now) / 1000).max(1);
let mut ids = existing.unwrap_or_default();
if !ids.split('\n').any(|id| id == record.session_id) { redis::pipe()
if !ids.is_empty() { .atomic()
ids.push('\n'); .cmd("SET")
} .arg(&key)
ids.push_str(&record.session_id); .arg(&value)
client.set(&self.index_key, &ids)?; .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(()) Ok(())
} }
fn purge_expired_sessions(&self) -> Result<(), Box<dyn std::error::Error>> { fn update_session_cas(
let mut client = self.client.lock().unwrap(); &self,
let ids = client.get(&self.index_key)?.unwrap_or_default(); record: &SessionRecord,
let now = current_time_ms(); old_last_hash: &[u8],
let mut remaining: Vec<String> = Vec::new(); ) -> Result<(), Box<dyn std::error::Error>> {
for id in ids.split('\n').filter(|id| !id.is_empty()) { let mut conn = self.pool.get()?;
if let Some(payload) = client.get(&self.session_key(id))? { let key = self.session_key(&record.session_id);
if let Ok(record) = serde_json::from_str::<SessionRecord>(&payload) {
if record.expires_at > now { // Watch key for concurrent modification
remaining.push(id.to_string()); 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(()) Ok(())
} }
fn stats(&self) -> Result<StoreStats, Box<dyn std::error::Error>> { 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 now = current_time_ms();
let mut sessions = 0; let mut conn = self.pool.get()?;
let mut expired_sessions = 0; let sessions: u64 = conn.zcard(&self.index_key)?;
let mut max_chain_length = 0; let expired_sessions: u64 = conn.zcount(&self.index_key, 0, now)?;
for id in ids.split('\n').filter(|id| !id.is_empty()) {
if let Some(payload) = client.get(&self.session_key(id))? { let max_chain_length_res: Vec<(String, u64)> =
if let Ok(record) = serde_json::from_str::<SessionRecord>(&payload) { conn.zrevrange_withscores("sessions:chain_lengths", 0, 0)?;
sessions += 1; let max_chain_length = max_chain_length_res
if record.expires_at < now { .first()
expired_sessions += 1; .map(|(_, score)| *score)
} .unwrap_or(0);
max_chain_length = max_chain_length.max(record.chain_length);
}
}
}
Ok(StoreStats { Ok(StoreStats {
sessions, sessions,
expired_sessions, expired_sessions,
@@ -385,6 +478,212 @@ impl ValkeyStore {
pub fn current_time_ms() -> u64 { pub fn current_time_ms() -> u64 {
SystemTime::now() SystemTime::now()
.duration_since(UNIX_EPOCH) .duration_since(UNIX_EPOCH)
.unwrap() .expect("system clock is before UNIX epoch; check system time")
.as_millis() as u64 .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::sync::Arc;
use std::thread;
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::sync::Arc;
use std::thread;
#[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);
}
}
+8
View File
@@ -1,6 +1,14 @@
use crate::config::Config; use crate::config::Config;
use shared::protocol::EntropyData; 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( pub fn validate_mouse(
data: &EntropyData, data: &EntropyData,
config: &Config, config: &Config,
+17
View File
@@ -4,6 +4,13 @@ use shared::{
vm_extensions::{self, ExecutionTrace, MutationError, MutationOrder}, 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> { pub fn generate_random_program(len_range: std::ops::RangeInclusive<usize>) -> Vec<u8> {
let mut rng = rand::thread_rng(); let mut rng = rand::thread_rng();
let count = rng.gen_range(len_range); let count = rng.gen_range(len_range);
@@ -47,6 +54,11 @@ pub fn generate_random_program(len_range: std::ops::RangeInclusive<usize>) -> Ve
ops 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)] #[allow(dead_code)]
pub fn execute_mutation_program( pub fn execute_mutation_program(
state: &mut GeneState, state: &mut GeneState,
@@ -55,6 +67,11 @@ pub fn execute_mutation_program(
vm_extensions::execute_program(state, 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)] #[allow(dead_code)]
pub fn execute_mutation_order( pub fn execute_mutation_order(
state: &mut GeneState, state: &mut GeneState,
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "shared" name = "shared"
version = "0.6.0" version = "1.0.1"
edition = "2021" edition = "2021"
[dependencies] [dependencies]
+1
View File
@@ -10,3 +10,4 @@ pub const MAX_MUTATION_ROUNDS: u8 = 10;
pub const MAX_MUTATION_INSTRUCTION_BUDGET: usize = 2048; pub const MAX_MUTATION_INSTRUCTION_BUDGET: usize = 2048;
pub const HASH_OPCODE_INSTRUCTION_COST: usize = 16; pub const HASH_OPCODE_INSTRUCTION_COST: usize = 16;
pub const SOFT_CAP_DURATION_MS: u128 = 50; pub const SOFT_CAP_DURATION_MS: u128 = 50;
pub const MAX_STACK_DEPTH: usize = 64;
+69
View File
@@ -55,6 +55,10 @@ impl std::fmt::Display for GeneError {
impl std::error::Error for GeneError {} impl std::error::Error for GeneError {}
/// Creates a new, blank `GeneState` with the specified gene buffer size.
///
/// # Arguments
/// * `gene_size` - The length of the gene byte buffer. Must be within `1..=MAX_GENE_SIZE`.
pub fn new_state(gene_size: usize) -> Result<GeneState, GeneError> { pub fn new_state(gene_size: usize) -> Result<GeneState, GeneError> {
if !(1..=MAX_GENE_SIZE).contains(&gene_size) { if !(1..=MAX_GENE_SIZE).contains(&gene_size) {
return Err(GeneError::InvalidGeneSize { size: gene_size }); return Err(GeneError::InvalidGeneSize { size: gene_size });
@@ -65,6 +69,7 @@ pub fn new_state(gene_size: usize) -> Result<GeneState, GeneError> {
}) })
} }
/// Creates a new `GeneState` with the default gene buffer size (`DEFAULT_GENE_SIZE`).
pub fn default_state() -> GeneState { pub fn default_state() -> GeneState {
GeneState { GeneState {
gene: vec![0; DEFAULT_GENE_SIZE], gene: vec![0; DEFAULT_GENE_SIZE],
@@ -72,6 +77,10 @@ pub fn default_state() -> GeneState {
} }
} }
/// Validates the structural invariants of the given `GeneState`.
///
/// Ensures the gene size is within valid bounds and the environment records are
/// properly sorted, non-empty, and free of duplicates.
pub fn validate_state(state: &GeneState) -> Result<(), GeneError> { pub fn validate_state(state: &GeneState) -> Result<(), GeneError> {
if !(1..=MAX_GENE_SIZE).contains(&state.gene.len()) { if !(1..=MAX_GENE_SIZE).contains(&state.gene.len()) {
return Err(GeneError::InvalidGeneSize { return Err(GeneError::InvalidGeneSize {
@@ -81,6 +90,13 @@ pub fn validate_state(state: &GeneState) -> Result<(), GeneError> {
validate_environment(&state.environment) validate_environment(&state.environment)
} }
/// Retrieves the quantity associated with a specific environment symbol.
///
/// Performs a binary search over the sorted environment records. Returns 0 if the symbol is missing.
///
/// # Arguments
/// * `state` - The gene state to query.
/// * `symbol` - The 16-bit key to search for.
pub fn get_env_quantity(state: &GeneState, symbol: u16) -> u32 { pub fn get_env_quantity(state: &GeneState, symbol: u16) -> u32 {
match state match state
.environment .environment
@@ -91,6 +107,14 @@ pub fn get_env_quantity(state: &GeneState, symbol: u16) -> u32 {
} }
} }
/// Sets the quantity of an environment symbol in a `GeneState`.
///
/// If quantity is 0, the record is removed. The environment is kept sorted alphabetically by symbol.
///
/// # Arguments
/// * `state` - The mutable gene state to update.
/// * `symbol` - The 16-bit key.
/// * `quantity` - The quantity to assign.
pub fn set_env_quantity( pub fn set_env_quantity(
state: &mut GeneState, state: &mut GeneState,
symbol: u16, symbol: u16,
@@ -125,6 +149,12 @@ pub fn set_env_quantity(
} }
} }
/// Adds a quantity to an environment symbol with saturating arithmetic.
///
/// # Arguments
/// * `state` - The mutable gene state.
/// * `symbol` - The 16-bit key.
/// * `quantity` - The quantity to add.
pub fn add_env_quantity( pub fn add_env_quantity(
state: &mut GeneState, state: &mut GeneState,
symbol: u16, symbol: u16,
@@ -136,6 +166,14 @@ pub fn add_env_quantity(
Ok(next) Ok(next)
} }
/// Subtracts a quantity from an environment symbol with saturating arithmetic.
///
/// If the resulting quantity drops to 0, the symbol is removed.
///
/// # Arguments
/// * `state` - The mutable gene state.
/// * `symbol` - The 16-bit key.
/// * `quantity` - The quantity to subtract.
pub fn sub_env_quantity( pub fn sub_env_quantity(
state: &mut GeneState, state: &mut GeneState,
symbol: u16, symbol: u16,
@@ -147,6 +185,12 @@ pub fn sub_env_quantity(
Ok(next) Ok(next)
} }
/// Encodes the environment records list into a compact byte slice.
///
/// Each record is written as a little-endian `u16` symbol followed by a little-endian `u32` quantity.
///
/// # Arguments
/// * `records` - The sorted environment records.
pub fn encode_environment(records: &[EnvironmentRecord]) -> Result<Vec<u8>, GeneError> { pub fn encode_environment(records: &[EnvironmentRecord]) -> Result<Vec<u8>, GeneError> {
validate_environment(records)?; validate_environment(records)?;
let mut out = Vec::with_capacity(records.len() * 6); let mut out = Vec::with_capacity(records.len() * 6);
@@ -157,6 +201,12 @@ pub fn encode_environment(records: &[EnvironmentRecord]) -> Result<Vec<u8>, Gene
Ok(out) Ok(out)
} }
/// Decodes environment records from a byte slice.
///
/// Validates that the length is a multiple of 6 and that records conform to sorting and quantity invariants.
///
/// # Arguments
/// * `blob` - The serialized byte slice.
pub fn decode_environment(blob: &[u8]) -> Result<Vec<EnvironmentRecord>, GeneError> { pub fn decode_environment(blob: &[u8]) -> Result<Vec<EnvironmentRecord>, GeneError> {
if !blob.len().is_multiple_of(6) { if !blob.len().is_multiple_of(6) {
return Err(GeneError::EnvironmentBlobLengthInvalid { len: blob.len() }); return Err(GeneError::EnvironmentBlobLengthInvalid { len: blob.len() });
@@ -180,6 +230,9 @@ pub fn decode_environment(blob: &[u8]) -> Result<Vec<EnvironmentRecord>, GeneErr
Ok(records) Ok(records)
} }
/// Computes the raw Blake3 cryptographic commitment of the `GeneState`.
///
/// Includes gene length, gene buffer, environment record count, and individual record key/values.
pub fn commitment(state: &GeneState) -> [u8; 32] { pub fn commitment(state: &GeneState) -> [u8; 32] {
let mut h = blake3::Hasher::new(); let mut h = blake3::Hasher::new();
h.update(b"chronoseal/gene/v1"); h.update(b"chronoseal/gene/v1");
@@ -193,10 +246,19 @@ pub fn commitment(state: &GeneState) -> [u8; 32] {
*h.finalize().as_bytes() *h.finalize().as_bytes()
} }
/// Computes the hex-encoded cryptographic commitment of the `GeneState`.
pub fn commitment_hex(state: &GeneState) -> String { pub fn commitment_hex(state: &GeneState) -> String {
hex::encode(commitment(state)) hex::encode(commitment(state))
} }
/// Computes a context-bound Blake3 cryptographic commitment of the `GeneState`.
///
/// Integrates `session_id` and the current `step` index into the hash to bind the commitment.
///
/// # Arguments
/// * `state` - The gene state.
/// * `session_id` - The client session ID.
/// * `step` - The mutation step index.
pub fn commitment_with_context(state: &GeneState, session_id: &str, step: u64) -> [u8; 32] { pub fn commitment_with_context(state: &GeneState, session_id: &str, step: u64) -> [u8; 32] {
let mut h = blake3::Hasher::new(); let mut h = blake3::Hasher::new();
h.update(b"chronoseal/gene/v1"); h.update(b"chronoseal/gene/v1");
@@ -206,10 +268,17 @@ pub fn commitment_with_context(state: &GeneState, session_id: &str, step: u64) -
*h.finalize().as_bytes() *h.finalize().as_bytes()
} }
/// Computes a context-bound, hex-encoded Blake3 cryptographic commitment of the `GeneState`.
///
/// # Arguments
/// * `state` - The gene state.
/// * `session_id` - The client session ID.
/// * `step` - The mutation step index.
pub fn commitment_hex_with_context(state: &GeneState, session_id: &str, step: u64) -> String { pub fn commitment_hex_with_context(state: &GeneState, session_id: &str, step: u64) -> String {
hex::encode(commitment_with_context(state, session_id, step)) hex::encode(commitment_with_context(state, session_id, step))
} }
/// Helper function to validate sorting, uniqueness, and non-zero properties of environment records.
fn validate_environment(records: &[EnvironmentRecord]) -> Result<(), GeneError> { fn validate_environment(records: &[EnvironmentRecord]) -> Result<(), GeneError> {
if records.len() > MAX_ENV_RECORDS { if records.len() > MAX_ENV_RECORDS {
return Err(GeneError::TooManyEnvironmentRecords { len: records.len() }); return Err(GeneError::TooManyEnvironmentRecords { len: records.len() });
+31 -8
View File
@@ -1,7 +1,15 @@
use crate::protocol::{EntropyData, StackState}; use crate::protocol::{EntropyData, StackState};
use blake3::Hasher; use blake3::Hasher;
/// Initial hash for a brand-new session: Blake3(session_id || pub_key || salt) /// Computes the initial hash for a brand-new attestation session.
///
/// The hash is constructed as:
/// `Blake3(session_id || pub_key || salt)`
///
/// # Arguments
/// * `session_id` - The unique hex-encoded identifier for the session.
/// * `pub_key` - The client's Ed25519 public key.
/// * `salt` - The initial server-issued salt.
pub fn initial_hash(session_id: &str, pub_key: &[u8], salt: &[u8]) -> Vec<u8> { pub fn initial_hash(session_id: &str, pub_key: &[u8], salt: &[u8]) -> Vec<u8> {
let mut h = Hasher::new(); let mut h = Hasher::new();
h.update(session_id.as_bytes()); h.update(session_id.as_bytes());
@@ -10,7 +18,18 @@ pub fn initial_hash(session_id: &str, pub_key: &[u8], salt: &[u8]) -> Vec<u8> {
h.finalize().as_bytes().to_vec() h.finalize().as_bytes().to_vec()
} }
/// Next hash in the chain: Blake3 with the salt mixed in (no keyed mode needed) /// Computes the next hash in the Blake3 attestation chain.
///
/// This mixes in the previous hash head, the client timestamp, the serialized entropy data,
/// the VM stack state, and the server-issued salt. Uses `serde_json::to_vec` to avoid
/// intermediate heap string allocations and UTF-8 verification checks.
///
/// # Arguments
/// * `prev_hash` - The previous hash-chain head.
/// * `timestamp` - The client-supplied heartbeat timestamp.
/// * `entropy` - The collected browser interaction entropy.
/// * `stack` - The final VM stack state after running the opcode program.
/// * `salt` - The server-issued salt for rotation.
pub fn next_chain_hash( pub fn next_chain_hash(
prev_hash: &[u8], prev_hash: &[u8],
timestamp: u64, timestamp: u64,
@@ -18,14 +37,13 @@ pub fn next_chain_hash(
stack: &StackState, stack: &StackState,
salt: &[u8], salt: &[u8],
) -> Vec<u8> { ) -> Vec<u8> {
let entropy_json = serde_json::to_string(entropy).unwrap(); let entropy_bytes = serde_json::to_vec(entropy).unwrap_or_default();
let stack_json = serde_json::to_string(stack).unwrap(); let stack_bytes = serde_json::to_vec(stack).unwrap_or_default();
let entropy_hash = blake3::hash(entropy_json.as_bytes()); let entropy_hash = blake3::hash(&entropy_bytes);
let stack_hash = blake3::hash(stack_json.as_bytes()); let stack_hash = blake3::hash(&stack_bytes);
let mut h = Hasher::new(); let mut h = Hasher::new();
// Mix the salt into the hash state
h.update(salt); h.update(salt);
h.update(prev_hash); h.update(prev_hash);
h.update(&timestamp.to_le_bytes()); h.update(&timestamp.to_le_bytes());
@@ -34,7 +52,12 @@ pub fn next_chain_hash(
h.finalize().as_bytes().to_vec() h.finalize().as_bytes().to_vec()
} }
/// Hash of all stack items for VM HASH opcode /// Computes a 32-bit FNV-like Blake3 hash of all stack elements.
///
/// This is used by the VM `HASH` opcode to fold the current stack state into a single value.
///
/// # Arguments
/// * `stack` - The list of u32 stack elements to hash.
pub fn hash_stack(stack: &[u32]) -> u32 { pub fn hash_stack(stack: &[u32]) -> u32 {
let data: Vec<u8> = stack.iter().flat_map(|x| x.to_le_bytes()).collect(); let data: Vec<u8> = stack.iter().flat_map(|x| x.to_le_bytes()).collect();
let hash = blake3::hash(&data); let hash = blake3::hash(&data);
+45
View File
@@ -1,73 +1,118 @@
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
/// Request payload sent by the client to initialize a new attestation session.
#[derive(Debug, Clone, Deserialize, Serialize)] #[derive(Debug, Clone, Deserialize, Serialize)]
pub struct InitRequest { pub struct InitRequest {
/// Hex-encoded 32-byte Ed25519 public verifying key generated by the client.
pub public_key: String, pub public_key: String,
} }
/// Response payload returned by the server upon successful session initialization.
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
pub struct InitResponse { pub struct InitResponse {
/// The unique hex-encoded session identifier.
pub session_id: String, pub session_id: String,
/// The initial server-issued salt to be mixed in the first heartbeat's hash.
pub salt: String, pub salt: String,
/// Base64-encoded initial VM program for client stack execution.
pub opcodes_b64: String, pub opcodes_b64: String,
/// The computed initial hash of the attestation chain.
pub initial_hash: String, pub initial_hash: String,
/// Timestamp in milliseconds indicating when the session expires.
pub expires_at: u64, pub expires_at: u64,
/// Minimum time in milliseconds allowed between subsequent heartbeats.
pub heartbeat_min_interval_ms: u64, pub heartbeat_min_interval_ms: u64,
/// Maximum time in milliseconds allowed between subsequent heartbeats.
pub heartbeat_max_interval_ms: u64, pub heartbeat_max_interval_ms: u64,
/// Size of the synthetic gene byte buffer.
pub gene_size: u32, pub gene_size: u32,
/// The current mutation step index (starts at 1).
pub mutation_step: u64, pub mutation_step: u64,
/// Base64-encoded initial gene mutation program.
pub mutation_order_b64: String, pub mutation_order_b64: String,
/// The number of mutation rounds configured on the server.
pub mutation_rounds: u8,
} }
/// Heartbeat request payload submitted periodically by the client to prove session continuity.
#[derive(Debug, Clone, Deserialize, Serialize)] #[derive(Debug, Clone, Deserialize, Serialize)]
pub struct HeartbeatRequest { pub struct HeartbeatRequest {
/// The session identifier.
pub session_id: String, pub session_id: String,
/// The expected hash from the previous heartbeat/initialization step.
pub prev_hash: String, pub prev_hash: String,
/// The client's current system timestamp in milliseconds.
pub timestamp: u64, pub timestamp: u64,
/// The collected client entropy data (such as mouse events).
pub entropy_data: EntropyData, pub entropy_data: EntropyData,
/// The final execution state of the client's VM stack program.
pub stack_state: StackState, pub stack_state: StackState,
/// The client's browser hardware and layout fingerprint.
pub fingerprint: Fingerprint, pub fingerprint: Fingerprint,
/// The mutation step index corresponding to the pending mutation.
pub mutation_step: u64, pub mutation_step: u64,
/// Hex-encoded commitment of the mutated gene state.
pub gene_commitment: String, pub gene_commitment: String,
/// Ed25519 signature of the canonical JSON-serialized payload.
pub signature: String, pub signature: String,
} }
/// Response payload returned by the server for heartbeat submissions.
///
/// In case of silent rejection, all fields except `status` are omitted.
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HeartbeatResponse { pub struct HeartbeatResponse {
/// Attestation status, typically "ok" even on silent failures.
pub status: String, pub status: String,
/// The next server-issued salt for hash chain progression.
#[serde(skip_serializing_if = "Option::is_none")] #[serde(skip_serializing_if = "Option::is_none")]
pub next_salt: Option<String>, pub next_salt: Option<String>,
/// The next expected mutation step index.
#[serde(skip_serializing_if = "Option::is_none")] #[serde(skip_serializing_if = "Option::is_none")]
pub next_mutation_step: Option<u64>, pub next_mutation_step: Option<u64>,
/// Base64-encoded next mutation program for client gene progression.
#[serde(skip_serializing_if = "Option::is_none")] #[serde(skip_serializing_if = "Option::is_none")]
pub next_mutation_order_b64: Option<String>, pub next_mutation_order_b64: Option<String>,
} }
/// Client browser fingerprint metadata used for basic sanity checks.
#[derive(Debug, Clone, Deserialize, Serialize)] #[derive(Debug, Clone, Deserialize, Serialize)]
pub struct Fingerprint { pub struct Fingerprint {
/// Aspect ratio of the client screen.
#[serde(rename = "aspectRatio")] #[serde(rename = "aspectRatio")]
pub aspect_ratio: String, pub aspect_ratio: String,
/// Device pixel ratio of the screen.
#[serde(rename = "devicePixelRatio")] #[serde(rename = "devicePixelRatio")]
pub device_pixel_ratio: String, pub device_pixel_ratio: String,
/// Number of logical processor cores available.
#[serde(rename = "hardwareConcurrency")] #[serde(rename = "hardwareConcurrency")]
pub hardware_concurrency: u32, pub hardware_concurrency: u32,
} }
/// Wrapper for browser-side entropy collection.
#[derive(Debug, Clone, Deserialize, Serialize)] #[derive(Debug, Clone, Deserialize, Serialize)]
pub struct EntropyData { pub struct EntropyData {
/// A chronological list of mouse movement events.
pub events: Vec<MouseEvent>, pub events: Vec<MouseEvent>,
} }
/// Information about a single mouse movement interaction.
#[derive(Deserialize, Serialize, Clone, Debug)] #[derive(Deserialize, Serialize, Clone, Debug)]
pub struct MouseEvent { pub struct MouseEvent {
/// Absolute horizontal coordinate of the cursor.
pub x: f64, pub x: f64,
/// Absolute vertical coordinate of the cursor.
pub y: f64, pub y: f64,
/// Relative timestamp in milliseconds of the event occurrence.
#[serde(rename = "t")] #[serde(rename = "t")]
pub timestamp_ms: f64, pub timestamp_ms: f64,
} }
/// The state of the VM stack machine after executing a program.
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StackState { pub struct StackState {
/// The elements remaining on the stack.
pub stack: Vec<u32>, pub stack: Vec<u32>,
/// The final instruction pointer location at program completion or termination.
pub ip: u16, pub ip: u16,
} }
+12
View File
@@ -25,6 +25,9 @@ pub fn execute(program: &[u8]) -> StackState {
program[ip + 3], program[ip + 3],
]); ]);
ip += 4; ip += 4;
if stack.len() >= crate::constants::MAX_STACK_DEPTH {
break;
}
stack.push(val); stack.push(val);
} }
0x01..=0x07 => { 0x01..=0x07 => {
@@ -43,6 +46,9 @@ pub fn execute(program: &[u8]) -> StackState {
0x07 => a.rotate_left(b % 32), 0x07 => a.rotate_left(b % 32),
_ => unreachable!(), _ => unreachable!(),
}; };
if stack.len() >= crate::constants::MAX_STACK_DEPTH {
break;
}
stack.push(r); stack.push(r);
} }
0x08 => { 0x08 => {
@@ -50,11 +56,17 @@ pub fn execute(program: &[u8]) -> StackState {
break; break;
} }
let a = stack.pop().unwrap(); let a = stack.pop().unwrap();
if stack.len() >= crate::constants::MAX_STACK_DEPTH {
break;
}
stack.push(!a); stack.push(!a);
} }
0x09 => { 0x09 => {
let r = crate::hashing::hash_stack(&stack); let r = crate::hashing::hash_stack(&stack);
stack.clear(); stack.clear();
if stack.len() >= crate::constants::MAX_STACK_DEPTH {
break;
}
stack.push(r); stack.push(r);
} }
_ => break, _ => break,
+65
View File
@@ -88,10 +88,18 @@ impl From<GeneError> for MutationError {
} }
} }
/// Encodes a `MutationOrder` into standard Base64 representation of its bytecode.
pub fn encode_order_b64(order: &MutationOrder) -> String { pub fn encode_order_b64(order: &MutationOrder) -> String {
base64::Engine::encode(&base64::engine::general_purpose::STANDARD, &order.program) base64::Engine::encode(&base64::engine::general_purpose::STANDARD, &order.program)
} }
/// Decodes a `MutationOrder` from its Base64 representation.
///
/// Validates that the decoded program size does not exceed the allowed maximum budget size.
///
/// # Arguments
/// * `step` - The step index associated with this mutation order.
/// * `b64` - The Base64 string containing the raw bytecode.
pub fn decode_order_b64(step: u64, b64: &str) -> Result<MutationOrder, MutationError> { pub fn decode_order_b64(step: u64, b64: &str) -> Result<MutationOrder, MutationError> {
let program = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, b64) let program = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, b64)
.map_err(MutationError::Base64)?; .map_err(MutationError::Base64)?;
@@ -101,11 +109,27 @@ pub fn decode_order_b64(step: u64, b64: &str) -> Result<MutationOrder, MutationE
Ok(MutationOrder { step, program }) Ok(MutationOrder { step, program })
} }
/// Generates a randomized `MutationOrder` program for a given step and gene size.
///
/// Uses thread-local random number generator.
///
/// # Arguments
/// * `step` - The step index.
/// * `gene_size` - The length of the gene byte buffer.
pub fn generate_order(step: u64, gene_size: usize) -> MutationOrder { pub fn generate_order(step: u64, gene_size: usize) -> MutationOrder {
let mut rng = rand::thread_rng(); let mut rng = rand::thread_rng();
generate_order_with_rng(&mut rng, step, gene_size) generate_order_with_rng(&mut rng, step, gene_size)
} }
/// Generates a randomized `MutationOrder` program using a specific custom RNG.
///
/// Builds a program containing between 20 and 36 mutation instructions (e.g. loads, point changes,
/// insertions, deletions, env modifications) and ensures a minimum number of finalize hash steps are included.
///
/// # Arguments
/// * `rng` - The random number generator.
/// * `step` - The step index.
/// * `gene_size` - The length of the gene byte buffer.
pub fn generate_order_with_rng<R: Rng + ?Sized>( pub fn generate_order_with_rng<R: Rng + ?Sized>(
rng: &mut R, rng: &mut R,
step: u64, step: u64,
@@ -213,10 +237,12 @@ pub fn generate_order_with_rng<R: Rng + ?Sized>(
MutationOrder { step, program } MutationOrder { step, program }
} }
/// Clones the `GeneState` and executes the mutation program for `DEFAULT_MUTATION_ROUNDS`.
pub fn apply_program_clone(state: &GeneState, program: &[u8]) -> Result<GeneState, MutationError> { pub fn apply_program_clone(state: &GeneState, program: &[u8]) -> Result<GeneState, MutationError> {
apply_program_clone_with_rounds(state, program, DEFAULT_MUTATION_ROUNDS) apply_program_clone_with_rounds(state, program, DEFAULT_MUTATION_ROUNDS)
} }
/// Clones the `GeneState` and executes the mutation program for a specific number of rounds.
pub fn apply_program_clone_with_rounds( pub fn apply_program_clone_with_rounds(
state: &GeneState, state: &GeneState,
program: &[u8], program: &[u8],
@@ -227,6 +253,7 @@ pub fn apply_program_clone_with_rounds(
Ok(next) Ok(next)
} }
/// Executes the mutation program on the mutable `GeneState` reference for a specific number of rounds.
pub fn apply_program_with_rounds( pub fn apply_program_with_rounds(
state: &mut GeneState, state: &mut GeneState,
program: &[u8], program: &[u8],
@@ -236,11 +263,21 @@ pub fn apply_program_with_rounds(
Ok(()) Ok(())
} }
/// Executes the mutation program on the mutable `GeneState` reference for `DEFAULT_MUTATION_ROUNDS`.
pub fn apply_program(state: &mut GeneState, program: &[u8]) -> Result<(), MutationError> { pub fn apply_program(state: &mut GeneState, program: &[u8]) -> Result<(), MutationError> {
let _ = execute_program_with_rounds(state, program, DEFAULT_MUTATION_ROUNDS)?; let _ = execute_program_with_rounds(state, program, DEFAULT_MUTATION_ROUNDS)?;
Ok(()) 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
/// programs from lagging the host server thread or client runtime.
///
/// # Arguments
/// * `state` - The mutable gene state buffer.
/// * `program` - The raw bytecode sequence.
/// * `rounds` - The requested number of execution rounds.
pub fn execute_program_with_rounds( pub fn execute_program_with_rounds(
state: &mut GeneState, state: &mut GeneState,
program: &[u8], program: &[u8],
@@ -317,6 +354,13 @@ fn estimate_program_cost(program: &[u8]) -> usize {
cost.max(1) cost.max(1)
} }
/// Executes the VM mutation program on the mutable `GeneState` reference.
///
/// This interprets VM mutation opcodes to modify the gene byte array and environment records.
///
/// # Arguments
/// * `state` - The mutable gene state to mutate.
/// * `program` - The raw instruction bytecode slice.
pub fn execute_program( pub fn execute_program(
state: &mut GeneState, state: &mut GeneState,
program: &[u8], program: &[u8],
@@ -835,4 +879,25 @@ mod tests {
"mutation execution too slow: {elapsed:?}" "mutation execution too slow: {elapsed:?}"
); );
} }
#[test]
fn test_vm_instruction_budget_soft_cap() {
let state = new_state(8).unwrap();
// Construct a program with 130 OP_FINALIZE_GENE_HASH instructions.
// HASH has HASH_OPCODE_INSTRUCTION_COST = 16.
// Total cost will be 130 * 16 = 2080, which exceeds MAX_MUTATION_INSTRUCTION_BUDGET (2048).
let program = vec![OP_FINALIZE_GENE_HASH; 130];
let cost = estimate_program_cost(&program);
assert!(cost >= 2080);
// Assert that the max allowed rounds is calculated as 1 since cost > budget.
let expected_rounds = std::cmp::max(1, MAX_MUTATION_INSTRUCTION_BUDGET / cost);
assert_eq!(expected_rounds, 1);
// Execute the program with a requested 10 rounds.
// The runtime should execute it successfully without panic, while applying the round limitation.
let mut test_state = state.clone();
let trace = execute_program_with_rounds(&mut test_state, &program, 10).unwrap();
assert_eq!(trace.final_ip, program.len());
}
} }
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "chronoseal-wasm" name = "chronoseal-wasm"
version = "0.6.0" version = "1.0.1"
edition = "2021" edition = "2021"
[lib] [lib]
+8 -2
View File
@@ -51,8 +51,14 @@ pub fn compute_next_hash(
) -> String { ) -> String {
let prev = hex::decode(prev_hash_hex).unwrap_or_default(); let prev = hex::decode(prev_hash_hex).unwrap_or_default();
let salt = hex::decode(salt_hex).unwrap_or_default(); let salt = hex::decode(salt_hex).unwrap_or_default();
let entropy = serde_json::from_str::<shared::protocol::EntropyData>(entropy_data_json).unwrap(); let entropy = match serde_json::from_str::<shared::protocol::EntropyData>(entropy_data_json) {
let stack = serde_json::from_str::<shared::protocol::StackState>(stack_state_json).unwrap(); Ok(v) => v,
Err(_) => return String::new(),
};
let stack = match serde_json::from_str::<shared::protocol::StackState>(stack_state_json) {
Ok(v) => v,
Err(_) => return String::new(),
};
let new = shared::hashing::next_chain_hash(&prev, timestamp, &entropy, &stack, &salt); let new = shared::hashing::next_chain_hash(&prev, timestamp, &entropy, &stack, &salt);
hex::encode(new) hex::encode(new)
} }
+9 -60
View File
@@ -4,70 +4,19 @@ use wasm_bindgen::prelude::*;
#[wasm_bindgen] #[wasm_bindgen]
pub fn run_program(program_b64: &str) -> JsValue { pub fn run_program(program_b64: &str) -> JsValue {
use base64::Engine; use base64::Engine;
let bytes = base64::engine::general_purpose::STANDARD let bytes = match base64::engine::general_purpose::STANDARD.decode(program_b64) {
.decode(program_b64) Ok(v) => v,
.unwrap(); Err(_) => return JsValue::NULL,
};
let state = execute(&bytes); 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 { fn execute(program: &[u8]) -> StackState {
let mut stack: Vec<u32> = Vec::new(); shared::vm::execute(program)
let mut ip: usize = 0;
while ip < program.len() {
let op = program[ip];
ip += 1;
match op {
0x00 => {
if ip + 4 > program.len() {
break;
}
let val = u32::from_le_bytes([
program[ip],
program[ip + 1],
program[ip + 2],
program[ip + 3],
]);
ip += 4;
stack.push(val);
}
0x01..=0x07 => {
if stack.len() < 2 {
break;
}
let b = stack.pop().unwrap();
let a = stack.pop().unwrap();
let r = match op {
0x01 => a.wrapping_add(b),
0x02 => a.wrapping_sub(b),
0x03 => a.wrapping_mul(b),
0x04 => a ^ b,
0x05 => a & b,
0x06 => a | b,
0x07 => a.rotate_left(b % 32),
_ => unreachable!(),
};
stack.push(r);
}
0x08 => {
if stack.is_empty() {
break;
}
let a = stack.pop().unwrap();
stack.push(!a);
}
0x09 => {
let r = shared::hashing::hash_stack(&stack);
stack.clear();
stack.push(r);
}
_ => break,
}
}
StackState {
stack,
ip: ip as u16,
}
} }
#[cfg(test)] #[cfg(test)]