diff --git a/neurosploit-rs/Cargo.lock b/neurosploit-rs/Cargo.lock index cd42bf8..8f2268d 100644 --- a/neurosploit-rs/Cargo.lock +++ b/neurosploit-rs/Cargo.lock @@ -91,6 +91,15 @@ version = "2.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b4388bee8683e3d04af747c73422af53102d2bd24d9eadb6cbc100baef4b43f8" +[[package]] +name = "block-buffer" +version = "0.10.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3078c7629b62d3f0439517fa394996acacc5cbc91c5a20d8c658e77abd503a71" +dependencies = [ + "generic-array", +] + [[package]] name = "bumpalo" version = "3.20.3" @@ -228,6 +237,15 @@ dependencies = [ "windows-sys 0.59.0", ] +[[package]] +name = "cpufeatures" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "59ed5838eebb26a2bb2e58f6d5b5316989ae9d08bab10e0e6d103e656d1b0280" +dependencies = [ + "libc", +] + [[package]] name = "crossterm" version = "0.28.1" @@ -253,6 +271,16 @@ dependencies = [ "winapi", ] +[[package]] +name = "crypto-common" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a" +dependencies = [ + "generic-array", + "typenum", +] + [[package]] name = "darling" version = "0.23.0" @@ -300,6 +328,17 @@ dependencies = [ "zeroize", ] +[[package]] +name = "digest" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" +dependencies = [ + "block-buffer", + "crypto-common", + "subtle", +] + [[package]] name = "displaydoc" version = "0.2.6" @@ -477,6 +516,16 @@ dependencies = [ "slab", ] +[[package]] +name = "generic-array" +version = "0.14.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a" +dependencies = [ + "typenum", + "version_check", +] + [[package]] name = "getrandom" version = "0.2.17" @@ -521,6 +570,15 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "hmac" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6c49c37c09c17a53d937dfbb742eb3a961d65a994e6bcdcf37e7399d0cc8ab5e" +dependencies = [ + "digest", +] + [[package]] name = "home" version = "0.5.12" @@ -891,11 +949,14 @@ name = "neurosploit-harness" version = "4.0.0" dependencies = [ "anyhow", + "base64", "futures", + "hmac", "regex", "reqwest", "serde", "serde_json", + "sha2", "tokio", "walkdir", ] @@ -1392,6 +1453,17 @@ dependencies = [ "serde", ] +[[package]] +name = "sha2" +version = "0.10.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + [[package]] name = "shell-words" version = "1.1.1" @@ -1720,6 +1792,12 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" +[[package]] +name = "typenum" +version = "1.20.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6f5e870be6c3b371b77fe0ee0bafb859fa4964b4404c27de1d380043c4dda20" + [[package]] name = "unicode-ident" version = "1.0.24" @@ -1785,6 +1863,12 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" +[[package]] +name = "version_check" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" + [[package]] name = "walkdir" version = "2.5.0" diff --git a/neurosploit-rs/crates/harness/Cargo.toml b/neurosploit-rs/crates/harness/Cargo.toml index dcc226b..1525018 100644 --- a/neurosploit-rs/crates/harness/Cargo.toml +++ b/neurosploit-rs/crates/harness/Cargo.toml @@ -16,4 +16,7 @@ reqwest.workspace = true anyhow.workspace = true futures.workspace = true walkdir = "2" +sha2 = "0.10" +hmac = "0.12" +base64 = "0.22" regex = "1" diff --git a/neurosploit-rs/crates/harness/src/audit.rs b/neurosploit-rs/crates/harness/src/audit.rs new file mode 100644 index 0000000..5afeb12 --- /dev/null +++ b/neurosploit-rs/crates/harness/src/audit.rs @@ -0,0 +1,576 @@ +//! Mandatory audit trail and hard kill conditions. +//! +//! An autonomous tool acting against someone else's systems has to be able to +//! answer, afterwards and in detail: *what did you do, to what, under whose +//! authority, and what happened?* Logs written for humans do not answer that — +//! they are prose, they are lossy, and they are trivially reordered. So every +//! action records one structured line: +//! +//! ```json +//! {"timestamp":"…","agent":"yaga","hypothesis":"H-023","action":"…", +//! "target":"…","policy_decision":"…","operator":null,"tool":"…", +//! "result":"…","evidence_hash":"…","capability_token":"…"} +//! ``` +//! +//! Two design decisions make it worth having: +//! +//! - **It is hash-chained.** Each record carries the hash of the one before it, +//! so a record cannot be removed or altered after the fact without breaking +//! every hash that follows. An append-only file that anyone can edit proves +//! nothing; [`AuditLog::verify`] is what turns it into evidence. +//! - **It records refusals too.** A trail containing only what happened cannot +//! show restraint. "The policy denied this, and the agent stopped" is exactly +//! what an operator needs to demonstrate afterwards. +//! +//! ## Hard kill conditions +//! +//! [`KillSwitch`] aborts the whole run, immediately and without negotiation, +//! when something happens that no amount of clever reasoning should be allowed +//! to continue through: the target stopped answering after our traffic, an +//! out-of-scope host was contacted, a forbidden industrial function code was +//! attempted, the capability token expired mid-run. These are *conditions*, not +//! heuristics — each one either happened or it did not, and each one ends the +//! engagement with an audit record explaining why. + +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use std::path::{Path, PathBuf}; +use std::sync::Mutex; + +/// One recorded action. The field set is fixed on purpose: an audit format that +/// varies per call site cannot be queried, and a record that omits the decision +/// or the authority does not answer the question the trail exists for. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct AuditRecord { + /// RFC3339 UTC. + pub timestamp: String, + /// Which agent acted. + pub agent: String, + /// The hypothesis this action was testing ("H-023"), so a trail can be read + /// as reasoning rather than as a list of requests. + pub hypothesis: String, + pub action: String, + pub target: String, + /// `allow` · `confirm: …` · `deny: …` — the policy's verdict, verbatim. + pub policy_decision: String, + /// Who approved, when approval was required. `None` means unattended, which + /// is a materially different claim from "someone approved it". + pub operator: Option, + pub tool: String, + pub result: String, + /// SHA-256 of the evidence this action produced. The evidence itself lives + /// with the run; the hash is what proves the two belong together. + pub evidence_hash: String, + /// Id of the grant that authorized it (never the token itself — the trail + /// is shared, and a token in it would be a credential leak). + pub capability_token: String, + /// Position in the chain, from 1. + #[serde(default)] + pub seq: u64, + /// Hash of the previous record; empty for the first. + #[serde(default)] + pub prev_hash: String, + /// Hash of this record's content plus `prev_hash`. + #[serde(default)] + pub hash: String, +} + +pub fn sha256_hex(data: &[u8]) -> String { + let mut h = Sha256::new(); + h.update(data); + format!("{:x}", h.finalize()) +} + +fn rfc3339_utc() -> String { + // The chain needs an ordered, unambiguous stamp and nothing more; pulling a + // date library in for that would be a dependency for cosmetics. + let secs = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0); + let days = secs / 86_400; + let tod = secs % 86_400; + let (y, m, d) = civil_from_days(days as i64); + format!("{y:04}-{m:02}-{d:02}T{:02}:{:02}:{:02}Z", tod / 3600, (tod % 3600) / 60, tod % 60) +} + +/// Howard Hinnant's days-from-civil, inverted. Exact for the whole proleptic +/// Gregorian range, which is more than a run log needs but costs nothing. +fn civil_from_days(z: i64) -> (i64, u32, u32) { + let z = z + 719_468; + let era = if z >= 0 { z } else { z - 146_096 } / 146_097; + let doe = (z - era * 146_097) as u64; + let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; + let y = yoe as i64 + era * 400; + let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); + let mp = (5 * doy + 2) / 153; + let d = (doy - (153 * mp + 2) / 5 + 1) as u32; + let m = if mp < 10 { mp + 3 } else { mp - 9 } as u32; + (if m <= 2 { y + 1 } else { y }, m, d) +} + +impl AuditRecord { + /// A record with the fixed fields filled and the chain fields empty — + /// [`AuditLog::append`] seals it. + pub fn new(agent: &str, action: &str, target: &str) -> AuditRecord { + AuditRecord { + timestamp: rfc3339_utc(), + agent: agent.to_string(), + hypothesis: String::new(), + action: action.to_string(), + target: target.to_string(), + policy_decision: String::new(), + operator: None, + tool: String::new(), + result: String::new(), + evidence_hash: String::new(), + capability_token: String::new(), + seq: 0, + prev_hash: String::new(), + hash: String::new(), + } + } + pub fn hypothesis(mut self, h: &str) -> Self { + self.hypothesis = h.to_string(); + self + } + pub fn decision(mut self, d: &str) -> Self { + self.policy_decision = d.to_string(); + self + } + pub fn operator(mut self, who: Option<&str>) -> Self { + self.operator = who.map(|s| s.to_string()); + self + } + pub fn tool(mut self, t: &str) -> Self { + self.tool = t.to_string(); + self + } + pub fn result(mut self, r: &str) -> Self { + self.result = r.to_string(); + self + } + pub fn evidence(mut self, bytes: &[u8]) -> Self { + self.evidence_hash = sha256_hex(bytes); + self + } + pub fn capability(mut self, id: &str) -> Self { + self.capability_token = id.to_string(); + self + } + + /// The bytes the chain hash covers: every meaningful field plus the + /// previous hash. `hash` itself is excluded, or it would be hashing itself. + fn digest_input(&self) -> String { + format!( + "{}|{}|{}|{}|{}|{}|{}|{}|{}|{}|{}|{}|{}", + self.seq, + self.prev_hash, + self.timestamp, + self.agent, + self.hypothesis, + self.action, + self.target, + self.policy_decision, + self.operator.clone().unwrap_or_else(|| "null".into()), + self.tool, + self.result, + self.evidence_hash, + self.capability_token + ) + } +} + +/// Append-only, hash-chained trail on disk (JSON Lines). +pub struct AuditLog { + path: PathBuf, + state: Mutex<(u64, String)>, // (last seq, last hash) +} + +impl AuditLog { + /// Open (or create) the trail, resuming the chain from what is already on + /// disk so a restarted run continues the same chain instead of forking it. + pub fn open(path: impl AsRef) -> AuditLog { + let path = path.as_ref().to_path_buf(); + if let Some(parent) = path.parent() { + let _ = std::fs::create_dir_all(parent); + } + let mut last = (0u64, String::new()); + if let Ok(text) = std::fs::read_to_string(&path) { + for line in text.lines().filter(|l| !l.trim().is_empty()) { + if let Ok(r) = serde_json::from_str::(line) { + last = (r.seq, r.hash); + } + } + } + AuditLog { path, state: Mutex::new(last) } + } + + /// Seal a record into the chain and write it. + pub fn append(&self, mut rec: AuditRecord) -> AuditRecord { + let mut guard = match self.state.lock() { + Ok(g) => g, + Err(p) => p.into_inner(), + }; + rec.seq = guard.0 + 1; + rec.prev_hash = guard.1.clone(); + rec.hash = sha256_hex(rec.digest_input().as_bytes()); + if let Ok(line) = serde_json::to_string(&rec) { + use std::io::Write; + if let Ok(mut f) = std::fs::OpenOptions::new().create(true).append(true).open(&self.path) { + let _ = writeln!(f, "{line}"); + } + } + *guard = (rec.seq, rec.hash.clone()); + rec + } + + pub fn read_all(&self) -> Vec { + std::fs::read_to_string(&self.path) + .map(|t| t.lines().filter_map(|l| serde_json::from_str(l).ok()).collect()) + .unwrap_or_default() + } + + /// Re-derive every hash. Returns the first sequence number that does not + /// match — a record that was altered or removed after it was written. + pub fn verify(&self) -> Result { + let records = self.read_all(); + let mut prev = String::new(); + for (i, r) in records.iter().enumerate() { + let expect_seq = i as u64 + 1; + if r.seq != expect_seq { + return Err(format!("record {} is out of sequence (claims #{}) — an entry was removed or reordered", i + 1, r.seq)); + } + if r.prev_hash != prev { + return Err(format!("record #{} does not follow the previous one — the chain was broken", r.seq)); + } + if r.hash != sha256_hex(r.digest_input().as_bytes()) { + return Err(format!("record #{} was altered after it was written", r.seq)); + } + prev = r.hash.clone(); + } + Ok(records.len()) + } + + pub fn path(&self) -> &Path { + &self.path + } +} + +/// Why a run was killed. Each variant is a condition that either occurred or +/// did not — none of them is a judgement call. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "kebab-case", tag = "condition", content = "detail")] +pub enum KillReason { + /// The target stopped answering after we started sending traffic. + TargetUnresponsive(String), + /// Error responses crossed the threshold — we are breaking it, not testing it. + ErrorRateExceeded(String), + /// A request left the authorized scope. + OutOfScope(String), + /// An industrial function code on the forbidden list was attempted. + ForbiddenFunctionCode(u16), + /// A safety instrumented system was addressed. + SafetySystemTouched(String), + /// The grant expired or was revoked while the run was in flight. + CapabilityInvalid(String), + /// The agent kept attempting things the policy refuses. + RepeatedPolicyViolations(usize), + /// Wall-clock or budget ceiling. + BudgetExhausted(String), + /// A human stopped it. + OperatorStop(String), +} + +impl KillReason { + pub fn explain(&self) -> String { + match self { + KillReason::TargetUnresponsive(d) => format!("target stopped responding after our traffic ({d}) — continuing risks an outage we caused"), + KillReason::ErrorRateExceeded(d) => format!("error rate above the threshold ({d}) — the target is failing, not revealing"), + KillReason::OutOfScope(h) => format!("a request was directed at {h}, outside the authorized scope"), + KillReason::ForbiddenFunctionCode(c) => format!("industrial function code {c} is on the forbidden list — it can stop a process"), + KillReason::SafetySystemTouched(h) => format!("a safety instrumented system was addressed ({h})"), + KillReason::CapabilityInvalid(d) => format!("the capability token is no longer valid ({d}) — authorization ended mid-run"), + KillReason::RepeatedPolicyViolations(n) => format!("{n} refused actions were attempted — the run is not respecting its policy"), + KillReason::BudgetExhausted(d) => format!("budget exhausted ({d})"), + KillReason::OperatorStop(w) => format!("stopped by {w}"), + } + } +} + +/// Conditions that end a run outright. +#[derive(Debug)] +pub struct KillSwitch { + /// Consecutive transport failures tolerated after the target has answered + /// at least once. + pub max_consecutive_failures: usize, + /// Fraction of 5xx responses that counts as breaking the target. + pub max_error_rate: f64, + /// Minimum responses before the error rate means anything — 2 failures out + /// of 2 requests is noise, not a trend. + pub error_rate_floor: usize, + /// Refused actions tolerated before the run is stopped. + pub max_policy_violations: usize, + state: Mutex, +} + +#[derive(Debug, Default)] +struct KillState { + responded_once: bool, + consecutive_failures: usize, + responses: usize, + errors: usize, + violations: usize, + killed: Option, +} + +impl Default for KillSwitch { + fn default() -> Self { + KillSwitch { + max_consecutive_failures: 5, + max_error_rate: 0.6, + error_rate_floor: 10, + max_policy_violations: 5, + state: Mutex::new(KillState::default()), + } + } +} + +impl KillSwitch { + /// OT runs are stopped far sooner: a device that misses a few requests may + /// already be struggling, and "let's see if it recovers" is not a decision + /// anyone should make on a live process. + pub fn ot() -> KillSwitch { + KillSwitch { + max_consecutive_failures: 2, + max_error_rate: 0.3, + error_rate_floor: 5, + max_policy_violations: 1, + state: Mutex::new(KillState::default()), + } + } + + fn lock(&self) -> std::sync::MutexGuard<'_, KillState> { + match self.state.lock() { + Ok(g) => g, + Err(p) => p.into_inner(), + } + } + + /// Record a response. Returns a kill reason the moment one is tripped. + pub fn note_response(&self, status: u16) -> Option { + let mut s = self.lock(); + s.responded_once = true; + s.consecutive_failures = 0; + s.responses += 1; + if status >= 500 { + s.errors += 1; + } + if s.responses >= self.error_rate_floor { + let rate = s.errors as f64 / s.responses as f64; + if rate > self.max_error_rate { + let r = KillReason::ErrorRateExceeded(format!("{:.0}% of {} responses were 5xx", rate * 100.0, s.responses)); + s.killed = Some(r.clone()); + return Some(r); + } + } + None + } + + /// Record a transport failure (connection refused, timeout). + pub fn note_failure(&self, what: &str) -> Option { + let mut s = self.lock(); + s.consecutive_failures += 1; + // Before the target ever answered, failures mean "wrong address" or + // "nothing listening" — not "we knocked it over". + if s.responded_once && s.consecutive_failures >= self.max_consecutive_failures { + let r = KillReason::TargetUnresponsive(format!("{} consecutive failures ({what})", s.consecutive_failures)); + s.killed = Some(r.clone()); + return Some(r); + } + None + } + + /// Record an action the policy refused. + pub fn note_violation(&self, _what: &str) -> Option { + let mut s = self.lock(); + s.violations += 1; + if s.violations >= self.max_policy_violations { + let r = KillReason::RepeatedPolicyViolations(s.violations); + s.killed = Some(r.clone()); + return Some(r); + } + None + } + + /// Conditions with no threshold: one occurrence ends the run. + pub fn trip(&self, reason: KillReason) -> KillReason { + let mut s = self.lock(); + s.killed = Some(reason.clone()); + reason + } + + pub fn killed(&self) -> Option { + self.lock().killed.clone() + } + + pub fn is_killed(&self) -> bool { + self.lock().killed.is_some() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn tmp(name: &str) -> PathBuf { + let mut p = std::env::temp_dir(); + p.push(format!("ns-audit-{}-{}.jsonl", name, std::process::id())); + let _ = std::fs::remove_file(&p); + p + } + + #[test] + fn a_record_carries_every_field_the_format_promises() { + let log = AuditLog::open(tmp("fields")); + let rec = log.append( + AuditRecord::new("yaga", "GET /admin", "https://app.test/admin") + .hypothesis("H-023") + .decision("allow") + .operator(None) + .tool("replay") + .result("200, 4kb") + .evidence(b"raw response bytes") + .capability("cap-001"), + ); + let json: serde_json::Value = serde_json::from_str(&serde_json::to_string(&rec).unwrap()).unwrap(); + for k in ["timestamp", "agent", "hypothesis", "action", "target", "policy_decision", "operator", "tool", "result", "evidence_hash", "capability_token"] { + assert!(json.get(k).is_some(), "{k} missing from the record"); + } + assert!(json["operator"].is_null(), "unattended must be null, not an empty string"); + assert_eq!(rec.evidence_hash, sha256_hex(b"raw response bytes")); + let _ = std::fs::remove_file(log.path()); + } + + #[test] + fn the_chain_verifies_and_notices_tampering() { + let path = tmp("chain"); + let log = AuditLog::open(&path); + for i in 0..5 { + log.append(AuditRecord::new("yaga", &format!("action {i}"), "https://app.test/").decision("allow")); + } + assert_eq!(log.verify(), Ok(5)); + + // Rewrite one line the way someone hiding an action would. + let text = std::fs::read_to_string(&path).unwrap(); + let mut lines: Vec = text.lines().map(String::from).collect(); + let mut rec: AuditRecord = serde_json::from_str(&lines[2]).unwrap(); + rec.target = "https://somewhere-else.test/".into(); + lines[2] = serde_json::to_string(&rec).unwrap(); + std::fs::write(&path, lines.join("\n") + "\n").unwrap(); + + let reopened = AuditLog::open(&path); + match reopened.verify() { + Err(e) => assert!(e.contains("altered"), "{e}"), + Ok(_) => panic!("an edited record must break the chain"), + } + let _ = std::fs::remove_file(&path); + } + + #[test] + fn removing_a_record_breaks_the_chain_too() { + let path = tmp("removed"); + let log = AuditLog::open(&path); + for i in 0..4 { + log.append(AuditRecord::new("yaga", &format!("a{i}"), "t")); + } + let text = std::fs::read_to_string(&path).unwrap(); + let kept: Vec<&str> = text.lines().enumerate().filter(|(i, _)| *i != 1).map(|(_, l)| l).collect(); + std::fs::write(&path, kept.join("\n") + "\n").unwrap(); + assert!(AuditLog::open(&path).verify().is_err(), "a deleted entry must not go unnoticed"); + let _ = std::fs::remove_file(&path); + } + + #[test] + fn reopening_continues_the_same_chain() { + let path = tmp("resume"); + { + let log = AuditLog::open(&path); + log.append(AuditRecord::new("a", "one", "t")); + log.append(AuditRecord::new("a", "two", "t")); + } + let log = AuditLog::open(&path); + let third = log.append(AuditRecord::new("a", "three", "t")); + assert_eq!(third.seq, 3, "a restart must not fork the chain"); + assert_eq!(log.verify(), Ok(3)); + let _ = std::fs::remove_file(&path); + } + + #[test] + fn refusals_are_recorded_not_just_actions() { + let log = AuditLog::open(tmp("deny")); + let rec = log.append(AuditRecord::new("yaga", "DELETE /orders/1", "https://app.test/orders/1").decision("deny: destructive actions are disabled")); + assert!(rec.policy_decision.starts_with("deny"), "a trail that omits refusals cannot show restraint"); + let _ = std::fs::remove_file(log.path()); + } + + #[test] + fn failures_before_the_first_response_are_not_an_outage_we_caused() { + let k = KillSwitch::default(); + for _ in 0..10 { + assert!(k.note_failure("connection refused").is_none(), "nothing listening is not the same as knocked over"); + } + assert!(!k.is_killed()); + } + + #[test] + fn losing_a_target_after_it_answered_kills_the_run() { + let k = KillSwitch::default(); + assert!(k.note_response(200).is_none()); + let mut tripped = None; + for _ in 0..k.max_consecutive_failures { + tripped = k.note_failure("timeout"); + } + match tripped { + Some(KillReason::TargetUnresponsive(_)) => {} + other => panic!("expected a kill, got {other:?}"), + } + assert!(k.is_killed()); + } + + #[test] + fn a_few_errors_are_not_a_trend() { + let k = KillSwitch::default(); + assert!(k.note_response(500).is_none()); + assert!(k.note_response(500).is_none(), "2 of 2 is noise, not a trend"); + for _ in 0..8 { + k.note_response(500); + } + assert!(k.is_killed(), "a sustained 5xx rate means we are breaking it"); + } + + #[test] + fn ot_stops_far_sooner_than_a_web_run() { + let k = KillSwitch::ot(); + k.note_response(200); + assert!(k.note_failure("timeout").is_none()); + assert!(k.note_failure("timeout").is_some(), "a PLC missing two requests already warrants stopping"); + } + + #[test] + fn conditions_without_a_threshold_trip_on_the_first_occurrence() { + let k = KillSwitch::default(); + let r = k.trip(KillReason::SafetySystemTouched("sis-01".into())); + assert!(r.explain().contains("safety instrumented")); + assert!(k.is_killed()); + } + + #[test] + fn the_timestamp_is_a_sane_rfc3339_date() { + let ts = rfc3339_utc(); + assert!(ts.ends_with('Z') && ts.len() == 20, "{ts}"); + let year: i64 = ts[..4].parse().unwrap(); + assert!((2020..2100).contains(&year), "{ts}"); + assert_eq!(civil_from_days(0), (1970, 1, 1)); + } +} diff --git a/neurosploit-rs/crates/harness/src/capability.rs b/neurosploit-rs/crates/harness/src/capability.rs new file mode 100644 index 0000000..556cd2b --- /dev/null +++ b/neurosploit-rs/crates/harness/src/capability.rs @@ -0,0 +1,378 @@ +//! Capability tokens — authorization the harness can verify, not just trust. +//! +//! Everything in [`crate::scope`] and [`crate::policy`] is configuration: the +//! operator types a scope, the harness enforces it. That is right for a local +//! run and not enough for a real engagement, where the person running the tool +//! and the person who authorized the test are different people, and the client +//! wants the boundary to be theirs rather than whatever was typed at the +//! keyboard. +//! +//! A capability token is that grant, made checkable: an HMAC-signed statement +//! of *who authorized what, against which hosts, in which environment, until +//! when*. The harness verifies the signature with a shared secret, refuses an +//! expired or not-yet-valid token, and then treats the token as the ceiling — +//! local configuration may narrow it and can never widen it. +//! +//! ```text +//! ns-cap.v1.. +//! ``` +//! +//! Two properties worth stating plainly: +//! +//! - **The token narrows, never widens.** [`CapabilityToken::constrain`] takes +//! the intersection with the local policy. A token that says `*.example.com` +//! does not authorize a run configured for one host to wander to the others. +//! - **It is not a secret store.** The signature proves the grant was issued by +//! someone holding the key; it does not encrypt anything. Never put +//! credentials in a token — its payload is readable by anyone holding it. + +use crate::policy::{ActionKind, Environment}; +use crate::scope::{Pattern, ScopePolicy}; +use base64::Engine; +use hmac::{Hmac, Mac}; +use serde::{Deserialize, Serialize}; +use sha2::Sha256; + +type HmacSha256 = Hmac; + +const PREFIX: &str = "ns-cap.v1."; + +/// The signed statement of what is authorized. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct Capability { + /// Stable id, quoted in every audit record so an action can be traced back + /// to the grant that permitted it. + pub id: String, + /// Who authorized the engagement (the client contact, ticket, or system). + pub issuer: String, + /// Who it was issued to (the tester or team). + pub subject: String, + /// Hosts / wildcards / CIDRs / URL prefixes this grant covers. + #[serde(default)] + pub scope: Vec, + /// Explicit exclusions inside that scope. + #[serde(default)] + pub exclude: Vec, + pub environment: Environment, + /// The strongest action this grant permits. + pub max_action: ActionKind, + /// Ceiling on `effective_risk` (see [`crate::policy::Risk`]). + pub max_risk: f64, + /// Unix seconds. A grant without an end is not a grant, it is a standing + /// permission nobody remembers issuing. + pub expires_at: u64, + #[serde(default)] + pub not_before: u64, + /// Free-text reference (engagement id, statement of work, ticket). + #[serde(default)] + pub reference: String, +} + +#[derive(Debug, Clone, PartialEq)] +pub enum TokenError { + Malformed(String), + BadSignature, + Expired { expired_at: u64, now: u64 }, + NotYetValid { starts_at: u64, now: u64 }, +} + +impl std::fmt::Display for TokenError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + TokenError::Malformed(w) => write!(f, "capability token is malformed: {w}"), + TokenError::BadSignature => write!(f, "capability token signature does not verify — it was not issued by the holder of this key, or it was altered"), + TokenError::Expired { expired_at, now } => write!(f, "capability token expired {}s ago", now.saturating_sub(*expired_at)), + TokenError::NotYetValid { starts_at, now } => write!(f, "capability token is not valid for another {}s", starts_at.saturating_sub(*now)), + } + } +} + +fn now() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0) +} + +fn b64() -> base64::engine::general_purpose::GeneralPurpose { + base64::engine::general_purpose::URL_SAFE_NO_PAD +} + +impl Capability { + /// Sign this grant with the issuer's key. + pub fn issue(&self, key: &[u8]) -> String { + let payload = serde_json::to_vec(self).unwrap_or_default(); + let body = b64().encode(&payload); + let sig = sign(key, body.as_bytes()); + format!("{PREFIX}{body}.{}", b64().encode(sig)) + } + + /// Verify a token and return the grant it carries. + /// + /// Signature first, then validity window: a well-formed token from the + /// wrong key must never get as far as having its claims read, or an + /// attacker-supplied "grant" decides what the harness believes. + pub fn verify(token: &str, key: &[u8]) -> Result { + let rest = token.trim().strip_prefix(PREFIX).ok_or_else(|| TokenError::Malformed("wrong prefix or version".into()))?; + let (body, sig_b64) = rest.rsplit_once('.').ok_or_else(|| TokenError::Malformed("missing signature".into()))?; + let expected = sign(key, body.as_bytes()); + let given = b64().decode(sig_b64).map_err(|_| TokenError::Malformed("signature is not base64url".into()))?; + if !constant_time_eq(&expected, &given) { + return Err(TokenError::BadSignature); + } + let payload = b64().decode(body).map_err(|_| TokenError::Malformed("payload is not base64url".into()))?; + let cap: Capability = serde_json::from_slice(&payload).map_err(|e| TokenError::Malformed(e.to_string()))?; + let t = now(); + if cap.not_before > 0 && t < cap.not_before { + return Err(TokenError::NotYetValid { starts_at: cap.not_before, now: t }); + } + if cap.expires_at > 0 && t > cap.expires_at { + return Err(TokenError::Expired { expired_at: cap.expires_at, now: t }); + } + Ok(cap) + } + + /// Read a token's claims WITHOUT verifying them — for displaying an + /// unverified token in a UI. Never use this to make a decision. + pub fn peek(token: &str) -> Option { + let rest = token.trim().strip_prefix(PREFIX)?; + let (body, _) = rest.rsplit_once('.')?; + serde_json::from_slice(&b64().decode(body).ok()?).ok() + } + + /// Intersect a local policy with this grant. + /// + /// The token is the ceiling: entries the operator configured that the grant + /// does not cover are dropped (and reported), and the grant's exclusions are + /// added. Configuration can only ever be narrower than the authorization it + /// runs under. + pub fn constrain(&self, local: &ScopePolicy) -> (ScopePolicy, Vec) { + let mut granted = ScopePolicy::default(); + for s in &self.scope { + granted.allow(s); + } + for e in &self.exclude { + granted.deny(e); + } + granted.soft = local.soft.clone(); + + let mut dropped = Vec::new(); + let mut kept: Vec = Vec::new(); + for p in &local.hard { + // A local entry survives only if the grant covers it. Host-shaped + // patterns are checked as a URL; anything the grant does not match + // is refused rather than silently kept. + let probe = format!("https://{}", p.as_text().trim_start_matches("*.")); + if granted.in_hard_scope(&probe) { + kept.push(p.clone()); + } else { + dropped.push(p.as_text()); + } + } + let mut out = if kept.is_empty() { + // The operator configured nothing the grant covers (or configured + // nothing at all) — fall back to exactly what was granted. + granted.clone() + } else { + let mut o = ScopePolicy::default(); + o.hard = kept; + o.soft = local.soft.clone(); + o + }; + for e in &self.exclude { + out.deny(e); + } + for e in &local.exclude { + out.deny(&e.as_text()); + } + (out, dropped) + } + + /// Is this action kind within the grant? + pub fn permits(&self, action: ActionKind) -> bool { + rank(action) <= rank(self.max_action) + } + + pub fn summary(&self) -> String { + let left = self.expires_at.saturating_sub(now()); + format!( + "{} · issued by {} to {} · {} · max {} · risk ≤ {:.1} · {}", + self.id, + self.issuer, + self.subject, + self.scope.join(","), + self.max_action.as_str(), + self.max_risk, + if self.expires_at == 0 { + "no expiry".to_string() + } else if left == 0 { + "EXPIRED".to_string() + } else { + format!("{}h left", left / 3600) + } + ) + } +} + +fn rank(a: ActionKind) -> u8 { + match a { + ActionKind::Read => 0, + ActionKind::Enumerate => 1, + ActionKind::Authenticate => 2, + ActionKind::ProbeExploit => 3, + ActionKind::Write => 4, + ActionKind::Disruptive => 5, + } +} + +fn sign(key: &[u8], body: &[u8]) -> Vec { + let mut mac = HmacSha256::new_from_slice(key).expect("HMAC accepts any key length"); + mac.update(body); + mac.finalize().into_bytes().to_vec() +} + +/// Compare without an early return, so a forged signature cannot be recovered +/// byte by byte from how long the check took. +fn constant_time_eq(a: &[u8], b: &[u8]) -> bool { + if a.len() != b.len() { + return false; + } + let mut diff = 0u8; + for (x, y) in a.iter().zip(b.iter()) { + diff |= x ^ y; + } + diff == 0 +} + +/// The signing key, from `NEUROSPLOIT_CAPABILITY_KEY` or a file path in +/// `NEUROSPLOIT_CAPABILITY_KEY_FILE`. Returns `None` when no key is +/// configured — in which case tokens cannot be verified and the harness must +/// say so rather than accepting them. +pub fn key_from_env() -> Option> { + if let Ok(k) = std::env::var("NEUROSPLOIT_CAPABILITY_KEY") { + if !k.trim().is_empty() { + return Some(k.into_bytes()); + } + } + if let Ok(p) = std::env::var("NEUROSPLOIT_CAPABILITY_KEY_FILE") { + if let Ok(bytes) = std::fs::read(p.trim()) { + if !bytes.is_empty() { + return Some(bytes); + } + } + } + None +} + +#[cfg(test)] +mod tests { + use super::*; + + fn cap() -> Capability { + Capability { + id: "cap-001".into(), + issuer: "acme-security@example.com".into(), + subject: "red-team".into(), + scope: vec!["*.example.com".into()], + exclude: vec!["payments.example.com".into()], + environment: Environment::Production, + max_action: ActionKind::ProbeExploit, + max_risk: 3.5, + expires_at: now() + 3600, + not_before: 0, + reference: "SOW-2026-14".into(), + } + } + + #[test] + fn a_token_round_trips_and_verifies() { + let t = cap().issue(b"secret-key"); + let back = Capability::verify(&t, b"secret-key").expect("must verify"); + assert_eq!(back, cap()); + } + + #[test] + fn a_token_signed_with_another_key_is_refused() { + let t = cap().issue(b"secret-key"); + assert_eq!(Capability::verify(&t, b"different-key"), Err(TokenError::BadSignature)); + } + + #[test] + fn altering_the_claims_breaks_the_signature() { + let t = cap().issue(b"secret-key"); + // Widen the scope by hand, the way an over-eager tester might. + let mut forged = cap(); + forged.scope = vec!["*".into()]; + let payload = b64().encode(serde_json::to_vec(&forged).unwrap()); + let sig = t.rsplit_once('.').unwrap().1; + let tampered = format!("{PREFIX}{payload}.{sig}"); + assert_eq!(Capability::verify(&tampered, b"secret-key"), Err(TokenError::BadSignature)); + } + + #[test] + fn an_expired_token_is_refused_even_with_a_valid_signature() { + let mut c = cap(); + c.expires_at = now() - 10; + let t = c.issue(b"secret-key"); + match Capability::verify(&t, b"secret-key") { + Err(TokenError::Expired { .. }) => {} + other => panic!("an expired grant is not a grant: {other:?}"), + } + } + + #[test] + fn a_token_that_has_not_started_is_refused() { + let mut c = cap(); + c.not_before = now() + 600; + let t = c.issue(b"secret-key"); + assert!(matches!(Capability::verify(&t, b"secret-key"), Err(TokenError::NotYetValid { .. }))); + } + + #[test] + fn the_grant_narrows_local_scope_and_never_widens_it() { + let c = cap(); + let mut local = ScopePolicy::default(); + local.allow("app.example.com"); + local.allow("unrelated.test"); // not covered by the grant + let (effective, dropped) = c.constrain(&local); + assert!(effective.in_hard_scope("https://app.example.com/x")); + assert!(!effective.in_hard_scope("https://unrelated.test/"), "a host outside the grant must not survive"); + assert_eq!(dropped, vec!["unrelated.test"]); + // The grant's own exclusion still applies to everything. + assert!(!effective.in_hard_scope("https://payments.example.com/")); + } + + #[test] + fn an_empty_local_policy_inherits_exactly_what_was_granted() { + let (effective, dropped) = cap().constrain(&ScopePolicy::default()); + assert!(dropped.is_empty()); + assert!(effective.in_hard_scope("https://api.example.com/")); + assert!(!effective.in_hard_scope("https://payments.example.com/")); + } + + #[test] + fn the_grant_caps_which_actions_are_permitted() { + let c = cap(); // max_action = ProbeExploit + assert!(c.permits(ActionKind::Read)); + assert!(c.permits(ActionKind::ProbeExploit)); + assert!(!c.permits(ActionKind::Write)); + assert!(!c.permits(ActionKind::Disruptive)); + } + + #[test] + fn peek_reads_claims_but_is_not_a_decision() { + let t = cap().issue(b"secret-key"); + assert_eq!(Capability::peek(&t).unwrap().id, "cap-001"); + // Even garbage signatures peek fine — which is exactly why peek must + // never gate anything. + let tampered = format!("{}.{}", t.rsplit_once('.').unwrap().0, b64().encode(b"nonsense")); + assert!(Capability::peek(&tampered).is_some()); + assert!(Capability::verify(&tampered, b"secret-key").is_err()); + } + + #[test] + fn malformed_input_is_reported_as_malformed_not_as_a_bad_signature() { + assert!(matches!(Capability::verify("not-a-token", b"k"), Err(TokenError::Malformed(_)))); + assert!(matches!(Capability::verify("ns-cap.v1.nosignature", b"k"), Err(TokenError::Malformed(_)))); + } +} diff --git a/neurosploit-rs/crates/harness/src/lib.rs b/neurosploit-rs/crates/harness/src/lib.rs index 3d02700..ea9d9f0 100644 --- a/neurosploit-rs/crates/harness/src/lib.rs +++ b/neurosploit-rs/crates/harness/src/lib.rs @@ -8,13 +8,16 @@ pub mod agents; pub mod attack_graph; +pub mod audit; pub mod belief; +pub mod capability; pub mod creds; pub mod grounding; pub mod hygiene; pub mod integrations; pub mod knowledge_graph; pub mod memory; +pub mod policy; pub mod pomdp; pub mod models; pub mod pipeline; @@ -37,6 +40,9 @@ pub use pipeline::run; pub use knowledge_graph::{EdgeKind, KnowledgeGraph, NodeKind}; pub use memory::{Memory, Query as MemoryQuery, Tier as MemoryTier}; pub use pool::{ModelPool, Task}; +pub use audit::{AuditLog, AuditRecord, KillReason, KillSwitch}; +pub use capability::{Capability, TokenError}; +pub use policy::{Act, ActionKind, BlastRadius, EngagementPolicy, Environment, Protocol, Risk, RiskDecision, SafetyPolicy}; pub use replay::{ReplayEngine, ReqSpec}; pub use scope::{Action as ScopeAction, Decision as ScopeDecision, ScopePolicy}; pub use types::{Finding, RunConfig}; diff --git a/neurosploit-rs/crates/harness/src/policy.rs b/neurosploit-rs/crates/harness/src/policy.rs new file mode 100644 index 0000000..7e5305e --- /dev/null +++ b/neurosploit-rs/crates/harness/src/policy.rs @@ -0,0 +1,826 @@ +//! Risk model and engagement policies. +//! +//! [`crate::scope`] answers *may we touch this address*. That is necessary and +//! nowhere near sufficient: an in-scope action can still be the wrong action. +//! Reading a web page and writing a holding register on a PLC are both "in +//! scope" against an authorized host, and only one of them can stop a physical +//! process. +//! +//! So risk is computed per action, from factors the operator declares: +//! +//! ```text +//! effective_risk = (action_risk +//! + asset_criticality +//! + protocol_risk +//! + privilege_level +//! + blast_radius) × environment_multiplier +//! ``` +//! +//! Each of the five terms is 0.0–1.0, so the sum is 0–5, and the multiplier +//! places that sum in its context: the same write is a different act on a lab +//! bench and on a live substation. The result is a number the [`SafetyPolicy`] +//! can refuse, and — more usefully — a number the operator can *read*, because +//! every term is named and comes from a declaration rather than a model's +//! judgement. +//! +//! Three policies sit on top, each answering a different question: +//! +//! - [`SafetyPolicy`] — *what may be done to the target?* Ceilings, approval +//! thresholds, and the hard prohibitions that make OT/ICS testing survivable. +//! - [`ReasoningPolicy`] — *how must the agent think?* Baseline before payload, +//! a bounded number of hypotheses, evidence before escalation, and explicit +//! stop conditions so a run ends on purpose rather than on exhaustion. +//! - [`ProofOfImpactPolicy`] — *what may be claimed?* The evidence a severity +//! has to carry before it is allowed to be that severity. +//! +//! ## Why OT/ICS is not "web testing with different ports" +//! +//! Industrial protocols were designed without authentication, on the assumption +//! of a physically isolated network. A Modbus write function is not an exploit +//! — it is the protocol working as intended, addressed to a device that may be +//! holding a valve. Scanners routinely crash PLCs simply by sending unexpected +//! data at line rate. So [`SafetyProfile::ot_conservative`] blocks writes and +//! fuzzing outright, caps the request rate to something the device tolerates, +//! and refuses the function codes that stop a CPU — and those refusals are not +//! advisory text in a prompt, they are checks. + +use serde::{Deserialize, Serialize}; + +fn clamp01(v: f64) -> f64 { + v.clamp(0.0, 1.0) +} + +/// What the action itself does, independent of where. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "kebab-case")] +pub enum ActionKind { + /// Read something already exposed. + Read, + /// Enumerate, fingerprint, map. + Enumerate, + /// Send a payload intended to prove a weakness, without changing state. + ProbeExploit, + /// Authenticate, create a session, use a credential. + Authenticate, + /// Write, update, or create application state. + Write, + /// Delete state, restart a service, change a device's operating mode. + Disruptive, +} + +impl ActionKind { + pub fn risk(self) -> f64 { + match self { + ActionKind::Read => 0.05, + ActionKind::Enumerate => 0.15, + ActionKind::ProbeExploit => 0.45, + ActionKind::Authenticate => 0.35, + ActionKind::Write => 0.8, + ActionKind::Disruptive => 1.0, + } + } + pub fn as_str(self) -> &'static str { + match self { + ActionKind::Read => "read", + ActionKind::Enumerate => "enumerate", + ActionKind::ProbeExploit => "probe-exploit", + ActionKind::Authenticate => "authenticate", + ActionKind::Write => "write", + ActionKind::Disruptive => "disruptive", + } + } +} + +/// The protocol carrying the action. Industrial protocols rank high not because +/// they are hard to speak but because they are trivial to speak — they have no +/// authentication and the devices behind them are fragile. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "kebab-case")] +pub enum Protocol { + Http, + Https, + Dns, + Smb, + Ssh, + Rdp, + Database, + /// Modbus, DNP3, S7comm, EtherNet/IP, BACnet, OPC-UA… + Industrial, + /// Safety instrumented systems — the layer that exists to stop the process + /// safely. Touching it is never routine. + SafetySystem, + Other, +} + +impl Protocol { + pub fn risk(self) -> f64 { + match self { + Protocol::Https => 0.1, + Protocol::Http => 0.15, + Protocol::Dns => 0.15, + Protocol::Database => 0.5, + Protocol::Smb => 0.45, + Protocol::Ssh => 0.4, + Protocol::Rdp => 0.45, + Protocol::Industrial => 0.9, + Protocol::SafetySystem => 1.0, + Protocol::Other => 0.3, + } + } + pub fn is_ot(self) -> bool { + matches!(self, Protocol::Industrial | Protocol::SafetySystem) + } + pub fn as_str(self) -> &'static str { + match self { + Protocol::Http => "http", + Protocol::Https => "https", + Protocol::Dns => "dns", + Protocol::Smb => "smb", + Protocol::Ssh => "ssh", + Protocol::Rdp => "rdp", + Protocol::Database => "database", + Protocol::Industrial => "industrial", + Protocol::SafetySystem => "safety-system", + Protocol::Other => "other", + } + } + /// Best-effort classification from a port, used when the operator did not + /// declare one. Industrial ports are recognised so an unlabelled 502/20000 + /// does not get treated as an ordinary service. + pub fn from_port(port: u16) -> Protocol { + match port { + 80 | 8080 | 8000 => Protocol::Http, + 443 | 8443 => Protocol::Https, + 53 => Protocol::Dns, + 22 => Protocol::Ssh, + 3389 => Protocol::Rdp, + 139 | 445 => Protocol::Smb, + 1433 | 3306 | 5432 | 1521 | 27017 => Protocol::Database, + // Modbus, DNP3, EtherNet/IP, S7comm(-plus), BACnet, OPC-UA, FL-net. + 502 | 802 | 20000 | 44818 | 2222 | 102 | 47808 | 4840 | 55000 => Protocol::Industrial, + _ => Protocol::Other, + } + } +} + +/// How far the consequences reach if the action goes wrong. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "kebab-case")] +pub enum BlastRadius { + /// One request, one object, trivially reversible. + SingleObject, + /// One account or one session. + SingleAccount, + /// One host or service instance. + SingleHost, + /// A shared service many consumers depend on. + SharedService, + /// A network segment, a cell, a production line. + Segment, + /// A physical process, or an organisation-wide system. + PhysicalProcess, +} + +impl BlastRadius { + pub fn risk(self) -> f64 { + match self { + BlastRadius::SingleObject => 0.05, + BlastRadius::SingleAccount => 0.2, + BlastRadius::SingleHost => 0.4, + BlastRadius::SharedService => 0.65, + BlastRadius::Segment => 0.85, + BlastRadius::PhysicalProcess => 1.0, + } + } +} + +/// Where this is running. The multiplier is the honest place for "same action, +/// different consequence" — it scales the whole sum rather than hiding inside +/// one term. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "kebab-case")] +pub enum Environment { + Lab, + Development, + Staging, + Production, + /// Production that runs a physical process. + OtProduction, +} + +impl Environment { + pub fn multiplier(self) -> f64 { + match self { + Environment::Lab => 0.3, + Environment::Development => 0.5, + Environment::Staging => 0.7, + Environment::Production => 1.0, + Environment::OtProduction => 1.6, + } + } + pub fn as_str(self) -> &'static str { + match self { + Environment::Lab => "lab", + Environment::Development => "development", + Environment::Staging => "staging", + Environment::Production => "production", + Environment::OtProduction => "ot-production", + } + } + pub fn parse(s: &str) -> Option { + Some(match s.trim().to_lowercase().as_str() { + "lab" => Environment::Lab, + "dev" | "development" => Environment::Development, + "stg" | "staging" => Environment::Staging, + "prod" | "production" => Environment::Production, + "ot" | "ot-prod" | "ot-production" | "ics" | "scada" => Environment::OtProduction, + _ => return None, + }) + } +} + +/// One action, described in the terms the risk formula needs. +#[derive(Debug, Clone)] +pub struct Act { + pub what: ActionKind, + pub protocol: Protocol, + pub blast: BlastRadius, + /// How privileged the identity performing it is, 0 (anonymous) to 1 (domain + /// admin / engineering workstation). + pub privilege: f64, + /// Declared importance of the asset, 0 (scratch) to 1 (crown jewels). + pub asset_criticality: f64, + /// Human-readable target, for the audit line. + pub target: String, +} + +impl Default for Act { + fn default() -> Self { + Act { + what: ActionKind::Read, + protocol: Protocol::Https, + blast: BlastRadius::SingleObject, + privilege: 0.2, + asset_criticality: 0.5, + target: String::new(), + } + } +} + +/// The computed risk, with every term kept so the number can be explained. +/// A score whose derivation is invisible gets argued with instead of acted on. +#[derive(Debug, Clone, PartialEq)] +pub struct Risk { + pub action_risk: f64, + pub asset_criticality: f64, + pub protocol_risk: f64, + pub privilege_level: f64, + pub blast_radius: f64, + pub environment_multiplier: f64, + pub effective: f64, +} + +impl Risk { + /// `(action + asset + protocol + privilege + blast) × environment`. + pub fn compute(act: &Act, env: Environment) -> Risk { + let action_risk = act.what.risk(); + let asset_criticality = clamp01(act.asset_criticality); + let protocol_risk = act.protocol.risk(); + let privilege_level = clamp01(act.privilege); + let blast_radius = act.blast.risk(); + let environment_multiplier = env.multiplier(); + let effective = + (action_risk + asset_criticality + protocol_risk + privilege_level + blast_radius) * environment_multiplier; + Risk { + action_risk, + asset_criticality, + protocol_risk, + privilege_level, + blast_radius, + environment_multiplier, + effective, + } + } + + /// The derivation, one line, for the audit log and the operator. + pub fn explain(&self) -> String { + format!( + "effective_risk {:.2} = (action {:.2} + asset {:.2} + protocol {:.2} + privilege {:.2} + blast {:.2}) × env {:.2}", + self.effective, + self.action_risk, + self.asset_criticality, + self.protocol_risk, + self.privilege_level, + self.blast_radius, + self.environment_multiplier + ) + } +} + +/// What the harness is allowed to do to the target. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SafetyPolicy { + pub environment: Environment, + /// Refuse any action above this effective risk. + pub max_effective_risk: f64, + /// Above this, act only with a human's explicit approval. + pub confirm_above: f64, + /// Writes of any kind. + pub allow_write: bool, + /// Actions that can stop or restart something. + pub allow_disruptive: bool, + /// Malformed/random input. On OT this is the classic way to crash a PLC + /// that has done nothing wrong. + pub allow_fuzzing: bool, + /// Ceiling on request rate, per protocol family. Industrial devices answer + /// one request at a time and fall over when treated like a web server. + pub max_requests_per_minute: u32, + /// Protocols that may not be touched at all. + pub forbidden_protocols: Vec, + /// Industrial function codes that must never be sent (Modbus write/ + /// diagnostic codes, S7 stop, DNP3 cold restart…). + pub forbidden_function_codes: Vec, + /// Free-text conditions for the operator's own record (maintenance window, + /// contact on call). Not enforceable, and kept separate so nobody mistakes + /// prose for a control. + pub notes: Vec, +} + +impl Default for SafetyPolicy { + fn default() -> Self { + SafetyPolicy::web_standard(Environment::Production) + } +} + +/// Named starting points an operator can pick and then adjust. +pub struct SafetyProfile; + +impl SafetyProfile { + /// Ordinary web application testing. + pub fn web_standard(env: Environment) -> SafetyPolicy { + SafetyPolicy::web_standard(env) + } + /// OT/ICS/SCADA: read-only, slow, and with the dangerous primitives removed + /// rather than discouraged. + pub fn ot_conservative() -> SafetyPolicy { + SafetyPolicy::ot_conservative() + } +} + +impl SafetyPolicy { + pub fn web_standard(environment: Environment) -> SafetyPolicy { + SafetyPolicy { + environment, + max_effective_risk: 3.5, + confirm_above: 2.5, + allow_write: environment != Environment::Production, + allow_disruptive: false, + allow_fuzzing: environment == Environment::Lab || environment == Environment::Development, + max_requests_per_minute: 240, + forbidden_protocols: vec![Protocol::SafetySystem], + forbidden_function_codes: Vec::new(), + notes: Vec::new(), + } + } + + /// The profile for live industrial environments. + /// + /// Everything here is a refusal rather than a warning, because the failure + /// mode is physical: a stopped CPU, a tripped line, a safety system that + /// was mid-test when it was needed. Discovery is still possible — passive + /// reads and enumeration — and that is usually where the findings are + /// anyway, since these protocols authenticate nothing. + pub fn ot_conservative() -> SafetyPolicy { + SafetyPolicy { + environment: Environment::OtProduction, + // Calibration matters here, and the obvious setting is wrong: a + // plain READ of a critical PLC scores 3.6 on this formula, so a + // tight ceiling refuses exactly the observation that OT findings + // come from. In an industrial environment it is the KIND of action + // that is forbidden (see the ProbeExploit/Write rules below), not + // the arithmetic. The ceiling catches the extremes; the low + // approval threshold means anything past trivial observation is a + // human's decision. + max_effective_risk: 5.0, + confirm_above: 3.0, + allow_write: false, + allow_disruptive: false, + allow_fuzzing: false, + // Roughly one request per second: what a small PLC tolerates while + // still serving its real traffic. + max_requests_per_minute: 60, + forbidden_protocols: vec![Protocol::SafetySystem], + forbidden_function_codes: vec![ + 5, // Modbus: write single coil + 6, // Modbus: write single register + 8, // Modbus: diagnostics (includes "restart communications") + 15, // Modbus: write multiple coils + 16, // Modbus: write multiple registers + 22, // Modbus: mask write register + 23, // Modbus: read/write multiple registers + 43, // Modbus: encapsulated interface transport + 0x29, // S7comm: PLC stop + 0x28, // S7comm: PLC start/warm restart + 13, // DNP3: cold restart + 14, // DNP3: warm restart + 18, // DNP3: stop application + ], + notes: vec![ + "OT profile: observation only. Any write, restart, or mode change requires a human, a maintenance window, and the process owner present.".into(), + ], + } + } + + /// The decision for one action. + pub fn check(&self, act: &Act) -> RiskDecision { + let risk = Risk::compute(act, self.environment); + + if self.forbidden_protocols.contains(&act.protocol) { + return RiskDecision::deny(risk, format!("{} is a forbidden protocol for this engagement", act.protocol.as_str())); + } + if act.what == ActionKind::Write && !self.allow_write { + return RiskDecision::deny(risk, "writes are disabled by the safety policy".into()); + } + if act.what == ActionKind::Disruptive && !self.allow_disruptive { + return RiskDecision::deny(risk, "disruptive actions (stop/restart/mode change) are disabled by the safety policy".into()); + } + if act.protocol.is_ot() && act.what == ActionKind::ProbeExploit && !self.allow_fuzzing { + // Exploit payloads at an industrial device are how scanners crash + // PLCs that have done nothing wrong: malformed input on a protocol + // with no input validation, answered by a CPU with no spare cycles. + return RiskDecision::deny(risk, "exploit payloads over industrial protocols are refused — malformed input is how these devices crash".into()); + } + if act.protocol.is_ot() && matches!(act.what, ActionKind::Write | ActionKind::Disruptive) { + // Belt and braces: even with writes enabled, an industrial write is + // its own decision and never a side effect of a permissive flag. + return RiskDecision::deny(risk, "writing over an industrial protocol requires an explicit, separate authorization".into()); + } + if risk.effective > self.max_effective_risk { + return RiskDecision::deny( + risk.clone(), + format!("{} exceeds the ceiling of {:.2}", risk.explain(), self.max_effective_risk), + ); + } + if risk.effective > self.confirm_above { + return RiskDecision::confirm(risk.clone(), format!("{} — above the approval threshold {:.2}", risk.explain(), self.confirm_above)); + } + RiskDecision::allow(risk) + } + + /// Is this industrial function code refused outright? + pub fn function_code_allowed(&self, code: u16) -> bool { + !self.forbidden_function_codes.contains(&code) + } + + pub fn summary(&self) -> String { + format!( + "env {} · ceiling {:.2} · confirm above {:.2} · write {} · disruptive {} · fuzz {} · {} req/min{}", + self.environment.as_str(), + self.max_effective_risk, + self.confirm_above, + if self.allow_write { "allowed" } else { "blocked" }, + if self.allow_disruptive { "allowed" } else { "blocked" }, + if self.allow_fuzzing { "allowed" } else { "blocked" }, + self.max_requests_per_minute, + if self.forbidden_function_codes.is_empty() { String::new() } else { format!(" · {} function code(s) blocked", self.forbidden_function_codes.len()) } + ) + } +} + +/// Allow / ask a human / refuse — with the arithmetic attached either way. +#[derive(Debug, Clone, PartialEq)] +pub enum RiskDecision { + Allow(Risk), + Confirm(Risk, String), + Deny(Risk, String), +} + +impl RiskDecision { + fn allow(r: Risk) -> Self { + RiskDecision::Allow(r) + } + fn confirm(r: Risk, why: String) -> Self { + RiskDecision::Confirm(r, why) + } + fn deny(r: Risk, why: String) -> Self { + RiskDecision::Deny(r, why) + } + pub fn allowed(&self) -> bool { + !matches!(self, RiskDecision::Deny(..)) + } + pub fn needs_human(&self) -> bool { + matches!(self, RiskDecision::Confirm(..)) + } + pub fn risk(&self) -> &Risk { + match self { + RiskDecision::Allow(r) | RiskDecision::Confirm(r, _) | RiskDecision::Deny(r, _) => r, + } + } + pub fn reason(&self) -> String { + match self { + RiskDecision::Allow(r) => r.explain(), + RiskDecision::Confirm(_, w) | RiskDecision::Deny(_, w) => w.clone(), + } + } +} + +/// How the agent is required to reason — the loop, made into rules. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ReasoningPolicy { + /// Capture the unmodified behaviour before sending a payload. Without it + /// there is nothing to compare against and every difference is a guess. + pub require_baseline_first: bool, + /// Hypotheses allowed in flight at once. Unbounded breadth is how a run + /// spends its budget touching everything shallowly. + pub max_open_hypotheses: usize, + /// Observations to gather before escalating from probe to exploit. + pub min_observations_before_exploit: usize, + /// Prefer the cheapest action that most reduces uncertainty (value of + /// information) over the most spectacular one. + pub voi_ordering: bool, + /// Give up on a hypothesis after this many failed attempts and write down + /// why, instead of retrying the same thing with different words. + pub max_attempts_per_hypothesis: usize, + /// End the run when no new information has arrived for this many rounds. + pub stop_after_idle_rounds: usize, + /// Re-derive nothing that the memory already knows. + pub consult_memory: bool, +} + +impl Default for ReasoningPolicy { + fn default() -> Self { + ReasoningPolicy { + require_baseline_first: true, + max_open_hypotheses: 8, + min_observations_before_exploit: 2, + voi_ordering: true, + max_attempts_per_hypothesis: 3, + stop_after_idle_rounds: 3, + consult_memory: true, + } + } +} + +impl ReasoningPolicy { + /// Rendered into prompts. These are expectations the harness also checks + /// where it can (baseline capture, evidence contract), not decoration. + pub fn prompt_block(&self) -> String { + let mut s = String::from("REASONING POLICY — how this engagement is expected to proceed:\n"); + if self.require_baseline_first { + s.push_str(" - Capture the baseline BEFORE sending any payload. A difference without a baseline is not evidence.\n"); + } + s.push_str(&format!(" - Keep at most {} hypotheses open; close one before opening another.\n", self.max_open_hypotheses)); + s.push_str(&format!(" - Gather at least {} independent observation(s) before escalating from probing to exploitation.\n", self.min_observations_before_exploit)); + if self.voi_ordering { + s.push_str(" - Choose the cheapest interaction that most reduces uncertainty, not the most impressive one.\n"); + } + s.push_str(&format!(" - After {} failed attempts at the same hypothesis, record WHY it failed and move on — do not retry it reworded.\n", self.max_attempts_per_hypothesis)); + s.push_str(&format!(" - If {} rounds pass with no new information, stop and report.\n", self.stop_after_idle_rounds)); + if self.consult_memory { + s.push_str(" - Use what is already known about this target (provided above) instead of re-deriving it.\n"); + } + s + } +} + +/// What a finding must carry before it may claim a given severity. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ProofOfImpactPolicy { + /// Critical/High need a deterministic validator verdict, not a vote. + pub require_validator_for_high: bool, + /// Repeats required before a behavioural difference counts. + pub min_reproductions: usize, + /// A state change must be read back; claiming a write that was never + /// verified is the most expensive kind of false positive. + pub require_read_back_for_writes: bool, + /// Impact must name data or a capability actually reached. + pub require_named_impact: bool, + /// Severity ceiling applied when proof is missing, instead of dropping the + /// finding — an unproven lead is still worth a human's time. + pub unproven_severity_cap: String, + /// Refuse to report at all when nothing at all backs the claim. + pub drop_unevidenced: bool, +} + +impl Default for ProofOfImpactPolicy { + fn default() -> Self { + ProofOfImpactPolicy { + require_validator_for_high: true, + min_reproductions: 2, + require_read_back_for_writes: true, + require_named_impact: true, + unproven_severity_cap: "Medium".into(), + drop_unevidenced: false, + } + } +} + +impl ProofOfImpactPolicy { + pub fn prompt_block(&self) -> String { + let mut s = String::from("PROOF-OF-IMPACT POLICY — what a claim must carry:\n"); + if self.require_validator_for_high { + s.push_str(" - High/Critical requires proof the harness can verify deterministically (see the evidence contract). A model's confidence is not proof.\n"); + } + s.push_str(&format!(" - A behavioural difference counts only if it reproduces at least {} time(s).\n", self.min_reproductions)); + if self.require_read_back_for_writes { + s.push_str(" - A state change must be READ BACK. A 200 on the write proves the request was accepted, not that anything changed.\n"); + } + if self.require_named_impact { + s.push_str(" - Impact must name the data or capability actually reached in THIS application. No generic consequences.\n"); + } + s.push_str(&format!(" - Without that proof the finding is reported at most as {} and flagged for review — do not inflate it.\n", self.unproven_severity_cap)); + s + } +} + +/// The three policies plus the environment, carried together. +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct EngagementPolicy { + #[serde(default)] + pub safety: SafetyPolicy, + #[serde(default)] + pub reasoning: ReasoningPolicy, + #[serde(default)] + pub proof: ProofOfImpactPolicy, +} + +impl EngagementPolicy { + pub fn web(env: Environment) -> Self { + EngagementPolicy { safety: SafetyPolicy::web_standard(env), ..Default::default() } + } + + /// OT: conservative safety, patient reasoning, and proof requirements that + /// do not push an agent toward "just try the write and see". + pub fn ot() -> Self { + EngagementPolicy { + safety: SafetyPolicy::ot_conservative(), + reasoning: ReasoningPolicy { + min_observations_before_exploit: 4, + max_open_hypotheses: 4, + ..Default::default() + }, + proof: ProofOfImpactPolicy { + // On a live process, "prove it by doing it" is not available, + // so a documented reachable capability is the ceiling of proof. + require_read_back_for_writes: false, + unproven_severity_cap: "High".into(), + ..Default::default() + }, + } + } + + pub fn prompt_block(&self) -> String { + let mut s = self.reasoning.prompt_block(); + s.push('\n'); + s.push_str(&self.proof.prompt_block()); + s.push('\n'); + s.push_str(&format!("SAFETY POLICY — enforced by the harness: {}\n", self.safety.summary())); + for n in &self.safety.notes { + s.push_str(&format!(" note: {n}\n")); + } + if self.safety.environment == Environment::OtProduction { + s.push_str( + " This is a LIVE INDUSTRIAL environment. Industrial protocols authenticate nothing: a write is not an exploit, it is the protocol working — addressed to a device that may be holding a valve. Observe, document reachability, and never act on the process.\n", + ); + } + s + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn the_formula_is_the_sum_of_five_terms_times_the_environment() { + let act = Act { + what: ActionKind::Read, // 0.05 + protocol: Protocol::Https, // 0.10 + blast: BlastRadius::SingleObject, // 0.05 + privilege: 0.2, + asset_criticality: 0.5, + target: "https://app.test".into(), + }; + let r = Risk::compute(&act, Environment::Production); // ×1.0 + assert!((r.effective - 0.90).abs() < 1e-9, "{}", r.explain()); + let lab = Risk::compute(&act, Environment::Lab); // ×0.3 + assert!((lab.effective - 0.27).abs() < 1e-9, "{}", lab.explain()); + } + + #[test] + fn the_same_action_costs_more_on_a_live_process() { + let act = Act { what: ActionKind::Enumerate, protocol: Protocol::Industrial, blast: BlastRadius::Segment, privilege: 0.2, asset_criticality: 0.9, target: "10.0.0.5".into() }; + let staging = Risk::compute(&act, Environment::Staging).effective; + let ot = Risk::compute(&act, Environment::OtProduction).effective; + assert!(ot > staging * 2.0, "staging {staging:.2} vs ot {ot:.2}"); + } + + #[test] + fn an_explanation_names_every_term() { + let r = Risk::compute(&Act::default(), Environment::Production); + for t in ["action", "asset", "protocol", "privilege", "blast", "env"] { + assert!(r.explain().contains(t), "{} missing from: {}", t, r.explain()); + } + } + + #[test] + fn ot_refuses_a_write_even_when_writes_were_enabled() { + let mut p = SafetyPolicy::ot_conservative(); + p.allow_write = true; // operator loosened the flag + let act = Act { what: ActionKind::Write, protocol: Protocol::Industrial, blast: BlastRadius::PhysicalProcess, privilege: 0.3, asset_criticality: 1.0, target: "plc-1".into() }; + match p.check(&act) { + RiskDecision::Deny(_, why) => assert!(why.contains("separate authorization"), "{why}"), + d => panic!("an industrial write must never ride on a permissive flag: {d:?}"), + } + } + + #[test] + fn ot_still_allows_looking_but_asks_first() { + let p = SafetyPolicy::ot_conservative(); + let act = Act { what: ActionKind::Read, protocol: Protocol::Industrial, blast: BlastRadius::SingleHost, privilege: 0.1, asset_criticality: 0.8, target: "plc-1".into() }; + let d = p.check(&act); + assert!(d.allowed(), "observation is where OT findings come from: {}", d.reason()); + assert!(d.needs_human(), "touching a live process at all is a human's call: {}", d.reason()); + } + + #[test] + fn ot_refuses_exploit_payloads_even_though_they_change_nothing() { + let p = SafetyPolicy::ot_conservative(); + let act = Act { what: ActionKind::ProbeExploit, protocol: Protocol::Industrial, blast: BlastRadius::SingleHost, privilege: 0.1, asset_criticality: 0.5, target: "plc-1".into() }; + match p.check(&act) { + RiskDecision::Deny(_, why) => assert!(why.contains("malformed input"), "{why}"), + d => panic!("malformed input is how PLCs crash: {d:?}"), + } + } + + #[test] + fn the_ceiling_refuses_and_the_threshold_asks() { + let p = SafetyPolicy::web_standard(Environment::Production); + let low = Act { what: ActionKind::Read, ..Default::default() }; + assert!(matches!(p.check(&low), RiskDecision::Allow(_))); + + let mid = Act { what: ActionKind::ProbeExploit, protocol: Protocol::Database, blast: BlastRadius::SharedService, privilege: 0.6, asset_criticality: 0.8, ..Default::default() }; + let d = p.check(&mid); + assert!(d.needs_human(), "{:?} → {}", d, d.reason()); + assert!(d.allowed(), "asking for approval is not a refusal"); + + let high = Act { what: ActionKind::Disruptive, protocol: Protocol::Industrial, blast: BlastRadius::PhysicalProcess, privilege: 1.0, asset_criticality: 1.0, ..Default::default() }; + assert!(!p.check(&high).allowed()); + } + + #[test] + fn a_safety_system_is_never_in_play() { + for env in [Environment::Lab, Environment::Production, Environment::OtProduction] { + let p = SafetyPolicy::web_standard(env); + let act = Act { what: ActionKind::Read, protocol: Protocol::SafetySystem, ..Default::default() }; + match p.check(&act) { + RiskDecision::Deny(_, why) => assert!(why.contains("forbidden protocol"), "{why}"), + d => panic!("safety instrumented systems are off limits in {}: {d:?}", env.as_str()), + } + } + } + + #[test] + fn dangerous_industrial_function_codes_are_refused() { + let p = SafetyPolicy::ot_conservative(); + for code in [5, 6, 8, 15, 16, 0x29] { + assert!(!p.function_code_allowed(code), "function code {code} must be blocked"); + } + for code in [1, 2, 3, 4] { + assert!(p.function_code_allowed(code), "read code {code} must stay available"); + } + } + + #[test] + fn industrial_ports_are_recognised_without_a_declaration() { + assert_eq!(Protocol::from_port(502), Protocol::Industrial); + assert_eq!(Protocol::from_port(20000), Protocol::Industrial); + assert_eq!(Protocol::from_port(102), Protocol::Industrial); + assert_eq!(Protocol::from_port(443), Protocol::Https); + assert!(Protocol::from_port(502).is_ot()); + } + + #[test] + fn ot_paces_itself() { + assert!(SafetyPolicy::ot_conservative().max_requests_per_minute <= 60); + assert!(!SafetyPolicy::ot_conservative().allow_fuzzing, "fuzzing a PLC crashes devices that have done nothing wrong"); + } + + #[test] + fn policies_render_into_a_prompt_that_states_the_rules() { + let p = EngagementPolicy::ot(); + let block = p.prompt_block(); + assert!(block.contains("REASONING POLICY")); + assert!(block.contains("PROOF-OF-IMPACT POLICY")); + assert!(block.contains("SAFETY POLICY")); + assert!(block.contains("LIVE INDUSTRIAL"), "an OT engagement must say so in the prompt"); + } + + #[test] + fn environment_parses_the_words_operators_actually_type() { + for (s, want) in [("prod", Environment::Production), ("SCADA", Environment::OtProduction), ("ics", Environment::OtProduction), ("staging", Environment::Staging), ("lab", Environment::Lab)] { + assert_eq!(Environment::parse(s), Some(want), "{s}"); + } + assert_eq!(Environment::parse("banana"), None); + } +}