Initial ChronoSeal release
This commit is contained in:
commit
a66debdece
48 files changed
+2362
No files matched your search
@@ -0,0 +1,20 @@
|
||||
[package]
|
||||
name = "antibot-server"
|
||||
version = "0.2.0"
|
||||
edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
shared = { path = "../shared" }
|
||||
axum = "0.7"
|
||||
tokio = { version = "1", features = ["full"] }
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
rusqlite = { version = "0.31", features = ["bundled"] }
|
||||
tracing = "0.1"
|
||||
tracing-subscriber = "0.3"
|
||||
tower = "0.4"
|
||||
tower-http = { version = "0.5", features = ["cors", "fs"] }
|
||||
hex = "0.4"
|
||||
base64 = "0.22"
|
||||
rand = "0.8"
|
||||
ed25519-dalek = "2"
|
||||
@@ -0,0 +1,14 @@
|
||||
use std::sync::Arc;
|
||||
use crate::session::AppState;
|
||||
|
||||
pub async fn cleanup_loop(state: Arc<AppState>) {
|
||||
loop {
|
||||
tokio::time::sleep(std::time::Duration::from_secs(60)).await;
|
||||
let db = state.db.lock().await; // this is infallible
|
||||
let now = crate::storage::current_time_ms();
|
||||
let _ = db.execute(
|
||||
"DELETE FROM sessions WHERE expires_at < ?1",
|
||||
rusqlite::params![now],
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
use ed25519_dalek::{VerifyingKey, Signature};
|
||||
use shared::protocol::HeartbeatRequest;
|
||||
|
||||
pub fn verify_signature(
|
||||
pub_key_bytes: &[u8],
|
||||
req: &HeartbeatRequest,
|
||||
) -> Result<(), Box<dyn std::error::Error>> {
|
||||
let pk = VerifyingKey::from_bytes(
|
||||
&pub_key_bytes.try_into().map_err(|_| "invalid pubkey")?,
|
||||
)?;
|
||||
let sig_bytes = hex::decode(&req.signature)?;
|
||||
let sig = Signature::from_slice(&sig_bytes)?;
|
||||
|
||||
// Build canonical JSON exactly as client signed (sorted keys, no extra spaces)
|
||||
let payload = serde_json::json!({
|
||||
"sessionId": req.session_id,
|
||||
"prevHash": req.prev_hash,
|
||||
"timestamp": req.timestamp,
|
||||
"entropyData": req.entropy_data,
|
||||
"stackState": req.stack_state,
|
||||
"fingerprint": req.fingerprint,
|
||||
});
|
||||
let message = serde_json::to_string(&payload)?;
|
||||
|
||||
pk.verify_strict(message.as_bytes(), &sig)?;
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,10 @@
|
||||
use shared::protocol::Fingerprint;
|
||||
|
||||
pub fn validate(fp: &Fingerprint) -> Result<(), Box<dyn std::error::Error>> {
|
||||
let ar: f64 = fp.aspect_ratio.parse().map_err(|_| "ar")?;
|
||||
if ar < 0.5 || ar > 3.0 { return Err("aspect ratio".into()); }
|
||||
let dpr: f64 = fp.device_pixel_ratio.parse().map_err(|_| "dpr")?;
|
||||
if dpr <= 0.0 || dpr > 5.0 { return Err("dpr".into()); }
|
||||
if fp.hardware_concurrency == 0 { return Err("hw".into()); }
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,47 @@
|
||||
mod cleanup;
|
||||
mod crypto;
|
||||
mod fingerprint;
|
||||
mod middleware;
|
||||
mod ratelimit;
|
||||
mod routes;
|
||||
mod session;
|
||||
mod storage;
|
||||
mod trust;
|
||||
mod vm;
|
||||
|
||||
use axum::Router;
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::Mutex;
|
||||
use tracing::info;
|
||||
|
||||
use session::AppState;
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() {
|
||||
tracing_subscriber::fmt::init();
|
||||
|
||||
let conn = storage::init_db().expect("DB init");
|
||||
let state = Arc::new(AppState {
|
||||
db: Mutex::new(conn),
|
||||
rate_limiter: Mutex::new(ratelimit::RateLimiter::new(
|
||||
shared::constants::RATE_LIMIT_COUNT,
|
||||
shared::constants::RATE_LIMIT_WINDOW_SECS,
|
||||
)),
|
||||
});
|
||||
|
||||
// Periodic cleanup
|
||||
let bg_state = state.clone();
|
||||
tokio::spawn(async move { cleanup::cleanup_loop(bg_state).await });
|
||||
|
||||
let app = Router::new()
|
||||
.route("/init", axum::routing::post(routes::init::handler))
|
||||
.route("/hb", axum::routing::post(routes::heartbeat::handler))
|
||||
.nest_service("/", tower_http::services::ServeDir::new("../frontend"))
|
||||
.layer(tower_http::cors::CorsLayer::permissive())
|
||||
.layer(axum::middleware::from_fn(middleware::log_request))
|
||||
.with_state(state);
|
||||
|
||||
let listener = tokio::net::TcpListener::bind("0.0.0.0:3000").await.unwrap();
|
||||
info!("Server running on :3000");
|
||||
axum::serve(listener, app).await.unwrap();
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
use axum::extract::Request;
|
||||
use axum::middleware::Next;
|
||||
use axum::response::Response;
|
||||
|
||||
pub async fn log_request(req: Request, next: Next) -> Response {
|
||||
let method = req.method().clone();
|
||||
let uri = req.uri().clone();
|
||||
let response = next.run(req).await;
|
||||
tracing::info!("{} {} -> {}", method, uri, response.status());
|
||||
response
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
use std::collections::HashMap;
|
||||
use std::time::Instant;
|
||||
|
||||
pub struct RateLimiter {
|
||||
buckets: HashMap<String, (u32, Instant)>,
|
||||
limit: u32,
|
||||
window_secs: u64,
|
||||
}
|
||||
|
||||
impl RateLimiter {
|
||||
pub fn new(limit: u32, window_secs: u64) -> Self {
|
||||
Self { buckets: HashMap::new(), limit, window_secs }
|
||||
}
|
||||
pub fn check(&mut self, key: &str) -> 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() >= self.window_secs {
|
||||
*entry = (1, now);
|
||||
true
|
||||
} else if entry.0 >= self.limit {
|
||||
false
|
||||
} else {
|
||||
entry.0 += 1;
|
||||
true
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
use axum::{extract::State, http::StatusCode, Json};
|
||||
use std::sync::Arc;
|
||||
use shared::protocol::{HeartbeatRequest, HeartbeatResponse};
|
||||
use crate::session::AppState;
|
||||
|
||||
pub async fn handler(
|
||||
State(state): State<Arc<AppState>>,
|
||||
Json(payload): Json<HeartbeatRequest>,
|
||||
) -> (StatusCode, Json<HeartbeatResponse>) {
|
||||
// Rate limiting
|
||||
{
|
||||
let mut rl = state.rate_limiter.lock().await;
|
||||
if !rl.check(&payload.session_id) {
|
||||
tracing::debug!("Rate limit hit: {}", payload.session_id);
|
||||
return (StatusCode::OK, Json(HeartbeatResponse { status: "ok".into(), next_salt: None }));
|
||||
}
|
||||
}
|
||||
|
||||
let db = state.db.lock().await;
|
||||
match crate::session::verify_heartbeat(&db, &payload) {
|
||||
Ok(next_salt) => (
|
||||
StatusCode::OK,
|
||||
Json(HeartbeatResponse { status: "ok".into(), next_salt: Some(next_salt) }),
|
||||
),
|
||||
Err(e) => {
|
||||
tracing::warn!("Heartbeat failed for {}: {}", payload.session_id, e);
|
||||
(StatusCode::OK, Json(HeartbeatResponse { status: "ok".into(), next_salt: None }))
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
use axum::{extract::State, http::StatusCode, Json};
|
||||
use std::sync::Arc;
|
||||
use shared::protocol::{InitRequest, InitResponse};
|
||||
use crate::session::AppState;
|
||||
|
||||
pub async fn handler(
|
||||
State(state): State<Arc<AppState>>,
|
||||
Json(payload): Json<InitRequest>,
|
||||
) -> Result<Json<InitResponse>, (StatusCode, String)> {
|
||||
let db = state.db.lock().await;
|
||||
crate::session::create_session(&db, &payload.public_key)
|
||||
.map(Json)
|
||||
.map_err(|e| {
|
||||
tracing::error!("Init error: {}", e);
|
||||
(StatusCode::INTERNAL_SERVER_ERROR, "Internal".into())
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,2 @@
|
||||
pub mod init;
|
||||
pub mod heartbeat;
|
||||
@@ -0,0 +1,98 @@
|
||||
pub struct AppState {
|
||||
pub db: tokio::sync::Mutex<rusqlite::Connection>,
|
||||
pub rate_limiter: tokio::sync::Mutex<crate::ratelimit::RateLimiter>,
|
||||
}
|
||||
|
||||
use rusqlite::params;
|
||||
use shared::protocol::{HeartbeatRequest, InitResponse};
|
||||
use crate::{crypto, trust, fingerprint, vm, storage};
|
||||
|
||||
pub fn create_session(
|
||||
conn: &rusqlite::Connection,
|
||||
pub_key_hex: &str,
|
||||
) -> Result<InitResponse, Box<dyn std::error::Error>> {
|
||||
let pub_key = hex::decode(pub_key_hex)?;
|
||||
if pub_key.len() != shared::constants::SESSION_ID_LEN {
|
||||
return Err("invalid pubkey len".into());
|
||||
}
|
||||
let session_id = hex::encode(rand::random::<[u8; shared::constants::SESSION_ID_LEN]>());
|
||||
let salt = rand::random::<[u8; shared::constants::SALT_LEN]>();
|
||||
let now = storage::current_time_ms();
|
||||
let expires_at = now + (shared::constants::EXPIRATION_MINUTES as u64) * 60 * 1000;
|
||||
|
||||
let initial_hash = shared::hashing::initial_hash(&session_id, &pub_key, &salt);
|
||||
|
||||
conn.execute(
|
||||
"INSERT INTO sessions (session_id, public_key, salt, last_hash, created_at, last_seen, expires_at)
|
||||
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
|
||||
params![session_id, pub_key, salt.to_vec(), initial_hash, now, now, expires_at],
|
||||
)?;
|
||||
|
||||
let opcodes = vm::generate_random_program(8..=16);
|
||||
let opcodes_b64 = base64::Engine::encode(&base64::engine::general_purpose::STANDARD, &opcodes);
|
||||
|
||||
Ok(InitResponse {
|
||||
session_id,
|
||||
salt: hex::encode(salt),
|
||||
opcodes_b64,
|
||||
initial_hash: hex::encode(&initial_hash),
|
||||
expires_at,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn verify_heartbeat(
|
||||
conn: &rusqlite::Connection,
|
||||
req: &HeartbeatRequest,
|
||||
) -> Result<String, Box<dyn std::error::Error>> {
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT public_key, salt, last_hash, expires_at FROM sessions WHERE session_id = ?1",
|
||||
)?;
|
||||
let (pub_key, salt, stored_last_hash, expires_at): (Vec<u8>, Vec<u8>, Vec<u8>, u64) =
|
||||
stmt.query_row(params![req.session_id], |row| {
|
||||
Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?))
|
||||
})?;
|
||||
|
||||
let now = storage::current_time_ms();
|
||||
if now > expires_at {
|
||||
return Err("expired".into());
|
||||
}
|
||||
|
||||
// 1. Verify signature
|
||||
crypto::verify_signature(&pub_key, req)?;
|
||||
|
||||
// 2. Check chain continuity
|
||||
if stored_last_hash != hex::decode(&req.prev_hash)? {
|
||||
return Err("chain broken".into());
|
||||
}
|
||||
|
||||
// 3. Time window
|
||||
let diff = (now as i64) - (req.timestamp as i64);
|
||||
if diff.abs() > shared::constants::MAX_TIMESTAMP_DRIFT_MS {
|
||||
return Err("timestamp drift".into());
|
||||
}
|
||||
|
||||
// 4. Trusted mouse & fingerprint
|
||||
trust::validate_mouse(&req.entropy_data)?;
|
||||
fingerprint::validate(&req.fingerprint)?;
|
||||
|
||||
// 5. Compute new hash
|
||||
let prev_hash_bytes = hex::decode(&req.prev_hash)?;
|
||||
let new_hash = shared::hashing::next_chain_hash(
|
||||
&prev_hash_bytes,
|
||||
req.timestamp,
|
||||
&req.entropy_data,
|
||||
&req.stack_state,
|
||||
&salt,
|
||||
);
|
||||
|
||||
// 6. New salt for client
|
||||
let next_salt = rand::random::<[u8; shared::constants::SALT_LEN]>();
|
||||
let next_salt_hex = hex::encode(next_salt);
|
||||
|
||||
conn.execute(
|
||||
"UPDATE sessions SET last_hash=?1, salt=?2, chain_length=chain_length+1, last_seen=?3 WHERE session_id=?4",
|
||||
params![new_hash, next_salt.to_vec(), now, req.session_id],
|
||||
)?;
|
||||
|
||||
Ok(next_salt_hex)
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
use rusqlite::Connection;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
pub fn init_db() -> Result<Connection, rusqlite::Error> {
|
||||
let conn = Connection::open_in_memory()?;
|
||||
conn.execute_batch(
|
||||
"CREATE TABLE IF NOT EXISTS sessions (
|
||||
session_id TEXT PRIMARY KEY,
|
||||
public_key BLOB NOT NULL,
|
||||
salt BLOB NOT NULL,
|
||||
last_hash BLOB NOT NULL,
|
||||
chain_length INTEGER NOT NULL DEFAULT 1,
|
||||
created_at INTEGER NOT NULL,
|
||||
last_seen INTEGER NOT NULL,
|
||||
expires_at INTEGER NOT NULL
|
||||
);",
|
||||
)?;
|
||||
Ok(conn)
|
||||
}
|
||||
|
||||
pub fn current_time_ms() -> u64 {
|
||||
SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_millis() as u64
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
use shared::protocol::EntropyData;
|
||||
|
||||
pub fn validate_mouse(data: &EntropyData) -> Result<(), Box<dyn std::error::Error>> {
|
||||
let events = &data.events;
|
||||
if events.len() < 3 {
|
||||
return Err("few events".into());
|
||||
}
|
||||
let mut total_dist = 0.0;
|
||||
let mut pauses = 0u32;
|
||||
for i in 1..events.len() {
|
||||
let p = &events[i-1];
|
||||
let c = &events[i];
|
||||
let dx = c.x - p.x;
|
||||
let dy = c.y - p.y;
|
||||
let dt = (c.timestamp_ms - p.timestamp_ms).max(1.0);
|
||||
let dist = (dx*dx + dy*dy).sqrt();
|
||||
total_dist += dist;
|
||||
if dist < 0.2 && dt > 50.0 { pauses += 1; }
|
||||
}
|
||||
if total_dist < shared::constants::MIN_MOUSE_TOTAL_DIST {
|
||||
return Err("insufficient distance".into());
|
||||
}
|
||||
let avg_speed = total_dist / events.len() as f64;
|
||||
if avg_speed > shared::constants::MAX_MOUSE_AVG_SPEED {
|
||||
return Err("speed too high".into());
|
||||
}
|
||||
if pauses < shared::constants::MIN_PAUSE_COUNT {
|
||||
return Err("no pause".into());
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
use rand::Rng;
|
||||
|
||||
pub fn generate_random_program(len_range: std::ops::RangeInclusive<usize>) -> Vec<u8> {
|
||||
// Same logic as earlier, using shared::hashing for HASH if needed
|
||||
let mut rng = rand::thread_rng();
|
||||
let count = rng.gen_range(len_range);
|
||||
let mut ops = Vec::new();
|
||||
let mut depth: i32 = 0;
|
||||
for _ in 0..count {
|
||||
if depth < 2 {
|
||||
ops.push(0x00); // PUSH
|
||||
let val = rng.gen::<u32>();
|
||||
ops.extend_from_slice(&val.to_le_bytes());
|
||||
depth += 1;
|
||||
} else {
|
||||
let op = rng.gen_range(0..10);
|
||||
match op {
|
||||
0x00 => {
|
||||
ops.push(0x00);
|
||||
let val = rng.gen::<u32>();
|
||||
ops.extend_from_slice(&val.to_le_bytes());
|
||||
depth += 1;
|
||||
}
|
||||
0x01..=0x08 => { ops.push(op as u8); depth -= 1; }
|
||||
0x09 => { ops.push(0x09); depth = 1; }
|
||||
_ => unreachable!(),
|
||||
}
|
||||
}
|
||||
}
|
||||
ops
|
||||
}
|
||||
|
||||
// Server does not need to execute the program; client does.
|
||||
Reference in new issue
Block a user