use crate::estate::{BootOrder, Estate, Label, Refusal, StorageKind};
use crate::faults::{Fault, Faults};
use crate::render;
use serde_json::{json, Value};
use std::sync::{Arc, Mutex};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpListener, TcpStream};
pub struct Mock {
pub estate: Mutex<Estate>,
pub heard: Mutex<std::collections::BTreeMap<String, Heard>>,
pub calls: Option<Mutex<Vec<(u16, String, String, String)>>>,
pub log: bool,
#[cfg(feature = "test-inject")]
pub injections: Mutex<Vec<Injection>>,
}
impl Mock {
fn build(estate: Estate, calls: Option<Mutex<Vec<(u16, String, String, String)>>>, log: bool) -> Arc<Mock> {
Arc::new(Mock {
estate: Mutex::new(estate),
heard: Mutex::default(),
calls,
log,
#[cfg(feature = "test-inject")]
injections: Mutex::default(),
})
}
pub fn new(estate: Estate) -> Arc<Mock> {
Mock::build(estate, None, false)
}
pub fn recording(estate: Estate) -> Arc<Mock> {
Mock::build(estate, Some(Mutex::default()), false)
}
pub fn logging(estate: Estate) -> Arc<Mock> {
Mock::build(estate, None, true)
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Heard {
pub requests: u64,
pub account_calls: u64,
}
impl Mock {
fn heard(&self, credential: &str, method: &str, path: &str) {
let key = crate::digest::sha256_hex(credential.as_bytes());
let mut h = self.heard.lock().unwrap();
let e = h.entry(key).or_default();
e.requests += 1;
if method == "GET" && path.trim_end_matches('/') == "/1.3/account" {
e.account_calls += 1;
}
}
pub fn heard_from(&self, digest: &str) -> Heard {
self.heard.lock().unwrap().get(digest).copied().unwrap_or_default()
}
}
pub async fn serve(mock: Arc<Mock>, addr: &str) -> std::io::Result<(u16, tokio::task::JoinHandle<()>)> {
let listener = TcpListener::bind(addr).await?;
let port = listener.local_addr()?.port();
mock.estate.lock().unwrap().upload_base = format!("http://127.0.0.1:{port}");
let h = tokio::spawn(async move {
loop {
let Ok((sock, _)) = listener.accept().await else { continue };
let _ = sock.set_nodelay(true);
let m = mock.clone();
tokio::spawn(async move {
let _ = connection(m, sock).await;
});
}
});
Ok((port, h))
}
struct Req {
method: String,
path: String,
query: Vec<(String, String)>,
authorized: bool,
body: Value,
raw: Vec<u8>,
}
pub(crate) enum Outcome {
Reply { status: u16, body: Value },
Reset,
}
fn reply(status: u16, body: Value) -> Outcome {
Outcome::Reply { status, body }
}
fn refuse(r: Refusal) -> Outcome {
reply(r.status, render::error(r.code, &r.message))
}
async fn connection(mock: Arc<Mock>, mut sock: TcpStream) -> std::io::Result<()> {
let mut buf: Vec<u8> = Vec::with_capacity(8 * 1024);
loop {
let head_end = loop {
if let Some(i) = find_crlfcrlf(&buf) {
break i;
}
let mut chunk = [0u8; 4096];
let n = sock.read(&mut chunk).await?;
if n == 0 {
return Ok(());
}
buf.extend_from_slice(&chunk[..n]);
};
let head = String::from_utf8_lossy(&buf[..head_end]).to_string();
let mut lines = head.split("\r\n");
let Some(start) = lines.next() else { return Ok(()) };
let mut parts = start.split_whitespace();
let method = parts.next().unwrap_or("").to_string();
let target = parts.next().unwrap_or("/").to_string();
let mut content_length = 0usize;
let mut authorized = false;
let mut credential = String::new();
let mut chunked = false;
let mut keep_alive = true;
for l in lines {
let Some((k, v)) = l.split_once(':') else { continue };
let (k, v) = (k.trim().to_ascii_lowercase(), v.trim());
match k.as_str() {
"content-length" => content_length = v.parse().unwrap_or(0),
"authorization" => {
authorized = !v.is_empty();
credential = v.strip_prefix("Bearer ").unwrap_or(v).trim().to_string();
}
"transfer-encoding" => chunked = v.eq_ignore_ascii_case("chunked"),
"connection" => keep_alive = !v.eq_ignore_ascii_case("close"),
_ => {}
}
}
let body_start = head_end + 4;
while buf.len() < body_start + content_length {
let mut chunk = [0u8; 4096];
let n = sock.read(&mut chunk).await?;
if n == 0 {
return Ok(());
}
buf.extend_from_slice(&chunk[..n]);
}
let raw: Vec<u8> = buf[body_start..body_start + content_length].to_vec();
buf.drain(..body_start + content_length);
let outcome = if chunked && target.starts_with("/uploader/") {
reply(411, render::error("LENGTH_REQUIRED", "Content-Length is required"))
} else if chunked {
reply(
400,
render::error(
"MOCK_UPCLOUD_CHUNKED_REQUEST",
"the real API is never sent a chunked request by this estate, and the mock refuses to guess",
),
)
} else {
answer(&mock, &method, &target, authorized.then_some(credential.as_str()), raw)
};
match outcome {
Outcome::Reset => return Ok(()),
Outcome::Reply { status, body } => {
let text = if body.is_null() { String::new() } else { body.to_string() };
let head = format!(
"HTTP/1.1 {status} {}\r\nContent-Type: {}\r\nContent-Length: {}\r\nConnection: {}\r\n\r\n",
reason(status),
if status == 403 && text.contains("correlation_id") {
"application/problem+json"
} else {
"application/json"
},
text.len(),
if keep_alive { "keep-alive" } else { "close" }
);
let mut out = Vec::with_capacity(head.len() + text.len());
out.extend_from_slice(head.as_bytes());
out.extend_from_slice(text.as_bytes());
sock.write_all(&out).await?;
sock.flush().await?;
if !keep_alive {
return Ok(());
}
}
}
}
}
fn reason(s: u16) -> &'static str {
match s {
200 => "OK",
201 => "Created",
202 => "Accepted",
204 => "No Content",
400 => "Bad Request",
401 => "Unauthorized",
402 => "Payment Required",
403 => "Forbidden",
404 => "Not Found",
409 => "Conflict",
411 => "Length Required",
412 => "Precondition Failed",
429 => "Too Many Requests",
502 => "Bad Gateway",
503 => "Service Unavailable",
511 => "Network Authentication Required",
_ => "Error",
}
}
fn find_crlfcrlf(b: &[u8]) -> Option<usize> {
b.windows(4).position(|w| w == b"\r\n\r\n")
}
fn split_target(t: &str) -> (String, Vec<(String, String)>) {
match t.split_once('?') {
None => (t.to_string(), vec![]),
Some((p, q)) => {
let params = q
.split('&')
.filter(|s| !s.is_empty())
.map(|kv| {
let (k, v) = kv.split_once('=').unwrap_or((kv, ""));
(percent_decode(k), percent_decode(v))
})
.collect();
(p.to_string(), params)
}
}
}
fn percent_decode(s: &str) -> String {
let b = s.as_bytes();
let mut out = Vec::with_capacity(b.len());
let mut i = 0;
while i < b.len() {
match b[i] {
b'%' if i + 2 < b.len() => {
let h = (hex(b[i + 1]), hex(b[i + 2]));
if let (Some(a), Some(c)) = h {
out.push(a * 16 + c);
i += 3;
continue;
}
out.push(b[i]);
i += 1;
}
b'+' => {
out.push(b' ');
i += 1;
}
c => {
out.push(c);
i += 1;
}
}
}
String::from_utf8_lossy(&out).to_string()
}
fn hex(c: u8) -> Option<u8> {
match c {
b'0'..=b'9' => Some(c - b'0'),
b'a'..=b'f' => Some(c - b'a' + 10),
b'A'..=b'F' => Some(c - b'A' + 10),
_ => None,
}
}
fn route(mock: &Arc<Mock>, r: Req) -> Outcome {
if r.path.starts_with("/mock/") {
return mock_door(mock, &r);
}
if let Some(uuid) = r.path.strip_prefix("/uploader/session/") {
let uuid = uuid.to_string();
let mut e = mock.estate.lock().unwrap();
e.tick();
return match e.upload(&uuid, &r.raw) {
Err(x) => refuse(x),
Ok(im) => reply(200, render::import(&im)),
};
}
if !r.authorized {
return reply(
401,
render::error("AUTHENTICATION_FAILED", "no Authorization header was sent"),
);
}
let Some(rest) = r.path.strip_prefix("/1.3") else {
return reply(404, render::not_implemented(&r.method, &r.path));
};
if rest.contains("//") {
return reply(404, render::error("NOT_FOUND", "Not found."));
}
{
let e = mock.estate.lock().unwrap();
if e.faults.fires(Fault::DeadToken) {
return reply(401, render::error("AUTHENTICATION_FAILED", "Authentication failed using the given username and password."));
}
if r.method == "GET" && e.faults.fires(Fault::ReadBadGateway) {
return reply(502, render::error("BAD_GATEWAY", "Bad Gateway"));
}
}
let seg: Vec<&str> = rest.trim_matches('/').split('/').filter(|s| !s.is_empty()).collect();
let mut e = mock.estate.lock().unwrap();
e.tick();
let m = r.method.as_str();
match (m, seg.as_slice()) {
("GET", ["price"]) => {
if e.faults.fires(Fault::PriceTransportReset) {
return Outcome::Reset;
}
let zone = e.zone.clone();
reply(200, render::price(&zone))
}
("GET", ["account"]) => {
if e.faults.fires(Fault::RevokedCredential) {
let cid = e.next_correlation_id();
return reply(403, render::auth_failed(&cid));
}
reply(200, json!({"account": {"username": "mock", "credits": 100_000.0}}))
}
("GET", ["zone"]) => reply(
200,
json!({"zones": {"zone": [{"id": e.zone, "description": "Stockholm #1", "public": "yes"}]}}),
),
("GET", ["server"]) => {
let labels = label_filters(&r.query);
let created = !e.faults.fires(Fault::WithholdCreatedField);
let rows: Vec<Value> = e
.servers_matching(&labels)
.into_iter()
.map(|s| render::server(s, false, created))
.collect();
reply(200, json!({"servers": {"server": rows}}))
}
("GET", ["server", uuid]) => {
if e.faults.fires(Fault::RevokedCredential) {
let cid = e.next_correlation_id();
return reply(403, render::auth_failed(&cid));
}
let created = !e.faults.fires(Fault::WithholdCreatedField);
match e.server(uuid) {
Some(_) if e.faults.fires(Fault::DetailNotFoundForListedServer) => {
reply(404, render::error("SERVER_NOT_FOUND", &format!("server {uuid} not found")))
}
Some(s) => reply(200, json!({"server": render::server(s, true, created)})),
None => reply(404, render::error("SERVER_NOT_FOUND", &format!("server {uuid} not found"))),
}
}
("GET", ["server", _uuid, "firewall_rule"]) | ("GET", ["server", _uuid, "firewall_rule", _]) => {
let uuid = seg[1];
if e.server(uuid).is_some() && !e.faults.fires(Fault::FirewallForbidden) {
let rules = e.rules(uuid).unwrap_or(&[]).to_vec();
if let Some(pos) = seg.get(3) {
return match rules.iter().find(|r| r.position == *pos) {
Some(r) => reply(200, json!({ "firewall_rule": crate::tf::render_rule(r) })),
None => reply(
404,
render::error("FIREWALL_RULE_NOT_FOUND", &format!("no rule at position {pos}")),
),
};
}
return reply(200, crate::tf::render_rules(&rules));
}
let cid = e.next_correlation_id();
reply(403, render::auth_failed(&cid))
}
("PUT", ["server", _uuid, "firewall_rule"]) | ("POST", ["server", _uuid, "firewall_rule"]) => {
let uuid = seg[1].to_string();
let mut rules = crate::tf::rules_from_body(&r.body);
if m == "POST" {
let mut existing = e.rules(&uuid).unwrap_or(&[]).to_vec();
existing.append(&mut rules);
rules = existing;
}
match e.set_rules(&uuid, rules) {
Err(x) => refuse(x),
Ok(()) => {
let rules = e.rules(&uuid).unwrap_or(&[]).to_vec();
reply(if m == "POST" { 201 } else { 200 }, crate::tf::render_rules(&rules))
}
}
}
("DELETE", ["server", _uuid, "firewall_rule", _pos]) => {
let uuid = seg[1].to_string();
let pos = seg[3].to_string();
let kept: Vec<_> = e.rules(&uuid).unwrap_or(&[]).iter().filter(|r| r.position != pos).cloned().collect();
match e.set_rules(&uuid, kept) {
Err(x) => refuse(x),
Ok(()) => reply(204, Value::Null),
}
}
("GET", ["plan"]) => reply(200, crate::tf::plans()),
("POST", ["server"]) => {
let b = &r.body["server"];
let plan = b["plan"].as_str().unwrap_or("1xCPU-1GB").to_string();
if !render::plan_known(&plan) {
return refuse(Refusal::new(400, "INVALID_PLAN", format!("no such plan: {plan}")));
}
let title = b["title"].as_str().unwrap_or("").to_string();
let hostname = b["hostname"].as_str().unwrap_or(&title).to_string();
let zone = b["zone"].as_str().unwrap_or(&e.zone).to_string();
let labels = read_server_labels(&b["labels"]);
let devs: Vec<Value> = b["storage_devices"]["storage_device"].as_array().cloned().unwrap_or_default();
let dev = devs.first().cloned().unwrap_or(Value::Null);
let disk_title = dev["title"].as_str().unwrap_or("boot").to_string();
let disk_gib = dev["size"].as_u64().unwrap_or(20);
let mut attach: Vec<(String, String, Option<String>)> = Vec::new();
for d in devs.iter().skip(1) {
let Some(st) = d["storage"].as_str().map(str::to_string) else { continue };
let kind = if d["type"].as_str() == Some("cdrom") { "cdrom" } else { "disk" }.to_string();
match e.storage(&st) {
None => {
return refuse(Refusal::new(404, "STORAGE_NOT_FOUND", format!("storage {st} not found")))
}
Some(x) if x.state != "online" => {
return refuse(Refusal::new(
409,
"STORAGE_STATE_ILLEGAL",
format!("storage {st} is {} — wait for online", x.state),
))
}
Some(_) => attach.push((st, kind, d["address"].as_str().map(str::to_string))),
}
}
match e.create_server(&title, &hostname, &plan, &zone, labels, &disk_title, disk_gib) {
Err(x) => refuse(x),
Ok(uuid) => {
let ifaces = crate::tf::interfaces_from_body(b);
let boot = b["boot_order"].as_str().and_then(BootOrder::parse);
let firewall_on = b["firewall"].as_str().unwrap_or("off") == "on";
let metadata = b["metadata"].as_str().unwrap_or("yes") != "no";
let tz = b["timezone"].as_str().unwrap_or("UTC").to_string();
if let Err(x) = e.configure_server(&uuid, ifaces, boot, firewall_on, metadata, &tz) {
return refuse(x);
}
for (st, kind, want) in &attach {
if let Err(x) = e.attach_at_create(&uuid, st, kind, want.as_deref()) {
return refuse(x);
}
}
if e.faults.fires(Fault::CommitThenDropReply) {
return Outcome::Reset;
}
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.server(&uuid).expect("just created");
reply(201, json!({"server": render::server(s, true, created)}))
}
}
}
("PUT", ["server", uuid]) => {
let b = &r.body["server"];
let plan = b["plan"].as_str().map(str::to_string);
let bo = b["boot_order"].as_str().and_then(parse_boot_order);
let labels = if b["labels"].is_null() { None } else { Some(read_server_labels(&b["labels"])) };
let ra = b["remote_access_enabled"].as_str().map(|s| s == "yes");
let rap = b["remote_access_password"].as_str().map(str::to_string);
let rename = (b["hostname"].as_str().map(str::to_string), b["title"].as_str().map(str::to_string));
if let Some(p) = &plan {
if !render::plan_known(p) {
return refuse(Refusal::new(400, "INVALID_PLAN", format!("no such plan: {p}")));
}
}
match e.modify_server(uuid, plan.as_deref(), bo, labels, ra, rap.as_deref()) {
Err(x) => refuse(x),
Ok(()) => {
e.rename_server(uuid, rename.0.as_deref(), rename.1.as_deref());
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.server(uuid).expect("modified");
reply(202, json!({"server": render::server(s, true, created)}))
}
}
}
("DELETE", ["server", uuid]) => {
let with_storages = r
.query
.iter()
.any(|(k, v)| k == "storages" && (v == "1" || v == "true"));
match e.delete_server(uuid, with_storages) {
Err(x) => refuse(x),
Ok(()) => reply(204, Value::Null),
}
}
("POST", ["server", uuid, "start"]) => match e.start_server(uuid) {
Err(x) => refuse(x),
Ok(()) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.server(uuid).expect("started");
reply(200, json!({"server": render::server(s, true, created)}))
}
},
("POST", ["server", uuid, "stop"]) => {
if !r.body["stop_server"]["timeout_action"].is_null() {
return refuse(Refusal::new(
400,
"INVALID_STOP_SERVER",
"stop_server has no attribute timeout_action (it belongs to restart_server)",
));
}
let hard = r.body["stop_server"]["stop_type"].as_str() == Some("hard");
match e.stop_server(uuid, hard) {
Err(x) => refuse(x),
Ok(()) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.server(uuid).expect("stopping");
reply(200, json!({"server": render::server(s, true, created)}))
}
}
}
("POST", ["server", uuid, "restart"]) => {
let uuid = uuid.to_string();
if let Err(x) = e.stop_server(&uuid, false) {
return refuse(x);
}
e.run_to_quiet();
match e.start_server(&uuid) {
Err(x) => refuse(x),
Ok(()) => reply(200, json!({"server": {"uuid": uuid}})),
}
}
("POST", ["server", uuid, "storage", "attach"]) => {
let d = &r.body["storage_device"];
let storage = d["storage"].as_str().unwrap_or("").to_string();
let kind = d["type"].as_str().unwrap_or("disk").to_string();
let want = d["address"].as_str().map(str::to_string);
match e.attach_at(uuid, &storage, &kind, want.as_deref()) {
Err(x) => refuse(x),
Ok(_) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.server(uuid).expect("attached");
reply(200, json!({"server": render::server(s, true, created)}))
}
}
}
("POST", ["server", uuid, "storage", "detach"]) => {
let address = r.body["storage_device"]["address"].as_str().unwrap_or("").to_string();
match e.detach(uuid, &address) {
Err(x) => refuse(x),
Ok(()) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.server(uuid).expect("detached");
reply(200, json!({"server": render::server(s, true, created)}))
}
}
}
("POST", ["server", uuid, "cdrom", "eject"]) => match e.eject(uuid) {
Err(x) => refuse(x),
Ok(()) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.server(uuid).expect("ejected");
reply(200, json!({"server": render::server(s, true, created)}))
}
},
("POST", ["server", uuid, "cdrom", "load"]) => {
let storage = r.body["storage_device"]["storage"].as_str().unwrap_or("").to_string();
match e.load_cdrom(uuid, &storage) {
Err(x) => refuse(x),
Ok(()) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.server(uuid).expect("loaded");
reply(200, json!({"server": render::server(s, true, created)}))
}
}
}
("GET", ["storage"]) | ("GET", ["storage", "private"]) => {
let private_only = seg.len() == 2;
let labels = label_filters(&r.query);
let created = !e.faults.fires(Fault::WithholdCreatedField);
let rows: Vec<Value> = e
.storages_matching(&labels, private_only)
.into_iter()
.map(|s| render::storage(&e, s, false, created))
.collect();
reply(200, json!({"storages": {"storage": rows}}))
}
("GET", ["storage", filter @ ("public" | "template" | "favorite")]) => {
let labels = label_filters(&r.query);
let created = !e.faults.fires(Fault::WithholdCreatedField);
let rows: Vec<Value> = e
.storages_matching(&labels, false)
.into_iter()
.filter(|s| match *filter {
"favorite" => false,
_ => s.kind == StorageKind::Template,
})
.map(|s| render::storage(&e, s, false, created))
.collect();
reply(200, json!({"storages": {"storage": rows}}))
}
("GET", ["storage", uuid]) => {
if e.faults.fires(Fault::RevokedCredential) {
let cid = e.next_correlation_id();
return reply(403, render::auth_failed(&cid));
}
let created = !e.faults.fires(Fault::WithholdCreatedField);
match e.storage(uuid) {
Some(s) => reply(200, json!({"storage": render::storage(&e, s, true, created)})),
None => reply(404, render::error("STORAGE_NOT_FOUND", &format!("storage {uuid} not found"))),
}
}
("POST", ["storage"]) => {
let b = &r.body["storage"];
let title = b["title"].as_str().unwrap_or("").to_string();
let size = b["size"].as_u64().or_else(|| b["size"].as_str().and_then(|s| s.parse().ok())).unwrap_or(0);
let tier = b["tier"].as_str().unwrap_or("maxiops").to_string();
let zone = b["zone"].as_str().unwrap_or(&e.zone).to_string();
let labels = read_flat_labels(&b["labels"]);
match e.create_storage(&title, size, &tier, &zone, labels) {
Err(x) => refuse(x),
Ok(uuid) => {
if e.faults.fires(Fault::CommitThenDropReply) {
return Outcome::Reset;
}
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.storage(&uuid).expect("just created");
reply(201, json!({"storage": render::storage(&e, s, true, created)}))
}
}
}
("PUT", ["storage", uuid]) => {
let b = &r.body["storage"];
let size = b["size"].as_u64().or_else(|| b["size"].as_str().and_then(|s| s.parse().ok()));
let title = b["title"].as_str().map(str::to_string);
match e.modify_storage(uuid, size, title.as_deref()) {
Err(x) => refuse(x),
Ok(()) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.storage(uuid).expect("modified");
reply(200, json!({"storage": render::storage(&e, s, true, created)}))
}
}
}
("DELETE", ["storage", uuid]) => match e.delete_storage(uuid) {
Err(x) => refuse(x),
Ok(()) => reply(204, Value::Null),
},
("POST", ["storage", uuid, "resize"]) => match e.resize_filesystem(uuid) {
Err(x) => refuse(x),
Ok(backup) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let b = e.storage(&backup).expect("just minted");
reply(200, json!({"resize_backup": render::storage(&e, b, true, created)}))
}
},
("POST", ["storage", uuid, "import"]) => {
let source = r.body["storage_import"]["source"].as_str().unwrap_or("direct_upload").to_string();
match e.start_import(uuid, &source) {
Err(x) => refuse(x),
Ok(im) => reply(201, json!({"storage_import": render::import(&im)})),
}
}
("GET", ["storage", uuid, "import"]) => match e.storage(uuid).and_then(|s| s.import.as_ref()) {
None => reply(404, render::error("STORAGE_IMPORT_NOT_FOUND", &format!("no import session on {uuid}"))),
Some(im) => reply(200, json!({"storage_import": render::import(im)})),
},
("POST", ["storage", uuid, "clone"]) => {
let title = r.body["storage"]["title"].as_str().unwrap_or("clone").to_string();
match e.clone_storage(uuid, &title) {
Err(x) => refuse(x),
Ok(new) => {
let created = !e.faults.fires(Fault::WithholdCreatedField);
let s = e.storage(&new).expect("just cloned");
reply(201, json!({"storage": render::storage(&e, s, true, created)}))
}
}
}
_ => reply(404, render::not_implemented(m, &r.path)),
}
}
pub fn probe(mock: &Arc<Mock>, method: &str, target: &str, body: Value) -> (u16, Value) {
let (path, query) = split_target(target);
let raw = if body.is_null() { vec![] } else { body.to_string().into_bytes() };
let r = Req { method: method.to_string(), path, query, authorized: true, body, raw };
match route(mock, r) {
Outcome::Reset => (0, Value::Null),
Outcome::Reply { status, body } => (status, body),
}
}
fn parse_boot_order(s: &str) -> Option<BootOrder> {
if let Some(b) = BootOrder::parse(s) {
return Some(b);
}
match s.split(',').next().map(str::trim) {
Some("cdrom") => Some(BootOrder::Cdrom),
Some("disk") => Some(BootOrder::Disk),
_ => None,
}
}
fn label_filters(q: &[(String, String)]) -> Vec<(String, String)> {
q.iter()
.filter(|(k, _)| k == "label")
.filter_map(|(_, v)| v.split_once('=').map(|(a, b)| (a.to_string(), b.to_string())))
.collect()
}
fn read_flat_labels(v: &Value) -> Vec<Label> {
v.as_array()
.map(|a| {
a.iter()
.filter_map(|l| {
Some(Label {
key: l["key"].as_str()?.to_string(),
value: l["value"].as_str().unwrap_or("").to_string(),
})
})
.collect()
})
.unwrap_or_default()
}
fn read_server_labels(v: &Value) -> Vec<Label> {
if v["label"].is_array() {
read_flat_labels(&v["label"])
} else {
read_flat_labels(v)
}
}
fn mock_door(mock: &Arc<Mock>, r: &Req) -> Outcome {
let seg: Vec<&str> = r.path.trim_matches('/').split('/').skip(1).collect();
let mut e = mock.estate.lock().unwrap();
match (r.method.as_str(), seg.as_slice()) {
("GET", ["seed"]) => reply(200, json!({"seed": e.seed(), "lays": e.lays(), "now_ms": e.clock.now_ms()})),
("GET", ["heard"]) => {
let d = q(&r.query, "token_sha256");
let h = mock.heard_from(&d);
reply(200, json!({"token_sha256": d, "requests": h.requests, "account_calls": h.account_calls}))
}
("GET", ["inbound", uuid]) => {
e.settle();
let port: u16 = q(&r.query, "port").parse().unwrap_or(22);
let proto = { let p = q(&r.query, "proto"); if p.is_empty() { "tcp".to_string() } else { p } };
match e.inbound(&q(&r.query, "from"), uuid, &proto, port) {
Err(x) => reply(x.status, json!({"error": x.message})),
Ok(reach) => reply(200, json!({"ok": reach.is_ok(), "why": reach.why()})),
}
}
("GET", ["udp-reply", uuid]) => {
let port: u16 = q(&r.query, "port").parse().unwrap_or(53);
match e.udp_reply_arrives(uuid, &q(&r.query, "from"), port) {
Err(x) => reply(x.status, json!({"error": x.message})),
Ok(yes) => reply(200, json!({"arrives": yes})),
}
}
("GET", ["disks", uuid]) => reply(
200,
json!({"disks": e.guest_disk_names(uuid).into_iter().map(|(a, n)| json!({"address": a, "name": n})).collect::<Vec<_>>()}),
),
("GET", ["estate"]) => {
if !r.query.iter().any(|(k, _)| k == "as_is") {
e.run_to_quiet();
}
let servers: Vec<Value> = e
.all_servers()
.map(|s| json!({"uuid": s.uuid, "title": s.title, "state": s.state, "plan": s.plan,
"guest": format!("{:?}", s.guest), "public_ip": s.public_ip,
"utility_ip": s.utility_ip, "vnc_port": s.vnc_port,
"reported_vnc_port": s.reported_vnc_port, "boot_order": s.boot_order.as_str()}))
.collect();
let storages: Vec<Value> = e
.all_storages()
.filter(|s| s.kind != StorageKind::Template)
.map(|s| json!({"uuid": s.uuid, "title": s.title, "state": s.state, "size": s.size_gib,
"type": s.kind.as_str(), "origin": s.origin,
"labels": s.labels.iter().map(|l| format!("{}={}", l.key, l.value)).collect::<Vec<_>>()}))
.collect();
reply(200, json!({"servers": servers, "storages": storages}))
}
("POST", ["advance", ms]) => {
let ms: u64 = ms.parse().unwrap_or(0);
e.clock.advance_ms(ms);
e.settle();
reply(200, json!({"now_ms": e.clock.now_ms()}))
}
("POST", ["fault", name, verb]) => match Fault::parse(name) {
None => reply(404, json!({"error": format!("no such fault: {name}")})),
Some(f) => {
match *verb {
"arm" => e.faults.arm(f),
"disarm" => e.faults.disarm(f),
_ => return reply(400, json!({"error": "arm or disarm"})),
}
reply(200, json!({"fault": f.name(), "armed": e.faults.is_armed(f)}))
}
},
("GET", ["clock", uuid]) => {
e.settle();
let Some(s) = e.server(uuid) else {
return reply(404, json!({"error": format!("no such server: {uuid}")}));
};
let t = crate::guest_clock::days_from_civil(2026, 9, 20) * 86_400;
let skew = s.clock_skew_ms(t);
let reds = crate::guest_clock::cascade(skew, crate::guest_clock::SkewWindow::default());
reply(
200,
json!({
"uuid": s.uuid,
"zone": s.zone,
"hypervisor_rtc": "UTC",
"guest_reads_rtc_as": format!("{:?}", s.rtc),
"guest_clock_skew_ms": skew,
"udp_reply_arrives": crate::guest_clock::udp_reply_arrives(&e.faults),
"red": reds.iter().map(|r| json!({"row": r.row, "why": r.why})).collect::<Vec<_>>(),
}),
)
}
("GET", ["reach", from]) => {
e.settle();
let dest = r.query.iter().find(|(k, _)| k == "dest").map(|(_, v)| v.clone()).unwrap_or_default();
let port: u16 = r
.query
.iter()
.find(|(k, _)| k == "port")
.and_then(|(_, v)| v.parse().ok())
.unwrap_or(22);
match e.reach(from, &dest, port) {
Err(x) => reply(x.status, json!({"error": x.message})),
Ok(reach) => reply(
200,
json!({
"from": from, "dest": dest, "port": port,
"ok": reach.is_ok(),
"why": reach.why(),
"kind": match &reach {
crate::net::Reach::Ok => "ok",
crate::net::Reach::NoRouteOutbound { .. } => "no-route-outbound",
crate::net::Reach::NoHairpin { .. } => "no-hairpin",
crate::net::Reach::Refused { .. } => "refused",
crate::net::Reach::Dropped { .. } => "dropped",
},
"inbound_ok": crate::net::inbound_reaches(true),
}),
),
}
}
("GET", ["dhcp", uuid]) => match e.server(uuid) {
None => reply(404, json!({"error": format!("no such server: {uuid}")})),
Some(s) => {
let o = s.dhcp_offer();
reply(
200,
json!({
"address": o.address,
"prefix": o.prefix,
"router": o.router,
"option_121": o.classless_static_routes.iter().map(|x| x.to_string()).collect::<Vec<_>>(),
"guest_dhcp_client": format!("{:?}", s.dhcp_client),
}),
)
}
},
("POST", ["dnat", on_server, port, to_address]) => {
let port: u16 = port.parse().unwrap_or(0);
e.dnat.push(crate::net::Dnat {
on_server: on_server.to_string(),
port,
to_address: to_address.to_string(),
to_port: port,
});
reply(200, json!({"rules": e.dnat.len()}))
}
("GET", ["hostkey", uuid]) => {
e.settle();
let mut keys = serde_json::Map::new();
for p in crate::estate::HostKeyPath::ALL {
match e.host_key_via(uuid, p) {
Err(x) => return reply(x.status, json!({"error": x.message})),
Ok(k) => {
keys.insert(p.name().to_string(), json!(k));
}
}
}
let distinct: std::collections::BTreeSet<&str> =
keys.values().filter_map(|v| v.as_str()).collect();
reply(
200,
json!({
"paths": keys,
"agree": distinct.len() == 1,
"verdict": if distinct.len() == 1 {
"one key on three paths: this is the machine that was re-imaged"
} else {
"the paths disagree: a name is answering for a machine that is not behind the DNAT"
},
}),
)
}
("GET", ["reimage", uuid]) => {
let (uart, wall) = e.timings.reimage_observers(e.seed(), uuid);
reply(
200,
json!({
"guest_uart_ms": uart,
"ladder_wall_ms": wall,
"ratio": wall / uart.max(1),
"note": "the installer is not slow; the provider is. create, media sync, firmware, boot order, DHCP.",
"observers": {
"guest_uart_ms": "the guest's own PID 1, from inside",
"ladder_wall_ms": "the ladder's `install-time`, wall, from outside"
}
}),
)
}
("GET", ["guest", uuid]) => {
e.tick();
match e.engine.evidence(uuid) {
Some(v) => reply(200, json!({"engine": e.engine.name(), "guest": v})),
None => reply(404, json!({"error": format!("engine {} has no machine for {uuid}", e.engine.name())})),
}
}
("POST", ["relay"]) => {
e.relay();
reply(200, json!({"lays": e.lays()}))
}
("POST", ["seed", s]) => {
let seed: u64 = s.parse().unwrap_or(0);
let speed = e.clock.speed_milli();
let engine = e.engine.clone();
engine.forget_all();
*e = Estate::new(crate::Clock::new(speed), Faults::seeded(seed), seed).with_engine(engine);
reply(200, json!({"seed": seed}))
}
_ => reply(404, json!({"error": format!("no such mock door: {}", r.path)})),
}
}
fn q(query: &[(String, String)], key: &str) -> String {
query.iter().find(|(k, _)| k == key).map(|(_, v)| v.clone()).unwrap_or_default()
}
pub(crate) fn answer(mock: &Arc<Mock>, method: &str, target: &str, credential: Option<&str>, raw: Vec<u8>) -> Outcome {
let (path, query) = split_target(target);
let body: Value = if raw.is_empty() { Value::Null } else { serde_json::from_slice(&raw).unwrap_or(Value::Null) };
let said = format!("{method} {path}");
let (m_rec, p_rec) = (method.to_string(), path.clone());
let authorized = credential.is_some();
if let Some(c) = credential {
mock.heard(c, method, &path);
}
#[cfg(feature = "test-inject")]
let injected = mock.injected(method, &path);
#[cfg(not(feature = "test-inject"))]
let injected: Option<Outcome> = None;
let out = match injected {
Some(o) => o,
None => route(mock, Req { method: method.to_string(), path, query, authorized, body, raw }),
};
if let Some(calls) = &mock.calls {
let (st, code) = match &out {
Outcome::Reset => (0, String::new()),
Outcome::Reply { status, body } => (*status, body["error"]["error_code"].as_str().unwrap_or("").to_string()),
};
calls.lock().unwrap().push((st, m_rec, p_rec, code));
}
if mock.log {
let states = {
let e = mock.estate.lock().unwrap();
e.all_servers().map(|s| format!("{}={}", &s.uuid[..4], s.state)).collect::<Vec<_>>().join(" ")
};
let code = match &out {
Outcome::Reset => 0,
Outcome::Reply { status, .. } => *status,
};
eprintln!(" {code:>3} {said:<52} [{states}]");
}
out
}
#[cfg(feature = "test-inject")]
pub struct Injection {
pub method: String,
pub matches: Box<dyn Fn(&str, &Estate) -> bool + Send + Sync>,
pub status: u16,
pub code: String,
pub message: String,
pub times: u32,
}
#[cfg(feature = "test-inject")]
impl Mock {
pub fn inject(&self, i: Injection) {
self.injections.lock().unwrap().push(i);
}
fn injected(&self, method: &str, path: &str) -> Option<Outcome> {
let e = self.estate.lock().unwrap();
let mut all = self.injections.lock().unwrap();
let hit = all.iter_mut().find(|i| i.times > 0 && i.method == method && (i.matches)(path, &e))?;
hit.times -= 1;
Some(reply(hit.status, render::error(&hit.code, &hit.message)))
}
}