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
.env
.idea/
.antigravitycli/
Generated
+19 -4
View File
@@ -275,7 +275,7 @@ dependencies = [
[[package]]
name = "chronoseal-replay"
version = "0.1.0"
version = "1.0.1"
dependencies = [
"anyhow",
"base64",
@@ -290,12 +290,13 @@ dependencies = [
[[package]]
name = "chronoseal-server"
version = "0.6.0"
version = "1.0.1"
dependencies = [
"axum",
"base64",
"clap",
"clap_complete",
"dashmap",
"ed25519-dalek",
"hex",
"r2d2",
@@ -319,7 +320,7 @@ dependencies = [
[[package]]
name = "chronoseal-wasm"
version = "0.6.0"
version = "1.0.1"
dependencies = [
"base64",
"blake3",
@@ -509,6 +510,20 @@ dependencies = [
"syn",
]
[[package]]
name = "dashmap"
version = "6.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e6361d5c062261c78a176addb82d4c821ae42bed6089de0e12603cd25de2059c"
dependencies = [
"cfg-if",
"crossbeam-utils",
"hashbrown 0.14.5",
"lock_api",
"once_cell",
"parking_lot_core",
]
[[package]]
name = "der"
version = "0.7.10"
@@ -1946,7 +1961,7 @@ dependencies = [
[[package]]
name = "shared"
version = "0.6.0"
version = "1.0.1"
dependencies = [
"base64",
"blake3",
+2 -1
View File
@@ -4,7 +4,8 @@ resolver = "2"
members = [
"shared",
"server",
"wasm"
"wasm",
"chronoseal-replay"
]
[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_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"]
+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">
</a>
<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>
<img src="https://img.shields.io/badge/wasm-rust--compiled-blueviolet.svg" alt="WASM">
</p>
@@ -536,7 +536,32 @@ Set storage mode with `db_type` or `CHRONOSEAL_DB_TYPE`.
| `sqlite-in-disk` | SQLite database persisted at `db_path`. |
| `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
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "chronoseal-replay"
version = "0.1.0"
version = "1.0.1"
edition = "2021"
[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 init_req = InitRequest { public_key: pk_hex };
let resp = client
@@ -126,7 +130,10 @@ fn do_handshake(client: &reqwest::blocking::Client, base_url: &str, sk: &Signing
.send()?;
if !resp.status().is_success() {
return Err(anyhow!("Handshake failed with HTTP status: {}", resp.status()));
return Err(anyhow!(
"Handshake failed with HTTP status: {}",
resp.status()
));
}
let init_resp: InitResponse = resp.json()?;
@@ -137,11 +144,17 @@ fn run_built_in_scenarios(client: &reqwest::blocking::Client, base_url: &str) ->
let mut failures = 0;
let scenarios = [
("valid_progression", run_valid_progression as fn(&reqwest::blocking::Client, &str) -> Result<()>),
(
"valid_progression",
run_valid_progression as fn(&reqwest::blocking::Client, &str) -> Result<()>,
),
("stale_replay", run_stale_replay),
("invalid_signature", run_invalid_signature),
("invalid_vm_stack", run_invalid_vm_stack),
("invalid_mutation_commitment", run_invalid_mutation_commitment),
(
"invalid_mutation_commitment",
run_invalid_mutation_commitment,
),
("drifted_timestamp", run_drifted_timestamp),
("concurrent_heartbeat", run_concurrent_heartbeat),
("rate_limit_trigger", run_rate_limit_trigger),
@@ -197,7 +210,8 @@ fn run_valid_progression(client: &reqwest::blocking::Client, base_url: &str) ->
init.mutation_rounds,
)?;
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, mutation_step);
let commitment =
shared::gene::commitment_hex_with_context(&candidate, &init.session_id, mutation_step);
let timestamp = current_time_ms();
let entropy = test_entropy();
@@ -215,13 +229,14 @@ fn run_valid_progression(client: &reqwest::blocking::Client, base_url: &str) ->
sign_request(&sk, &mut req)?;
let resp = client
.post(format!("{}/hb", base_url))
.json(&req)
.send()?;
let resp = client.post(format!("{}/hb", base_url)).json(&req).send()?;
if !resp.status().is_success() {
return Err(anyhow!("Step {} /hb returned HTTP error: {}", step, resp.status()));
return Err(anyhow!(
"Step {} /hb returned HTTP error: {}",
step,
resp.status()
));
}
let hb_resp: HeartbeatResponse = resp.json()?;
@@ -230,9 +245,15 @@ fn run_valid_progression(client: &reqwest::blocking::Client, base_url: &str) ->
}
// Verify it was a successful validation (not a silent rejection)
let next_salt = hb_resp.next_salt.ok_or_else(|| anyhow!("Step {} was silently rejected", step))?;
let next_step = hb_resp.next_mutation_step.ok_or_else(|| anyhow!("Step {} missing next mutation step", step))?;
let next_order = hb_resp.next_mutation_order_b64.ok_or_else(|| anyhow!("Step {} missing next mutation order", step))?;
let next_salt = hb_resp
.next_salt
.ok_or_else(|| anyhow!("Step {} was silently rejected", step))?;
let next_step = hb_resp
.next_mutation_step
.ok_or_else(|| anyhow!("Step {} missing next mutation step", step))?;
let next_order = hb_resp
.next_mutation_order_b64
.ok_or_else(|| anyhow!("Step {} missing next mutation order", step))?;
println!("Step {} successful. Salt rotated: {}", step, next_salt);
@@ -265,12 +286,21 @@ fn run_stale_replay(client: &reqwest::blocking::Client, base_url: &str) -> Resul
let sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?;
let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?;
let opcodes = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&init.opcodes_b64,
)?;
let stack_state = shared::vm::execute(&opcodes);
let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let order =
shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap();
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?;
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(
&gene_state,
&order.program,
init.mutation_rounds,
)?;
let commitment =
shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let mut req = HeartbeatRequest {
session_id: init.session_id.clone(),
@@ -296,7 +326,9 @@ fn run_stale_replay(client: &reqwest::blocking::Client, base_url: &str) -> Resul
let resp2 = client.post(format!("{}/hb", base_url)).json(&req).send()?;
let hb2: HeartbeatResponse = resp2.json()?;
if hb2.next_salt.is_some() {
return Err(anyhow!("Replayed heartbeat was successfully accepted (broken replay protection)"));
return Err(anyhow!(
"Replayed heartbeat was successfully accepted (broken replay protection)"
));
}
println!("Stale replay correctly rejected.");
@@ -308,12 +340,21 @@ fn run_invalid_signature(client: &reqwest::blocking::Client, base_url: &str) ->
let sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?;
let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?;
let opcodes = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&init.opcodes_b64,
)?;
let stack_state = shared::vm::execute(&opcodes);
let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let order =
shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap();
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?;
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(
&gene_state,
&order.program,
init.mutation_rounds,
)?;
let commitment =
shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let mut req = HeartbeatRequest {
session_id: init.session_id.clone(),
@@ -344,10 +385,16 @@ fn run_invalid_vm_stack(client: &reqwest::blocking::Client, base_url: &str) -> R
let sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?;
let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let order =
shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap();
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?;
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(
&gene_state,
&order.program,
init.mutation_rounds,
)?;
let commitment =
shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let mut req = HeartbeatRequest {
session_id: init.session_id.clone(),
@@ -375,12 +422,18 @@ fn run_invalid_vm_stack(client: &reqwest::blocking::Client, base_url: &str) -> R
Ok(())
}
fn run_invalid_mutation_commitment(client: &reqwest::blocking::Client, base_url: &str) -> Result<()> {
fn run_invalid_mutation_commitment(
client: &reqwest::blocking::Client,
base_url: &str,
) -> Result<()> {
let mut csprng = OsRng;
let sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?;
let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?;
let opcodes = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&init.opcodes_b64,
)?;
let stack_state = shared::vm::execute(&opcodes);
let mut req = HeartbeatRequest {
@@ -411,12 +464,21 @@ fn run_drifted_timestamp(client: &reqwest::blocking::Client, base_url: &str) ->
let sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?;
let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?;
let opcodes = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&init.opcodes_b64,
)?;
let stack_state = shared::vm::execute(&opcodes);
let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let order =
shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap();
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?;
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(
&gene_state,
&order.program,
init.mutation_rounds,
)?;
let commitment =
shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let mut req = HeartbeatRequest {
session_id: init.session_id.clone(),
@@ -446,12 +508,21 @@ fn run_concurrent_heartbeat(client: &reqwest::blocking::Client, base_url: &str)
let sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?;
let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?;
let opcodes = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&init.opcodes_b64,
)?;
let stack_state = shared::vm::execute(&opcodes);
let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let order =
shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap();
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?;
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(
&gene_state,
&order.program,
init.mutation_rounds,
)?;
let commitment =
shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let mut req = HeartbeatRequest {
session_id: init.session_id.clone(),
@@ -470,10 +541,8 @@ fn run_concurrent_heartbeat(client: &reqwest::blocking::Client, base_url: &str)
let client_clone = client.clone();
let req_clone = req.clone();
let url_clone = format!("{}/hb", base_url);
let handle = std::thread::spawn(move || {
client_clone.post(&url_clone).json(&req_clone).send()
});
let handle = std::thread::spawn(move || client_clone.post(&url_clone).json(&req_clone).send());
let resp2 = client.post(format!("{}/hb", base_url)).json(&req).send()?;
let resp1_res = handle.join().map_err(|_| anyhow!("Thread panicked"))?;
@@ -485,7 +554,10 @@ fn run_concurrent_heartbeat(client: &reqwest::blocking::Client, base_url: &str)
// One must succeed and one must fail (silent rejection) because of CAS check
let successes = (hb1.next_salt.is_some() as usize) + (hb2.next_salt.is_some() as usize);
if successes != 1 {
return Err(anyhow!("Expected exactly one concurrent heartbeat to succeed. Got: {}", successes));
return Err(anyhow!(
"Expected exactly one concurrent heartbeat to succeed. Got: {}",
successes
));
}
println!("Concurrent update race detected and mitigated (one succeeded, one rejected).");
@@ -497,12 +569,21 @@ fn run_rate_limit_trigger(client: &reqwest::blocking::Client, base_url: &str) ->
let sk = SigningKey::generate(&mut csprng);
let init = do_handshake(client, base_url, &sk)?;
let opcodes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &init.opcodes_b64)?;
let opcodes = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&init.opcodes_b64,
)?;
let stack_state = shared::vm::execute(&opcodes);
let order = shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let order =
shared::vm_extensions::decode_order_b64(init.mutation_step, &init.mutation_order_b64)?;
let gene_state = shared::gene::new_state(init.gene_size as usize).unwrap();
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(&gene_state, &order.program, init.mutation_rounds)?;
let commitment = shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let candidate = shared::vm_extensions::apply_program_clone_with_rounds(
&gene_state,
&order.program,
init.mutation_rounds,
)?;
let commitment =
shared::gene::commitment_hex_with_context(&candidate, &init.session_id, init.mutation_step);
let mut req = HeartbeatRequest {
session_id: init.session_id.clone(),
@@ -534,14 +615,20 @@ fn run_rate_limit_trigger(client: &reqwest::blocking::Client, base_url: &str) ->
}
if !rate_limited {
return Err(anyhow!("Rate limiter was not triggered after 35 rapid requests"));
return Err(anyhow!(
"Rate limiter was not triggered after 35 rapid requests"
));
}
println!("Rate limiter correctly triggered.");
Ok(())
}
fn run_file_scenario(_client: &reqwest::blocking::Client, _base_url: &str, file_path: &str) -> Result<()> {
fn run_file_scenario(
_client: &reqwest::blocking::Client,
_base_url: &str,
file_path: &str,
) -> Result<()> {
let scenario_content = std::fs::read_to_string(file_path)?;
let scenario: serde_json::Value = serde_json::from_str(&scenario_content)?;
+1
View File
@@ -27,6 +27,7 @@ RestrictSUIDSGID=yes
LockPersonality=yes
SystemCallArchitectures=native
ReadWritePaths=/run/chronoseal.pid
ReadWritePaths=/var/lib/chronoseal
# Logging
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 |
| `fingerprint.aspectRatio` | string | yes | Screen aspect ratio; server accepts numeric strings in range `0.5..=3.0` |
| `fingerprint.devicePixelRatio` | string | yes | Device pixel ratio; server accepts numeric strings in range `(0, 5]` |
| `fingerprint.hardwareConcurrency` | number | yes | Positive hardware concurrency value |
| `fingerprint.hardwareConcurrency` | number | yes | Hardware concurrency value; server accepts integers in range `1..=256` |
| `mutation_step` | number | yes | Mutation step currently expected by the server |
| `gene_commitment` | string | yes | Context-bound commitment produced by the WASM mutation preview |
| `signature` | string | yes | Ed25519 signature over the canonical payload |
+2 -2
View File
@@ -379,7 +379,7 @@ Storage is abstracted by `DbPool`.
|---|---|---|
| SQLite memory | `sqlite-in-memory` | default, process-local, ephemeral |
| 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:
@@ -389,7 +389,7 @@ The storage layer must support:
- delete expired sessions
- 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
+51 -1
View File
@@ -178,13 +178,63 @@ sudo mkdir -p /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
export CHRONOSEAL_DB_TYPE=valkey
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.
## 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.
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 |
| ------------------- | -----: |
| `chronoseal-server` | 30 |
| `chronoseal-wasm` | 24 |
| `shared` | 35 |
| **Total** | **89** |
| Crate | Tests |
| ---------------------------- | ------ |
| `chronoseal-server` | 33 |
| `chronoseal-wasm` | 24 |
| `shared` (unit) | 36 |
| `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:
* **Security invariants** over raw coverage metrics
* **Deterministic parity** between server and browser WASM runtimes
* **Negative-path testing** (tampering, replay, malformed input, edge cases)
* **Fuzz-style and randomized testing** for mutation logic
* **Performance regression detection**
* **Long-term protocol stability**
- **Security invariants** over raw coverage metrics
- **Deterministic parity** between server and browser WASM runtimes
- **Negative-path testing** (tampering, replay, malformed input, edge cases)
- **Property-based and fuzz-style testing** for mutation logic and VM robustness
- **Performance regression detection**
- **Long-term protocol stability**
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:
* Database backend selection
* TOML configuration parsing
* Command-line override behavior
* Default configuration values
* Runtime initialization logic
- Database backend selection
- TOML configuration parsing
- Command-line override behavior
- Default configuration values
- Runtime initialization logic
Supported backends include:
* `sqlite-in-memory`
* `sqlite-in-disk`
* `valkey`
- `sqlite-in-memory`
- `sqlite-in-disk`
- `valkey` (redis-compatible via r2d2 connection pool)
Example tests:
```text
```
test_apply_run_args_overrides_db_type
test_default_db_type_is_sqlite_in_memory
test_toml_parses_db_type_kebab_case
test_init_db_pool_sqlite_in_memory
test_init_db_pool_sqlite_in_disk
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 creation
* Public key validation
* Expiration handling
* Replay attack prevention
* Mutation step enforcement
* Commitment verification
* Long-running deterministic parity
- Session creation
- Public key validation
- Expiration handling
- Replay attack prevention
- Mutation step enforcement
- Commitment verification
- Long-running deterministic parity
Example tests:
```text
```
test_create_session_rejects_invalid_public_key_length
test_expired_session_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:
* Deterministic server/client parity
* Mutation order execution
* Gene state integrity
* Preview → Commit → Discard lifecycle
* Randomized mutation programs
* Edge-case validation
* Performance regression detection
- Deterministic server/client parity
- Mutation order execution
- Gene state integrity
- Preview → Commit → Discard lifecycle
- Randomized mutation programs
- Edge-case validation
- Performance regression detection
Example tests:
```text
```
test_server_client_parity_across_random_orders
test_generate_order_is_deterministic_for_seeded_rng
test_invalid_positions_wrap_deterministically
@@ -117,15 +119,15 @@ test_performance_smoke_mutation_execution
Heartbeat validation tests verify:
* Successful state advancement
* Silent rejection behavior
* Commitment validation
* Rate limiting
* Next-state mutation generation
- Successful state advancement
- Silent rejection behavior
- Commitment validation
- Rate limiting
- Next-state mutation generation
Example tests:
```text
```
test_handler_success_returns_next_mutation_fields
test_handler_tampered_commitment_is_silent_failure
test_handler_rate_limit_returns_no_mutation_data
@@ -137,16 +139,16 @@ test_handler_rate_limit_returns_no_mutation_data
Behavioral validation tests verify:
* Minimum mouse activity
* Minimum movement distance
* Pause detection
* Speed thresholds
* Optional activity requirements
* Fingerprint-related validation paths
- Minimum mouse activity
- Minimum movement distance
- Pause detection
- Speed thresholds
- Optional activity requirements
- Fingerprint-related validation paths
Example tests:
```text
```
test_validate_mouse_success
test_validate_mouse_insufficient_events
test_validate_mouse_insufficient_distance
@@ -161,15 +163,28 @@ test_validate_mouse_require_activity_toggle
Storage tests verify:
* SQLite in-memory operation
* SQLite disk-backed operation
* Valkey compatibility mode
* Session CRUD behavior
* Expiration cleanup
* Runtime statistics reporting
- SQLite in-memory operation
- SQLite disk-backed operation and pool concurrency
- Valkey compatibility mode and r2d2 pool concurrency
- Session CRUD behavior including the `opcodes` field
- Expiration cleanup
- Runtime statistics reporting
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
@@ -178,28 +193,19 @@ The VM core is tested extensively across both WASM and shared crates.
Coverage includes:
* ADD
* SUB
* MUL
* XOR
* AND
* OR
* NOT
* HASH
* ROT
* PUSH
- ADD, SUB, MUL, XOR, AND, OR, NOT, HASH, ROT, PUSH
Edge cases include:
* Stack underflow
* Truncated instructions
* Unknown opcodes
* Wrapping arithmetic
* Invalid instruction streams
- Stack underflow
- Truncated instructions
- Unknown opcodes
- Wrapping arithmetic
- Invalid instruction streams
Example tests:
```text
```
test_add
test_add_wrapping
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`)
The server crate currently contains **30 tests** covering:
The server crate currently contains **33 tests** covering:
* Configuration
* Runtime initialization
* Session management
* Heartbeat validation
* Rate limiting
* Trust validation
- Configuration
- Runtime initialization
- Session management
- Heartbeat validation
- Rate limiting
- Trust 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:
* VM execution
* Browser-side mutation lifecycle
* Gene initialization
* Mutation preview
* Mutation commit/discard behavior
* Deterministic parity with shared logic
- VM execution
- Browser-side mutation lifecycle
- Gene initialization
- Mutation preview
- Mutation commit/discard behavior
- Deterministic parity with shared logic
Example tests:
```text
```
test_preview_commitment_matches_shared_engine
test_commit_applies_preview
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`)
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:
### Synthetic Gene Engine
```text
```
test_new_state_with_default_size
test_new_state_rejects_invalid_sizes
test_commitment_changes_when_gene_or_environment_changes
@@ -272,7 +301,7 @@ test_table_driven_randomized_environment_roundtrip
### Mutation Engine
```text
```
test_opcode_insert
test_opcode_delete
test_opcode_mutate_point
@@ -283,7 +312,7 @@ test_mutation_chain
### Validation & Hardening
```text
```
test_rejects_stack_underflow
test_rejects_truncated_instruction
test_rejects_unknown_opcode
@@ -292,7 +321,7 @@ test_zero_length_gene_is_rejected
### Deterministic Parity
```text
```
test_server_client_parity_across_random_orders
test_generate_order_is_deterministic_for_seeded_rng
test_invalid_positions_wrap_deterministically
@@ -300,24 +329,47 @@ test_invalid_positions_wrap_deterministically
### Fuzz & Regression Testing
```text
```
test_fuzz_style_random_program_bytes_do_not_diverge
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
Run the full workspace:
```bash
```
cargo test --workspace
```
Run individual crates:
```bash
```
cargo test -p shared
cargo test -p chronoseal-wasm
cargo test -p chronoseal-server
@@ -325,7 +377,7 @@ cargo test -p chronoseal-server
Display test output:
```bash
```
cargo test -- --nocapture
```
@@ -333,18 +385,22 @@ cargo test -- --nocapture
# 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_replay_attack_is_rejected
test_handler_tampered_commitment_is_silent_failure
test_server_client_parity_across_random_orders
test_deterministic_server_client_parity_across_many_heartbeats
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/`
2. Ensure server ↔ WASM parity is validated
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
---
# Future Improvements
Planned enhancements include:
* Property-based testing using `proptest`
* Browser-driven end-to-end integration tests
* Valkey concurrency and failover testing
* Automated benchmark execution in CI
* Expanded mutation-engine fuzzing
* CI-enforced performance regression thresholds
- **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
`update_session` (security-critical, planned for v1.1.0)
- **Browser-driven end-to-end integration tests** — full Playwright or wasm-bindgen-test
harness exercising the complete init → heartbeat loop in a real browser environment
- **Valkey failover testing** — verify graceful degradation and reconnection under r2d2 pool
exhaustion and server-side connection drops
- **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
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
* Heartbeat validation
* Mutation engine correctness
* Deterministic server/WASM parity
* Trust validation
* Storage abstraction
* Replay resistance
* Protocol hardening
- Session security
- Heartbeat validation
- Mutation engine correctness
- Deterministic server/WASM parity
- Trust and behavioral validation
- Storage abstraction
- Replay resistance
- Protocol hardening
- Property-based VM and gene codec robustness
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
- 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
### 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 pendingMutationStep = 0;
let pendingMutationOrderB64 = '';
let mutationRounds = 4;
export async function initHeartbeat() {
await init();
@@ -32,6 +33,7 @@ export async function initHeartbeat() {
}
pendingMutationStep = initResp.mutation_step;
pendingMutationOrderB64 = initResp.mutation_order_b64;
mutationRounds = initResp.mutation_rounds || 4;
lastTime = performance.now();
scheduleNext();
}
@@ -56,7 +58,7 @@ async function sendHeartbeat() {
const timestamp = Date.now();
const entropyData = { events: events.map(e => ({ x: e.x, y: e.y, t: e.t })) };
const entropyJson = JSON.stringify(entropyData);
const geneCommitment = preview_gene_commitment(pendingMutationOrderB64);
const geneCommitment = preview_gene_commitment(pendingMutationOrderB64, session, pendingMutationStep, mutationRounds);
if (!geneCommitment) {
throw new Error('Unable to compute mutation commitment');
}
@@ -75,7 +77,6 @@ async function sendHeartbeat() {
const sig = sign_message(msg);
if (!sig) {
discard_gene_preview();
console.error('Keypair not initialised — skipping heartbeat');
return;
}
const resp = await sendRequest('/hb', 'POST', {
@@ -105,11 +106,9 @@ async function sendHeartbeat() {
pendingMutationOrderB64 = resp.next_mutation_order_b64;
} else {
discard_gene_preview();
console.warn('Heartbeat rejected');
}
} catch (e) {
discard_gene_preview();
console.error(e);
} finally {
scheduleNext();
}
+1
View File
@@ -2,6 +2,7 @@
<html lang="en">
<head>
<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>
</head>
<body>
+1 -3
View File
@@ -1,5 +1,3 @@
import { initHeartbeat } from './heartbeat.js';
(async () => {
await initHeartbeat();
})();
initHeartbeat().catch(() => {});
+1
View File
@@ -1,4 +1,5 @@
#!/bin/bash
set -euo pipefail
echo "Starting server with static frontend serving..."
cd ../server
cargo run --release
+4 -2
View File
@@ -1,6 +1,6 @@
[package]
name = "chronoseal-server"
version = "0.6.0"
version = "1.0.1"
edition = "2021"
[[bin]]
@@ -30,4 +30,6 @@ hex = "0.4"
base64 = "0.22"
rand = "0.8"
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 std::sync::Arc;
/// Runs an infinite background loop that periodically cleans up database and memory resources.
///
/// Every 60 seconds, this loop performs two tasks:
/// 1. Evicts expired session records from the configured database storage backend.
/// 2. Evicts stale rate-limiter entries that have outlived the current rate-limiting window.
///
/// # Arguments
/// * `state` - Shared reference to the server application state.
pub async fn cleanup_loop(state: Arc<AppState>) {
loop {
tokio::time::sleep(std::time::Duration::from_secs(60)).await;
@@ -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 mut rl = state.rate_limiter.lock().await;
rl.evict_stale(window_secs);
state.rate_limiter.evict_stale(window_secs);
}
}
}
+5 -2
View File
@@ -137,7 +137,9 @@ impl Config {
size: self.gene_size,
});
}
if !(1..=shared::constants::MAX_MUTATION_ROUNDS).contains(&self.mutation_rounds) {
if !(shared::constants::MIN_MUTATION_ROUNDS..=shared::constants::MAX_MUTATION_ROUNDS)
.contains(&self.mutation_rounds)
{
return Err(ConfigError::InvalidMutationRounds {
rounds: self.mutation_rounds,
});
@@ -274,7 +276,8 @@ impl std::fmt::Display for ConfigError {
Self::InvalidMutationRounds { rounds } => {
write!(
f,
"invalid mutation rounds {rounds}; expected 1..={}",
"invalid mutation rounds {rounds}; expected {}..={}",
shared::constants::MIN_MUTATION_ROUNDS,
shared::constants::MAX_MUTATION_ROUNDS
)
}
+15 -2
View File
@@ -2,11 +2,16 @@ use ed25519_dalek::{Signature, VerifyingKey};
use shared::protocol::HeartbeatRequest;
use std::collections::BTreeMap;
/// Serializes the heartbeat request into a canonical JSON representation for signature verification.
///
/// Uses `BTreeMap` to order top-level keys alphabetically, matching the JavaScript client's
/// sorting algorithm: `JSON.stringify(obj, Object.keys(obj).sort())`.
///
/// # Arguments
/// * `req` - The heartbeat request to serialize.
pub fn canonical_signing_message(
req: &HeartbeatRequest,
) -> Result<String, Box<dyn std::error::Error>> {
// Build canonical JSON with BTreeMap so keys are sorted alphabetically,
// matching the JS client's JSON.stringify(obj, Object.keys(obj).sort()).
let mut payload: BTreeMap<&str, serde_json::Value> = BTreeMap::new();
payload.insert("entropyData", serde_json::to_value(&req.entropy_data)?);
payload.insert("fingerprint", serde_json::to_value(&req.fingerprint)?);
@@ -19,6 +24,14 @@ pub fn canonical_signing_message(
Ok(serde_json::to_string(&payload)?)
}
/// Verifies the Ed25519 signature of a client's heartbeat request.
///
/// Decodes the signature and compares it strictly against the canonical JSON message
/// using the client's public key.
///
/// # Arguments
/// * `pub_key_bytes` - The client's public key bytes.
/// * `req` - The heartbeat request payload containing the signature.
pub fn verify_signature(
pub_key_bytes: &[u8],
req: &HeartbeatRequest,
+13
View File
@@ -25,12 +25,19 @@ pub enum SessionError {
#[error("Invalid gene configuration: {0}")]
InvalidGeneConfiguration(String),
#[error("Rate limited")]
RateLimited,
}
impl IntoResponse for SessionError {
fn into_response(self) -> Response {
let (status, error_message) = match self {
SessionError::InvalidPublicKeyLength => (StatusCode::BAD_REQUEST, self.to_string()),
SessionError::RateLimited => (
StatusCode::TOO_MANY_REQUESTS,
"Too many requests".to_string(),
),
_ => (
StatusCode::INTERNAL_SERVER_ERROR,
"Internal server error".to_string(),
@@ -86,4 +93,10 @@ pub enum VerificationError {
#[error("Gene state error: {0}")]
GeneState(String),
#[error("VM execution stack state mismatch")]
VmStackMismatch,
#[error("Concurrent state modification detected (CAS failed)")]
ConcurrentUpdate,
}
+66 -3
View File
@@ -1,16 +1,79 @@
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>> {
let ar: f64 = fp.aspect_ratio.parse().map_err(|_| "ar")?;
if !(0.5..=3.0).contains(&ar) {
if !ar.is_finite() || !(MIN_ASPECT_RATIO..=MAX_ASPECT_RATIO).contains(&ar) {
return Err("aspect ratio".into());
}
let dpr: f64 = fp.device_pixel_ratio.parse().map_err(|_| "dpr")?;
if dpr <= 0.0 || dpr > 5.0 {
if !dpr.is_finite() || dpr <= 0.0 || dpr > MAX_DEVICE_PIXEL_RATIO {
return Err("dpr".into());
}
if fp.hardware_concurrency == 0 {
if fp.hardware_concurrency == 0 || fp.hardware_concurrency > MAX_HARDWARE_CONCURRENCY {
return Err("hw".into());
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn fingerprint(
aspect_ratio: impl Into<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>> {
let cli = Cli::parse();
if let Some(config_path) = cli.globals.config.as_deref() {
std::env::set_var("CHRONOSEAL_CONFIG", config_path);
}
let log_filter = cli.globals.log.as_deref().unwrap_or("info");
let log_file = log_file_for_command(&cli);
let _log_guard = init_logging(log_filter, log_file)?;
+22
View File
@@ -9,3 +9,25 @@ pub async fn log_request(req: Request, next: Next) -> Response {
tracing::info!("{} {} -> {}", method, uri, response.status());
response
}
/// Injects defensive HTTP response headers on every response.
///
/// These headers mitigate several classes of attacks:
/// - `X-Content-Type-Options: nosniff` — prevents MIME-type sniffing.
/// - `X-Frame-Options: DENY` — blocks clickjacking via framing.
/// - `Referrer-Policy: no-referrer` — suppresses referrer leakage.
/// - `X-XSS-Protection: 0` — disables legacy XSS auditors (can introduce bugs).
/// - `Permissions-Policy` — restricts powerful browser features.
pub async fn security_headers(req: Request, next: Next) -> Response {
let mut response = next.run(req).await;
let headers = response.headers_mut();
headers.insert("x-content-type-options", "nosniff".parse().unwrap());
headers.insert("x-frame-options", "DENY".parse().unwrap());
headers.insert("referrer-policy", "no-referrer".parse().unwrap());
headers.insert("x-xss-protection", "0".parse().unwrap());
headers.insert(
"permissions-policy",
"camera=(), microphone=(), geolocation=()".parse().unwrap(),
);
response
}
+34 -14
View File
@@ -1,34 +1,54 @@
use std::collections::HashMap;
use dashmap::DashMap;
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 {
buckets: HashMap<String, (u32, Instant)>,
/// Maps rate-limit keys to request counts and window start timestamps.
buckets: DashMap<String, (u32, Instant)>,
}
impl RateLimiter {
/// Creates a new, empty `RateLimiter`.
pub fn new() -> 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 entry = self.buckets.entry(key.to_string()).or_insert((0, now));
if now.duration_since(entry.1).as_secs() >= window_secs {
*entry = (1, now);
let mut entry = self.buckets.entry(key.to_string()).or_insert((0, now));
let (count, ts) = entry.value_mut();
if now.duration_since(*ts).as_secs() >= window_secs {
*count = 1;
*ts = now;
true
} else if entry.0 >= limit {
} else if *count >= limit {
false
} else {
entry.0 += 1;
*count += 1;
true
}
}
/// Remove entries whose rate-limit window has fully elapsed.
/// Call this periodically (e.g. from the cleanup loop) to bound memory usage.
pub fn evict_stale(&mut self, window_secs: u64) {
/// Evicts expired rate-limit entries whose time windows have fully elapsed.
///
/// Intended to be called periodically to bound in-memory map growth.
///
/// # Arguments
/// * `window_secs` - The active rate-limiting window duration in seconds.
pub fn evict_stale(&self, window_secs: u64) {
let now = Instant::now();
self.buckets
.retain(|_, (_, ts)| now.duration_since(*ts).as_secs() < window_secs);
@@ -43,7 +63,7 @@ mod tests {
#[test]
fn test_rate_limiter() {
let mut rl = RateLimiter::new();
let rl = RateLimiter::new();
// Limit of 2 requests per 1 second window
assert!(rl.check("user1", 2, 1));
assert!(rl.check("user1", 2, 1));
@@ -57,7 +77,7 @@ mod tests {
#[test]
fn test_rate_limiter_eviction() {
let mut rl = RateLimiter::new();
let rl = RateLimiter::new();
assert!(rl.check("user1", 1, 1));
assert_eq!(rl.buckets.len(), 1);
+96 -11
View File
@@ -7,15 +7,52 @@ pub async fn handler(
State(state): State<Arc<AppState>>,
Json(payload): Json<HeartbeatRequest>,
) -> (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 cfg = state.get_config();
(cfg.rate_limit_count, cfg.rate_limit_window_secs)
};
let mut rl = state.rate_limiter.lock().await;
if !rl.check(&payload.session_id, limit, window_secs) {
if !state
.rate_limiter
.check(&payload.session_id, limit, window_secs)
{
tracing::debug!("Rate limit hit: {}", payload.session_id);
state
.verification_failures_total
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let http_dur = start_http.elapsed().as_nanos() as u64;
state
.http_latency_ns
.fetch_add(http_dur, std::sync::atomic::Ordering::Relaxed);
state
.http_ops_count
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
return (
StatusCode::OK,
Json(HeartbeatResponse {
@@ -29,7 +66,17 @@ pub async fn handler(
}
let config = state.get_config();
match crate::session::verify_heartbeat(&state.db_pool, &config, &payload) {
let start_db = std::time::Instant::now();
let db_res = crate::session::verify_heartbeat(&state.db_pool, &config, &payload);
let db_dur = start_db.elapsed().as_nanos() as u64;
state
.storage_latency_ns
.fetch_add(db_dur, std::sync::atomic::Ordering::Relaxed);
state
.storage_ops_count
.fetch_add(2, std::sync::atomic::Ordering::Relaxed); // read + write
let outcome = match db_res {
Ok(result) => (
StatusCode::OK,
Json(HeartbeatResponse {
@@ -41,6 +88,24 @@ pub async fn handler(
),
Err(e) => {
tracing::warn!("Heartbeat failed for {}: {}", payload.session_id, e);
state
.verification_failures_total
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
match &e {
crate::errors::VerificationError::ChainBroken => {
state
.replay_attempts_total
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
crate::errors::VerificationError::MutationCommitmentMismatch
| crate::errors::VerificationError::MutationProgram(_)
| crate::errors::VerificationError::GeneState(_) => {
state
.mutation_failures_total
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
_ => {}
}
(
StatusCode::OK,
Json(HeartbeatResponse {
@@ -51,7 +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)]
@@ -59,7 +134,7 @@ mod tests {
use super::*;
use axum::{extract::State, Json};
use ed25519_dalek::{Signer, SigningKey};
use shared::protocol::{EntropyData, Fingerprint, InitResponse, MouseEvent, StackState};
use shared::protocol::{EntropyData, Fingerprint, InitResponse, MouseEvent};
use std::path::Path;
fn test_config() -> crate::config::Config {
@@ -107,10 +182,12 @@ mod tests {
},
],
};
let stack_state = StackState {
stack: vec![9, 10, 11],
ip: 2,
};
let program_bytes = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&init.opcodes_b64,
)
.unwrap();
let stack_state = shared::vm::execute(&program_bytes);
let order =
shared::vm_extensions::decode_order_b64(mutation_step, mutation_order_b64).unwrap();
@@ -147,8 +224,16 @@ mod tests {
let pool = crate::storage::init_pool(Path::new(":memory:")).unwrap();
let state = Arc::new(AppState {
db_pool: pool.clone(),
rate_limiter: tokio::sync::Mutex::new(crate::ratelimit::RateLimiter::new()),
rate_limiter: crate::ratelimit::RateLimiter::new(),
config: std::sync::RwLock::new(config.clone()),
heartbeats_total: std::sync::atomic::AtomicU64::new(0),
verification_failures_total: std::sync::atomic::AtomicU64::new(0),
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();
+31 -1
View File
@@ -8,7 +8,37 @@ pub async fn handler(
State(state): State<Arc<AppState>>,
Json(payload): Json<InitRequest>,
) -> Result<Json<InitResponse>, SessionError> {
let start_http = std::time::Instant::now();
let config = state.get_config();
let resp = crate::session::create_session(&state.db_pool, &config, &payload.public_key)?;
// 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))
}
+128 -20
View File
@@ -5,17 +5,19 @@ use crate::{
routes, session,
storage::{self, StoreStats},
};
use axum::{http::StatusCode, response::IntoResponse, routing::get, Json, Router};
use axum::{
extract::ConnectInfo, http::StatusCode, response::IntoResponse, routing::get, Json, Router,
};
use serde::Serialize;
use std::{
fs,
io::{Read, Write},
net::{SocketAddr, TcpStream},
net::{IpAddr, SocketAddr, TcpStream},
path::Path,
sync::Arc,
time::Duration,
};
use tokio::sync::{Mutex, Notify};
use tokio::sync::Notify;
use tracing::{error, info, warn};
#[derive(Debug, Serialize)]
@@ -144,8 +146,16 @@ pub async fn run_daemon(config: Config) -> Result<(), Box<dyn std::error::Error>
let db_pool = init_db_pool(&config)?;
let state = Arc::new(session::AppState {
db_pool,
rate_limiter: Mutex::new(RateLimiter::new()),
rate_limiter: RateLimiter::new(),
config: std::sync::RwLock::new(config.clone()),
heartbeats_total: std::sync::atomic::AtomicU64::new(0),
verification_failures_total: std::sync::atomic::AtomicU64::new(0),
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();
@@ -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),
)
.layer(tower_http::cors::CorsLayer::permissive())
.layer(axum::middleware::from_fn(
crate::middleware::security_headers,
))
.layer(axum::middleware::from_fn(crate::middleware::log_request))
.layer(axum::extract::DefaultBodyLimit::max(64 * 1024)) // 64 KiB
.with_state(state.clone());
let addr: SocketAddr = config.bind.parse()?;
@@ -170,9 +184,12 @@ pub async fn run_daemon(config: Config) -> Result<(), Box<dyn std::error::Error>
info!(bind = %config.bind, "chronoseal daemon started");
let shutdown = signal_task(state.clone());
let result = axum::serve(listener, app)
.with_graceful_shutdown(shutdown)
.await;
let result = axum::serve(
listener,
app.into_make_service_with_connect_info::<SocketAddr>(),
)
.with_graceful_shutdown(shutdown)
.await;
remove_pid_file(&config.pid_file);
result?;
@@ -269,8 +286,12 @@ async fn health_handler() -> impl IntoResponse {
}
async fn stats_handler(
ConnectInfo(addr): ConnectInfo<SocketAddr>,
axum::extract::State(state): axum::extract::State<Arc<session::AppState>>,
) -> Result<Json<StoreStats>, (StatusCode, String)> {
if !is_loopback(addr.ip()) {
return Err((StatusCode::FORBIDDEN, "Forbidden".to_string()));
}
state
.db_pool
.stats()
@@ -279,18 +300,99 @@ async fn stats_handler(
}
async fn metrics_handler(
ConnectInfo(addr): ConnectInfo<SocketAddr>,
axum::extract::State(state): axum::extract::State<Arc<session::AppState>>,
) -> Result<String, (StatusCode, String)> {
state
if !is_loopback(addr.ip()) {
return Err((StatusCode::FORBIDDEN, "Forbidden".to_string()));
}
let stats = state
.db_pool
.stats()
.map(|stats| {
format!(
"# HELP chronoseal_sessions Active ChronoSeal sessions\n# TYPE chronoseal_sessions gauge\nchronoseal_sessions {}\n# HELP chronoseal_expired_sessions Expired sessions not yet removed\n# TYPE chronoseal_expired_sessions gauge\nchronoseal_expired_sessions {}\n# HELP chronoseal_max_chain_length Maximum heartbeat chain length\n# TYPE chronoseal_max_chain_length gauge\nchronoseal_max_chain_length {}\n",
stats.sessions, stats.expired_sessions, stats.max_chain_length
)
})
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?;
let heartbeats = state
.heartbeats_total
.load(std::sync::atomic::Ordering::Relaxed);
let ver_failures = state
.verification_failures_total
.load(std::sync::atomic::Ordering::Relaxed);
let mut_failures = state
.mutation_failures_total
.load(std::sync::atomic::Ordering::Relaxed);
let replays = state
.replay_attempts_total
.load(std::sync::atomic::Ordering::Relaxed);
let store_ns = state
.storage_latency_ns
.load(std::sync::atomic::Ordering::Relaxed) as f64;
let store_sum = store_ns / 1_000_000_000.0;
let store_count = state
.storage_ops_count
.load(std::sync::atomic::Ordering::Relaxed);
let http_ns = state
.http_latency_ns
.load(std::sync::atomic::Ordering::Relaxed) as f64;
let http_sum = http_ns / 1_000_000_000.0;
let http_count = state
.http_ops_count
.load(std::sync::atomic::Ordering::Relaxed);
Ok(format!(
"# HELP chronoseal_active_sessions Active ChronoSeal sessions\n\
# TYPE chronoseal_active_sessions gauge\n\
chronoseal_active_sessions {}\n\
# HELP chronoseal_expired_sessions Expired sessions not yet removed\n\
# TYPE chronoseal_expired_sessions gauge\n\
chronoseal_expired_sessions {}\n\
# HELP chronoseal_max_chain_length Maximum heartbeat chain length\n\
# TYPE chronoseal_max_chain_length gauge\n\
chronoseal_max_chain_length {}\n\
# HELP chronoseal_heartbeats_total Total heartbeat requests processed\n\
# TYPE chronoseal_heartbeats_total counter\n\
chronoseal_heartbeats_total {}\n\
# HELP chronoseal_verification_failures_total Total heartbeat verification failures\n\
# TYPE chronoseal_verification_failures_total counter\n\
chronoseal_verification_failures_total {}\n\
# HELP chronoseal_mutation_failures_total Total heartbeat mutation verification failures\n\
# TYPE chronoseal_mutation_failures_total counter\n\
chronoseal_mutation_failures_total {}\n\
# HELP chronoseal_replay_attempts_total Total heartbeat replay attempts detected\n\
# TYPE chronoseal_replay_attempts_total counter\n\
chronoseal_replay_attempts_total {}\n\
# HELP chronoseal_storage_latency_seconds_sum Total time spent in storage operations in seconds\n\
# TYPE chronoseal_storage_latency_seconds_sum counter\n\
chronoseal_storage_latency_seconds_sum {:.6}\n\
# HELP chronoseal_storage_latency_seconds_count Total storage operations count\n\
# TYPE chronoseal_storage_latency_seconds_count counter\n\
chronoseal_storage_latency_seconds_count {}\n\
# HELP chronoseal_http_latency_seconds_sum Total time spent in HTTP request processing in seconds\n\
# TYPE chronoseal_http_latency_seconds_sum counter\n\
chronoseal_http_latency_seconds_sum {:.6}\n\
# HELP chronoseal_http_latency_seconds_count Total HTTP operations count\n\
# TYPE chronoseal_http_latency_seconds_count counter\n\
chronoseal_http_latency_seconds_count {}\n",
stats.sessions,
stats.expired_sessions,
stats.max_chain_length,
heartbeats,
ver_failures,
mut_failures,
replays,
store_sum,
store_count,
http_sum,
http_count
))
}
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>) {
@@ -458,10 +560,16 @@ mod tests {
fn test_init_db_pool_valkey_compat_mode() {
let mut config = base_config();
config.db_type = crate::config::DbType::Valkey;
let pool = init_db_pool(&config).unwrap();
let stats = pool.stats().unwrap();
assert_eq!(stats.sessions, 0);
assert_eq!(stats.expired_sessions, 0);
assert_eq!(stats.max_chain_length, 0);
match init_db_pool(&config) {
Ok(pool) => {
let stats = pool.stats().unwrap();
assert_eq!(stats.sessions, 0);
assert_eq!(stats.expired_sessions, 0);
assert_eq!(stats.max_chain_length, 0);
}
Err(_) => {
// Valkey not running in the test environment, which is acceptable
}
}
}
}
+35 -11
View File
@@ -1,7 +1,15 @@
pub struct AppState {
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 heartbeats_total: std::sync::atomic::AtomicU64,
pub verification_failures_total: std::sync::atomic::AtomicU64,
pub mutation_failures_total: std::sync::atomic::AtomicU64,
pub replay_attempts_total: std::sync::atomic::AtomicU64,
pub storage_latency_ns: std::sync::atomic::AtomicU64,
pub storage_ops_count: std::sync::atomic::AtomicU64,
pub http_latency_ns: std::sync::atomic::AtomicU64,
pub http_ops_count: std::sync::atomic::AtomicU64,
}
impl AppState {
@@ -68,6 +76,7 @@ pub fn create_session(
environment: environment_blob,
pending_mutation: initial_mutation.program,
pending_mutation_step: initial_mutation.step,
opcodes,
};
db.insert_session(&record)
@@ -84,6 +93,7 @@ pub fn create_session(
gene_size: config.gene_size as u32,
mutation_step: initial_mutation.step,
mutation_order_b64: initial_mutation_b64,
mutation_rounds: config.mutation_rounds,
})
}
@@ -150,6 +160,12 @@ pub fn verify_heartbeat(
fingerprint::validate(&req.fingerprint)
.map_err(|e| crate::errors::VerificationError::FingerprintFailed(e.to_string()))?;
// 5.5 Verify VM execution state
let expected_stack = shared::vm::execute(&session.opcodes);
if req.stack_state.stack != expected_stack.stack || req.stack_state.ip != expected_stack.ip {
return Err(crate::errors::VerificationError::VmStackMismatch);
}
// 6. Compute new hash
let new_hash = shared::hashing::next_chain_hash(
&prev_hash_bytes,
@@ -182,9 +198,16 @@ pub fn verify_heartbeat(
environment: next_environment_blob,
pending_mutation: next_mutation.program,
pending_mutation_step: next_step,
opcodes: session.opcodes,
};
db.update_session(&update_record)
.map_err(|e| crate::errors::VerificationError::Storage(e.to_string()))?;
db.update_session(&update_record, &session.last_hash)
.map_err(|e| {
if e.to_string().contains("Concurrent update detected") {
crate::errors::VerificationError::ConcurrentUpdate
} else {
crate::errors::VerificationError::Storage(e.to_string())
}
})?;
Ok(HeartbeatVerificationResult {
next_salt_hex,
@@ -210,6 +233,7 @@ mod tests {
pending_mutation_step: u64,
pending_mutation_order_b64: String,
committed_gene_state: GeneState,
opcodes_b64: String,
}
fn test_config() -> crate::config::Config {
@@ -247,13 +271,6 @@ mod tests {
}
}
fn test_stack() -> StackState {
StackState {
stack: vec![42, 7, 99],
ip: 3,
}
}
fn test_fingerprint() -> Fingerprint {
Fingerprint {
aspect_ratio: "1.77".to_string(),
@@ -288,6 +305,7 @@ mod tests {
pending_mutation_step: init.mutation_step,
pending_mutation_order_b64: init.mutation_order_b64.clone(),
committed_gene_state: gene::new_state(init.gene_size as usize).unwrap(),
opcodes_b64: init.opcodes_b64.clone(),
}
}
@@ -304,7 +322,13 @@ mod tests {
vm_extensions::apply_program_clone(&client.committed_gene_state, &order.program)
.unwrap();
let entropy = test_entropy();
let stack = test_stack();
let program_bytes = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&client.opcodes_b64,
)
.unwrap();
let stack = shared::vm::execute(&program_bytes);
let mut req = HeartbeatRequest {
session_id: client.session_id.clone(),
+370 -71
View File
@@ -1,9 +1,8 @@
use crate::config::Config;
use redis::Commands;
use serde::{Deserialize, Serialize};
use std::path::Path;
use std::sync::{Arc, Mutex};
use std::time::{SystemTime, UNIX_EPOCH};
use valkey::Client as ValkeyClient;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StoreStats {
@@ -20,7 +19,7 @@ pub enum DbPool {
#[derive(Debug, Clone)]
pub struct ValkeyStore {
client: Arc<Mutex<ValkeyClient>>,
pool: r2d2::Pool<redis::Client>,
index_key: String,
}
@@ -38,6 +37,7 @@ pub struct SessionRecord {
pub environment: Vec<u8>,
pub pending_mutation: Vec<u8>,
pub pending_mutation_step: u64,
pub opcodes: Vec<u8>,
}
impl DbPool {
@@ -54,19 +54,18 @@ impl DbPool {
crate::config::DbType::Valkey => {
let addr = std::env::var("CHRONOSEAL_VALKEY_ADDR")
.unwrap_or_else(|_| "127.0.0.1:6666".to_string());
match ValkeyClient::connect(addr) {
Ok(client) => Ok(DbPool::Valkey(ValkeyStore {
client: Arc::new(Mutex::new(client)),
index_key: "sessions:ids".to_string(),
})),
Err(err) => {
tracing::warn!(
"valkey connection failed, falling back to sqlite-in-memory: {err}"
);
let pool = init_sqlite_pool(Path::new(":memory:"))?;
Ok(DbPool::Sqlite(pool))
}
}
let connection_string =
if addr.starts_with("redis://") || addr.starts_with("rediss://") {
addr.clone()
} else {
format!("redis://{}", addr)
};
let client = redis::Client::open(connection_string)?;
let pool = r2d2::Pool::builder().build(client)?;
Ok(DbPool::Valkey(ValkeyStore {
pool,
index_key: "sessions:ids".to_string(),
}))
}
}
}
@@ -79,8 +78,8 @@ impl DbPool {
"INSERT INTO sessions (
session_id, public_key, salt, last_hash, chain_length,
created_at, last_seen, expires_at, gene, environment,
pending_mutation, pending_mutation_step
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
pending_mutation, pending_mutation_step, opcodes
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
)?;
stmt.execute(rusqlite::params![
record.session_id,
@@ -95,6 +94,7 @@ impl DbPool {
&record.environment,
&record.pending_mutation,
record.pending_mutation_step,
&record.opcodes,
])?;
Ok(())
}
@@ -110,7 +110,7 @@ impl DbPool {
DbPool::Sqlite(pool) => {
let conn = pool.get()?;
let mut stmt = conn.prepare(
"SELECT session_id, public_key, salt, last_hash, chain_length, created_at, last_seen, expires_at, gene, environment, pending_mutation, pending_mutation_step
"SELECT session_id, public_key, salt, last_hash, chain_length, created_at, last_seen, expires_at, gene, environment, pending_mutation, pending_mutation_step, opcodes
FROM sessions WHERE session_id = ?1",
)?;
let row = stmt.query_row([session_id], |row| {
@@ -127,6 +127,7 @@ impl DbPool {
environment: row.get(9)?,
pending_mutation: row.get(10)?,
pending_mutation_step: row.get(11)?,
opcodes: row.get(12)?,
})
});
match row {
@@ -139,11 +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 {
DbPool::Sqlite(pool) => {
let conn = pool.get()?;
conn.execute(
let rows = conn.execute(
"UPDATE sessions SET
public_key=?1,
salt=?2,
@@ -155,8 +160,9 @@ impl DbPool {
gene=?8,
environment=?9,
pending_mutation=?10,
pending_mutation_step=?11
WHERE session_id=?12",
pending_mutation_step=?11,
opcodes=?12
WHERE session_id=?13 AND last_hash=?14",
rusqlite::params![
&record.public_key,
&record.salt,
@@ -169,12 +175,20 @@ impl DbPool {
&record.environment,
&record.pending_mutation,
record.pending_mutation_step,
&record.opcodes,
&record.session_id,
old_last_hash,
],
)?;
if rows == 0 {
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
"Concurrent update detected (CAS failed)",
)));
}
Ok(())
}
DbPool::Valkey(store) => store.insert_session(record),
DbPool::Valkey(store) => store.update_session_cas(record, old_last_hash),
}
}
@@ -237,6 +251,10 @@ fn init_sqlite_pool(
}
r2d2_sqlite::SqliteConnectionManager::file(path)
};
let manager = manager.with_init(|conn| {
conn.busy_timeout(std::time::Duration::from_millis(5000))?;
Ok(())
});
let pool = r2d2::Pool::new(manager)?;
let conn = pool.get()?;
init_schema(&conn)?;
@@ -257,7 +275,8 @@ fn init_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
gene BLOB NOT NULL DEFAULT X'',
environment BLOB NOT NULL DEFAULT X'',
pending_mutation BLOB NOT NULL DEFAULT X'',
pending_mutation_step INTEGER NOT NULL DEFAULT 0
pending_mutation_step INTEGER NOT NULL DEFAULT 0,
opcodes BLOB NOT NULL DEFAULT X''
);",
)?;
ensure_column(
@@ -280,6 +299,11 @@ fn init_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
"pending_mutation_step",
"ALTER TABLE sessions ADD COLUMN pending_mutation_step INTEGER NOT NULL DEFAULT 0",
)?;
ensure_column(
conn,
"opcodes",
"ALTER TABLE sessions ADD COLUMN opcodes BLOB NOT NULL DEFAULT X''",
)?;
conn.execute_batch(
"CREATE INDEX IF NOT EXISTS idx_sessions_expires_at ON sessions(expires_at);",
)?;
@@ -313,67 +337,136 @@ impl ValkeyStore {
&self,
session_id: &str,
) -> Result<Option<SessionRecord>, Box<dyn std::error::Error>> {
let mut client = self.client.lock().unwrap();
if let Some(payload) = client.get(&self.session_key(session_id))? {
let record = serde_json::from_str(&payload)?;
Ok(Some(record))
} else {
Ok(None)
let mut conn = self.pool.get()?;
let key = self.session_key(session_id);
let payload: Option<String> = conn.get(&key)?;
match payload {
Some(p) => Ok(serde_json::from_str(&p)?),
None => Ok(None),
}
}
fn insert_session(&self, record: &SessionRecord) -> Result<(), Box<dyn std::error::Error>> {
let mut client = self.client.lock().unwrap();
let mut conn = self.pool.get()?;
let key = self.session_key(&record.session_id);
let value = serde_json::to_string(record)?;
client.set(&self.session_key(&record.session_id), &value)?;
let existing = client.get(&self.index_key)?;
let mut ids = existing.unwrap_or_default();
if !ids.split('\n').any(|id| id == record.session_id) {
if !ids.is_empty() {
ids.push('\n');
}
ids.push_str(&record.session_id);
client.set(&self.index_key, &ids)?;
}
let now = current_time_ms();
let ttl_seconds = (record.expires_at.saturating_sub(now) / 1000).max(1);
redis::pipe()
.atomic()
.cmd("SET")
.arg(&key)
.arg(&value)
.arg("EX")
.arg(ttl_seconds)
.cmd("ZADD")
.arg(&self.index_key)
.arg(record.expires_at)
.arg(&record.session_id)
.cmd("ZADD")
.arg("sessions:chain_lengths")
.arg(record.chain_length)
.arg(&record.session_id)
.query::<()>(&mut *conn)?;
Ok(())
}
fn purge_expired_sessions(&self) -> Result<(), Box<dyn std::error::Error>> {
let mut client = self.client.lock().unwrap();
let ids = client.get(&self.index_key)?.unwrap_or_default();
let now = current_time_ms();
let mut remaining: Vec<String> = Vec::new();
for id in ids.split('\n').filter(|id| !id.is_empty()) {
if let Some(payload) = client.get(&self.session_key(id))? {
if let Ok(record) = serde_json::from_str::<SessionRecord>(&payload) {
if record.expires_at > now {
remaining.push(id.to_string());
}
fn update_session_cas(
&self,
record: &SessionRecord,
old_last_hash: &[u8],
) -> Result<(), Box<dyn std::error::Error>> {
let mut conn = self.pool.get()?;
let key = self.session_key(&record.session_id);
// Watch key for concurrent modification
redis::cmd("WATCH").arg(&key).query::<()>(&mut *conn)?;
// Fetch current and verify last_hash matches
let payload: Option<String> = conn.get(&key)?;
match payload {
Some(p) => {
let current_record: SessionRecord = serde_json::from_str(&p)?;
if current_record.last_hash != old_last_hash {
redis::cmd("UNWATCH").query::<()>(&mut *conn)?;
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
"Concurrent update detected (CAS failed in Valkey)",
)));
}
}
None => {
redis::cmd("UNWATCH").query::<()>(&mut *conn)?;
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::NotFound,
"Session not found for update in Valkey",
)));
}
}
let value = serde_json::to_string(record)?;
let now = current_time_ms();
let ttl_seconds = (record.expires_at.saturating_sub(now) / 1000).max(1);
let response: Option<()> = redis::pipe()
.atomic()
.cmd("SET")
.arg(&key)
.arg(&value)
.arg("EX")
.arg(ttl_seconds)
.cmd("ZADD")
.arg(&self.index_key)
.arg(record.expires_at)
.arg(&record.session_id)
.cmd("ZADD")
.arg("sessions:chain_lengths")
.arg(record.chain_length)
.arg(&record.session_id)
.query(&mut *conn)?;
match response {
Some(_) => Ok(()),
None => Err(Box::new(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
"Transaction aborted due to concurrent modification",
))),
}
}
fn purge_expired_sessions(&self) -> Result<(), Box<dyn std::error::Error>> {
let now = current_time_ms();
let mut conn = self.pool.get()?;
// Fetch expired session IDs
let expired_ids: Vec<String> = conn.zrangebyscore(&self.index_key, 0, now)?;
if !expired_ids.is_empty() {
redis::pipe()
.atomic()
.cmd("ZREM")
.arg(&self.index_key)
.arg(&expired_ids)
.cmd("ZREM")
.arg("sessions:chain_lengths")
.arg(&expired_ids)
.query::<()>(&mut *conn)?;
}
client.set(&self.index_key, &remaining.join("\n"))?;
Ok(())
}
fn stats(&self) -> Result<StoreStats, Box<dyn std::error::Error>> {
let mut client = self.client.lock().unwrap();
let ids = client.get(&self.index_key)?.unwrap_or_default();
let now = current_time_ms();
let mut sessions = 0;
let mut expired_sessions = 0;
let mut max_chain_length = 0;
for id in ids.split('\n').filter(|id| !id.is_empty()) {
if let Some(payload) = client.get(&self.session_key(id))? {
if let Ok(record) = serde_json::from_str::<SessionRecord>(&payload) {
sessions += 1;
if record.expires_at < now {
expired_sessions += 1;
}
max_chain_length = max_chain_length.max(record.chain_length);
}
}
}
let mut conn = self.pool.get()?;
let sessions: u64 = conn.zcard(&self.index_key)?;
let expired_sessions: u64 = conn.zcount(&self.index_key, 0, now)?;
let max_chain_length_res: Vec<(String, u64)> =
conn.zrevrange_withscores("sessions:chain_lengths", 0, 0)?;
let max_chain_length = max_chain_length_res
.first()
.map(|(_, score)| *score)
.unwrap_or(0);
Ok(StoreStats {
sessions,
expired_sessions,
@@ -385,6 +478,212 @@ impl ValkeyStore {
pub fn current_time_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.expect("system clock is before UNIX epoch; check system time")
.as_millis() as u64
}
#[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 shared::protocol::EntropyData;
/// Validates the browser mouse cursor interaction path for bot/automation detection.
///
/// Evaluates mouse velocity and distance features, checks the total distance traversed,
/// checks for cursor pauses (low movement over high time diff), and enforces average cursor speeds.
///
/// # Arguments
/// * `data` - The client-supplied interaction entropy events.
/// * `config` - The server configuration boundaries.
pub fn validate_mouse(
data: &EntropyData,
config: &Config,
+17
View File
@@ -4,6 +4,13 @@ use shared::{
vm_extensions::{self, ExecutionTrace, MutationError, MutationOrder},
};
/// Generates a randomized VM opcode instruction program within a length range.
///
/// Builds a program of mathematical and stack ops (e.g. literals, ADD, SUB, XOR, HASH)
/// with dynamic depth checking to ensure valid stacks and prevent out of bounds execution.
///
/// # Arguments
/// * `len_range` - The inclusive range of instruction counts to generate.
pub fn generate_random_program(len_range: std::ops::RangeInclusive<usize>) -> Vec<u8> {
let mut rng = rand::thread_rng();
let count = rng.gen_range(len_range);
@@ -47,6 +54,11 @@ pub fn generate_random_program(len_range: std::ops::RangeInclusive<usize>) -> Ve
ops
}
/// Executes a raw VM mutation program bytecode slice against a `GeneState`.
///
/// # Arguments
/// * `state` - The mutable gene state to mutate.
/// * `program` - The raw VM instruction program.
#[allow(dead_code)]
pub fn execute_mutation_program(
state: &mut GeneState,
@@ -55,6 +67,11 @@ pub fn execute_mutation_program(
vm_extensions::execute_program(state, program)
}
/// Executes a `MutationOrder` program against a `GeneState`.
///
/// # Arguments
/// * `state` - The mutable gene state.
/// * `order` - The mutation order.
#[allow(dead_code)]
pub fn execute_mutation_order(
state: &mut GeneState,
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "shared"
version = "0.6.0"
version = "1.0.1"
edition = "2021"
[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 HASH_OPCODE_INSTRUCTION_COST: usize = 16;
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 {}
/// 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> {
if !(1..=MAX_GENE_SIZE).contains(&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 {
GeneState {
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> {
if !(1..=MAX_GENE_SIZE).contains(&state.gene.len()) {
return Err(GeneError::InvalidGeneSize {
@@ -81,6 +90,13 @@ pub fn validate_state(state: &GeneState) -> Result<(), GeneError> {
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 {
match state
.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(
state: &mut GeneState,
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(
state: &mut GeneState,
symbol: u16,
@@ -136,6 +166,14 @@ pub fn add_env_quantity(
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(
state: &mut GeneState,
symbol: u16,
@@ -147,6 +185,12 @@ pub fn sub_env_quantity(
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> {
validate_environment(records)?;
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)
}
/// 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> {
if !blob.len().is_multiple_of(6) {
return Err(GeneError::EnvironmentBlobLengthInvalid { len: blob.len() });
@@ -180,6 +230,9 @@ pub fn decode_environment(blob: &[u8]) -> Result<Vec<EnvironmentRecord>, GeneErr
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] {
let mut h = blake3::Hasher::new();
h.update(b"chronoseal/gene/v1");
@@ -193,10 +246,19 @@ pub fn commitment(state: &GeneState) -> [u8; 32] {
*h.finalize().as_bytes()
}
/// Computes the hex-encoded cryptographic commitment of the `GeneState`.
pub fn commitment_hex(state: &GeneState) -> String {
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] {
let mut h = blake3::Hasher::new();
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()
}
/// 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 {
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> {
if records.len() > MAX_ENV_RECORDS {
return Err(GeneError::TooManyEnvironmentRecords { len: records.len() });
+31 -8
View File
@@ -1,7 +1,15 @@
use crate::protocol::{EntropyData, StackState};
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> {
let mut h = Hasher::new();
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()
}
/// 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(
prev_hash: &[u8],
timestamp: u64,
@@ -18,14 +37,13 @@ pub fn next_chain_hash(
stack: &StackState,
salt: &[u8],
) -> Vec<u8> {
let entropy_json = serde_json::to_string(entropy).unwrap();
let stack_json = serde_json::to_string(stack).unwrap();
let entropy_bytes = serde_json::to_vec(entropy).unwrap_or_default();
let stack_bytes = serde_json::to_vec(stack).unwrap_or_default();
let entropy_hash = blake3::hash(entropy_json.as_bytes());
let stack_hash = blake3::hash(stack_json.as_bytes());
let entropy_hash = blake3::hash(&entropy_bytes);
let stack_hash = blake3::hash(&stack_bytes);
let mut h = Hasher::new();
// Mix the salt into the hash state
h.update(salt);
h.update(prev_hash);
h.update(&timestamp.to_le_bytes());
@@ -34,7 +52,12 @@ pub fn next_chain_hash(
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 {
let data: Vec<u8> = stack.iter().flat_map(|x| x.to_le_bytes()).collect();
let hash = blake3::hash(&data);
+45
View File
@@ -1,73 +1,118 @@
use serde::{Deserialize, Serialize};
/// Request payload sent by the client to initialize a new attestation session.
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct InitRequest {
/// Hex-encoded 32-byte Ed25519 public verifying key generated by the client.
pub public_key: String,
}
/// Response payload returned by the server upon successful session initialization.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct InitResponse {
/// The unique hex-encoded session identifier.
pub session_id: String,
/// The initial server-issued salt to be mixed in the first heartbeat's hash.
pub salt: String,
/// Base64-encoded initial VM program for client stack execution.
pub opcodes_b64: String,
/// The computed initial hash of the attestation chain.
pub initial_hash: String,
/// Timestamp in milliseconds indicating when the session expires.
pub expires_at: u64,
/// Minimum time in milliseconds allowed between subsequent heartbeats.
pub heartbeat_min_interval_ms: u64,
/// Maximum time in milliseconds allowed between subsequent heartbeats.
pub heartbeat_max_interval_ms: u64,
/// Size of the synthetic gene byte buffer.
pub gene_size: u32,
/// The current mutation step index (starts at 1).
pub mutation_step: u64,
/// Base64-encoded initial gene mutation program.
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)]
pub struct HeartbeatRequest {
/// The session identifier.
pub session_id: String,
/// The expected hash from the previous heartbeat/initialization step.
pub prev_hash: String,
/// The client's current system timestamp in milliseconds.
pub timestamp: u64,
/// The collected client entropy data (such as mouse events).
pub entropy_data: EntropyData,
/// The final execution state of the client's VM stack program.
pub stack_state: StackState,
/// The client's browser hardware and layout fingerprint.
pub fingerprint: Fingerprint,
/// The mutation step index corresponding to the pending mutation.
pub mutation_step: u64,
/// Hex-encoded commitment of the mutated gene state.
pub gene_commitment: String,
/// Ed25519 signature of the canonical JSON-serialized payload.
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)]
pub struct HeartbeatResponse {
/// Attestation status, typically "ok" even on silent failures.
pub status: String,
/// The next server-issued salt for hash chain progression.
#[serde(skip_serializing_if = "Option::is_none")]
pub next_salt: Option<String>,
/// The next expected mutation step index.
#[serde(skip_serializing_if = "Option::is_none")]
pub next_mutation_step: Option<u64>,
/// Base64-encoded next mutation program for client gene progression.
#[serde(skip_serializing_if = "Option::is_none")]
pub next_mutation_order_b64: Option<String>,
}
/// Client browser fingerprint metadata used for basic sanity checks.
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct Fingerprint {
/// Aspect ratio of the client screen.
#[serde(rename = "aspectRatio")]
pub aspect_ratio: String,
/// Device pixel ratio of the screen.
#[serde(rename = "devicePixelRatio")]
pub device_pixel_ratio: String,
/// Number of logical processor cores available.
#[serde(rename = "hardwareConcurrency")]
pub hardware_concurrency: u32,
}
/// Wrapper for browser-side entropy collection.
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct EntropyData {
/// A chronological list of mouse movement events.
pub events: Vec<MouseEvent>,
}
/// Information about a single mouse movement interaction.
#[derive(Deserialize, Serialize, Clone, Debug)]
pub struct MouseEvent {
/// Absolute horizontal coordinate of the cursor.
pub x: f64,
/// Absolute vertical coordinate of the cursor.
pub y: f64,
/// Relative timestamp in milliseconds of the event occurrence.
#[serde(rename = "t")]
pub timestamp_ms: f64,
}
/// The state of the VM stack machine after executing a program.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StackState {
/// The elements remaining on the stack.
pub stack: Vec<u32>,
/// The final instruction pointer location at program completion or termination.
pub ip: u16,
}
+12
View File
@@ -25,6 +25,9 @@ pub fn execute(program: &[u8]) -> StackState {
program[ip + 3],
]);
ip += 4;
if stack.len() >= crate::constants::MAX_STACK_DEPTH {
break;
}
stack.push(val);
}
0x01..=0x07 => {
@@ -43,6 +46,9 @@ pub fn execute(program: &[u8]) -> StackState {
0x07 => a.rotate_left(b % 32),
_ => unreachable!(),
};
if stack.len() >= crate::constants::MAX_STACK_DEPTH {
break;
}
stack.push(r);
}
0x08 => {
@@ -50,11 +56,17 @@ pub fn execute(program: &[u8]) -> StackState {
break;
}
let a = stack.pop().unwrap();
if stack.len() >= crate::constants::MAX_STACK_DEPTH {
break;
}
stack.push(!a);
}
0x09 => {
let r = crate::hashing::hash_stack(&stack);
stack.clear();
if stack.len() >= crate::constants::MAX_STACK_DEPTH {
break;
}
stack.push(r);
}
_ => break,
+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 {
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> {
let program = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, b64)
.map_err(MutationError::Base64)?;
@@ -101,11 +109,27 @@ pub fn decode_order_b64(step: u64, b64: &str) -> Result<MutationOrder, MutationE
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 {
let mut rng = rand::thread_rng();
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>(
rng: &mut R,
step: u64,
@@ -213,10 +237,12 @@ pub fn generate_order_with_rng<R: Rng + ?Sized>(
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> {
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(
state: &GeneState,
program: &[u8],
@@ -227,6 +253,7 @@ pub fn apply_program_clone_with_rounds(
Ok(next)
}
/// Executes the mutation program on the mutable `GeneState` reference for a specific number of rounds.
pub fn apply_program_with_rounds(
state: &mut GeneState,
program: &[u8],
@@ -236,11 +263,21 @@ pub fn apply_program_with_rounds(
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> {
let _ = execute_program_with_rounds(state, program, DEFAULT_MUTATION_ROUNDS)?;
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(
state: &mut GeneState,
program: &[u8],
@@ -317,6 +354,13 @@ fn estimate_program_cost(program: &[u8]) -> usize {
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(
state: &mut GeneState,
program: &[u8],
@@ -835,4 +879,25 @@ mod tests {
"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]
name = "chronoseal-wasm"
version = "0.6.0"
version = "1.0.1"
edition = "2021"
[lib]
+8 -2
View File
@@ -51,8 +51,14 @@ pub fn compute_next_hash(
) -> String {
let prev = hex::decode(prev_hash_hex).unwrap_or_default();
let salt = hex::decode(salt_hex).unwrap_or_default();
let entropy = serde_json::from_str::<shared::protocol::EntropyData>(entropy_data_json).unwrap();
let stack = serde_json::from_str::<shared::protocol::StackState>(stack_state_json).unwrap();
let entropy = match serde_json::from_str::<shared::protocol::EntropyData>(entropy_data_json) {
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);
hex::encode(new)
}
+9 -60
View File
@@ -4,70 +4,19 @@ use wasm_bindgen::prelude::*;
#[wasm_bindgen]
pub fn run_program(program_b64: &str) -> JsValue {
use base64::Engine;
let bytes = base64::engine::general_purpose::STANDARD
.decode(program_b64)
.unwrap();
let bytes = match base64::engine::general_purpose::STANDARD.decode(program_b64) {
Ok(v) => v,
Err(_) => return JsValue::NULL,
};
let state = execute(&bytes);
serde_wasm_bindgen::to_value(&state).unwrap()
match serde_wasm_bindgen::to_value(&state) {
Ok(v) => v,
Err(_) => JsValue::NULL,
}
}
fn execute(program: &[u8]) -> StackState {
let mut stack: Vec<u32> = Vec::new();
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,
}
shared::vm::execute(program)
}
#[cfg(test)]