use std::io::Write as _;
use std::path::{Path, PathBuf};
use std::process::{Command, Output};
fn private_dir(tag: &str) -> PathBuf {
use std::os::unix::fs::PermissionsExt;
let dir = std::env::temp_dir().join(format!("gw-cli-{}-{tag}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("create dir");
std::fs::set_permissions(&dir, std::fs::Permissions::from_mode(0o700)).expect("chmod");
dir
}
fn gw(socket: &Path, line: &str) -> Output {
gw_env(line, &[("GWK_SOCKET_PATH", &socket.to_string_lossy())])
}
fn gw_env(line: &str, env: &[(&str, &str)]) -> Output {
let mut command = Command::new(env!("CARGO_BIN_EXE_gw"));
command.args(line.split_whitespace());
for (name, value) in env {
command.env(name, value);
}
command.output().expect("run gw")
}
fn code(output: &Output) -> i32 {
output.status.code().expect("gw exited on a signal")
}
fn json(output: &Output) -> serde_json::Value {
let text = String::from_utf8_lossy(&output.stdout);
serde_json::from_str(text.trim()).unwrap_or_else(|e| panic!("stdout is not JSON ({e}): {text}"))
}
#[test]
fn the_command_tree_is_printed_as_prose_and_everything_else_as_json() {
let dir = private_dir("help");
let socket = dir.join("absent.sock");
let help = gw(&socket, "--help");
assert_eq!(code(&help), 0);
let text = String::from_utf8_lossy(&help.stdout);
assert!(text.starts_with("gw —"), "{text}");
assert!(text.contains("gw kernel health"), "{text}");
let info = gw(&socket, "build-info");
assert_eq!(code(&info), 0);
let info = json(&info);
assert_eq!(info["type"], "build_info");
assert!(info["contract_version"].is_number(), "{info}");
assert!(info.get("public_revision").is_some(), "{info}");
if let Some(revision) = info["public_revision"].as_str() {
assert_eq!(revision.len(), 40, "{revision}");
}
}
#[test]
fn a_mistake_exits_two_and_says_what_it_was() {
let dir = private_dir("usage");
let socket = dir.join("absent.sock");
for line in ["nonsense", "projection list tsak", "kernel activate"] {
let out = gw(&socket, line);
assert_eq!(
code(&out),
2,
"{line}: {:?}",
String::from_utf8_lossy(&out.stdout)
);
let answer = json(&out);
assert_eq!(answer["type"], "error");
assert_eq!(answer["code"], "validation");
}
let answer = json(&gw(&socket, "projection list tsak"));
let message = answer["message"].as_str().expect("message");
assert!(message.contains("attention_item"), "{message}");
}
#[test]
fn a_daemon_that_is_not_there_is_unavailable_rather_than_a_crash() {
let dir = private_dir("absent");
let socket = dir.join("absent.sock");
let out = gw(&socket, "kernel health");
assert_eq!(code(&out), 5);
let answer = json(&out);
assert_eq!(answer["code"], "storage");
assert!(
answer["message"]
.as_str()
.expect("message")
.contains("absent.sock"),
"{answer}"
);
}
const RUNTIME_ROLE: &str = "gwk_cli_runtime";
const TEST_KEK: [u8; 32] = [0x5a; 32];
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn the_binary_talks_to_a_real_daemon() {
use gwk_kernel::blob::store::PgBlobStore;
use gwk_kernel::config::{ADMIN_DATABASE_URL_ENV, AdminConfig, BlobConfig, RUNTIME_ROLE_ENV};
use gwk_kernel::wire::listen::Listener;
use gwk_kernel::wire::serve::Daemon;
use gwk_kernel::{PgEventStore, admin, connect_pool};
use secrecy::SecretString;
use sqlx::PgPool;
let admin_url = std::env::var("GWK_TEST_ADMIN_DATABASE_URL")
.expect("GWK_TEST_ADMIN_DATABASE_URL must point at a PostgreSQL superuser DSN");
let maintenance = PgPool::connect(&admin_url).await.expect("maintenance");
let _ = sqlx::raw_sql(sqlx::AssertSqlSafe(format!(
"CREATE ROLE {RUNTIME_ROLE} NOLOGIN;"
)))
.execute(&maintenance)
.await;
let name = format!("gwk_cli_{}", std::process::id());
let url = {
let (prefix, _) = admin_url.rsplit_once('/').expect("a /database suffix");
format!("{prefix}/{name}")
};
let _ = sqlx::raw_sql(sqlx::AssertSqlSafe(format!(
"DROP DATABASE IF EXISTS {name} WITH (FORCE);"
)))
.execute(&maintenance)
.await;
sqlx::raw_sql(sqlx::AssertSqlSafe(format!("CREATE DATABASE {name};")))
.execute(&maintenance)
.await
.expect("create database");
let pool = connect_pool(&SecretString::from(url.clone()), 8)
.await
.expect("connect");
let config = AdminConfig::from_lookup(move |key| match key {
ADMIN_DATABASE_URL_ENV => Some(url.clone()),
RUNTIME_ROLE_ENV => Some(RUNTIME_ROLE.to_owned()),
_ => None,
})
.expect("admin config");
admin::init(&pool, &config).await.expect("init");
let dir = private_dir("live");
let blob_root = dir.join("blobs");
let blobs = PgBlobStore::open(
pool.clone(),
BlobConfig::new(blob_root.clone(), TEST_KEK, "kek-test".to_owned()).expect("blob config"),
)
.await
.expect("blob store");
let store = PgEventStore::open(pool).await.expect("store");
store
.ensure_genesis(&"a1b2c3d4e5".repeat(4))
.await
.expect("genesis");
let socket = dir.join("gwk.sock");
let listener = Listener::bind(&socket).await.expect("bind");
let daemon = std::sync::Arc::new(
Daemon::new(store.with_blobs(blobs), "a1b2c3d4e5".repeat(4)).expect("daemon"),
);
let (stop, stopped) = tokio::sync::oneshot::channel::<()>();
let serving = tokio::spawn(async move {
let _ = gwk_kernel::wire::serve::run(listener, daemon, async move {
let _ = stopped.await;
})
.await;
});
let socket_for_blocking = socket.clone();
let answers = tokio::task::spawn_blocking(move || {
let socket = socket_for_blocking.as_path();
let health = gw(socket, "kernel health");
let status = gw(socket, "kernel status");
let events = gw(socket, "event read --limit 1");
let tasks = gw(socket, "projection list task");
let missing = gw(socket, "projection get task t-nope");
let payload = socket.with_file_name("payload.bin");
let mut file = std::fs::File::create(&payload).expect("create");
let plaintext: Vec<u8> = (0..5_000u32).map(|i| (i % 251) as u8).collect();
file.write_all(&plaintext).expect("write");
drop(file);
let put = gw(
socket,
&format!(
"blob put --file {} --media-type text/plain",
payload.display()
),
);
(
health, status, events, tasks, missing, put, plaintext, payload,
)
})
.await
.expect("join");
let (health, status, events, tasks, missing, put, plaintext, payload) = answers;
assert_eq!(
code(&health),
0,
"{:?}",
String::from_utf8_lossy(&health.stderr)
);
assert_eq!(
json(&health),
serde_json::json!({"type": "health", "ready": true, "sealed": true})
);
assert_eq!(code(&status), 0);
let status = json(&status);
assert_eq!(status["type"], "status");
assert_eq!(status["public_revision"], "a1b2c3d4e5".repeat(4));
assert_eq!(code(&events), 0);
let events = json(&events);
assert_eq!(
events["events"].as_array().expect("events").len(),
1,
"{events}"
);
assert_eq!(code(&tasks), 0);
assert_eq!(json(&tasks)["records"], serde_json::json!([]));
assert_eq!(code(&missing), 4);
assert_eq!(json(&missing)["code"], "not_found");
assert_eq!(code(&put), 0, "{:?}", String::from_utf8_lossy(&put.stderr));
let put = json(&put);
assert_eq!(put["type"], "blob_committed");
assert_eq!(put["deduplicated"], false);
let address = put["descriptor"]["address"]
.as_str()
.expect("an address")
.to_owned();
let back = payload.with_file_name("back.bin");
let (stat, got) = tokio::task::spawn_blocking({
let socket = socket.clone();
let back = back.clone();
let address = address.clone();
move || {
let stat = gw(&socket, &format!("blob stat {address}"));
let got = gw(
&socket,
&format!("blob get {address} --output {}", back.display()),
);
(stat, got)
}
})
.await
.expect("join");
assert_eq!(code(&stat), 0);
assert_eq!(json(&stat)["descriptor"]["media_type"], "text/plain");
assert_eq!(code(&got), 0, "{:?}", String::from_utf8_lossy(&got.stderr));
assert_eq!(
std::fs::read(&back).expect("read back"),
plaintext,
"the blob did not survive the round trip"
);
let _ = stop.send(());
serving.await.expect("join");
let _ = std::fs::remove_dir_all(&dir);
let _ = sqlx::raw_sql(sqlx::AssertSqlSafe(format!(
"DROP DATABASE IF EXISTS {name} WITH (FORCE);"
)))
.execute(&maintenance)
.await;
}
const DAEMON_ROLE: &str = "gwk_cli_daemon";
const DAEMON_PASSWORD: &str = "ci";
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn the_admin_door_initializes_a_database_and_the_daemon_serves_it() {
use base64::prelude::{BASE64_STANDARD, Engine as _};
use sqlx::PgPool;
let admin_url = std::env::var("GWK_TEST_ADMIN_DATABASE_URL")
.expect("GWK_TEST_ADMIN_DATABASE_URL must point at a PostgreSQL superuser DSN");
let maintenance = PgPool::connect(&admin_url).await.expect("maintenance");
for statement in [
format!("CREATE ROLE {DAEMON_ROLE} LOGIN PASSWORD '{DAEMON_PASSWORD}';"),
format!(
"ALTER ROLE {DAEMON_ROLE} LOGIN PASSWORD '{DAEMON_PASSWORD}' NOSUPERUSER NOCREATEDB NOCREATEROLE NOBYPASSRLS;"
),
] {
let _ = sqlx::raw_sql(sqlx::AssertSqlSafe(statement))
.execute(&maintenance)
.await;
}
let live = format!("gwk_daemon_{}", std::process::id());
let scratch = format!("gwk_scratch_{}", std::process::id());
for name in [&live, &scratch] {
let _ = sqlx::raw_sql(sqlx::AssertSqlSafe(format!(
"DROP DATABASE IF EXISTS {name} WITH (FORCE);"
)))
.execute(&maintenance)
.await;
sqlx::raw_sql(sqlx::AssertSqlSafe(format!("CREATE DATABASE {name};")))
.execute(&maintenance)
.await
.expect("create database");
}
let (prefix, _) = admin_url.rsplit_once('/').expect("a /database suffix");
let admin_dsn = format!("{prefix}/{live}");
let runtime_dsn = {
let (scheme, rest) = prefix.split_once("://").expect("a scheme");
let host = rest.rsplit_once('@').map_or(rest, |(_, host)| host);
format!("{scheme}://{DAEMON_ROLE}:{DAEMON_PASSWORD}@{host}/{live}")
};
let dir = private_dir("daemon");
let socket = dir.join("gwk.sock");
let blob_root = dir.join("blobs");
let revision = "a1b2c3d4e5".repeat(4);
let kek = BASE64_STANDARD.encode([0x5a_u8; 32]);
let blob_root_text = blob_root.to_string_lossy().into_owned();
let admin_env: Vec<(&str, &str)> = vec![
("GWK_ADMIN_DATABASE_URL", &admin_dsn),
("GWK_RUNTIME_ROLE", DAEMON_ROLE),
("GWK_PUBLIC_REVISION", &revision),
("GWK_BLOB_ROOT", &blob_root_text),
("GWK_BLOB_KEK", &kek),
("GWK_BLOB_KEK_ID", "kek-test"),
];
let stamp = json(&gw_env("build-info", &admin_env))["public_revision"]
.as_str()
.map(str::to_owned);
let reported = stamp.clone().unwrap_or_else(|| revision.clone());
assert!(
reported.len() == 40 && reported.bytes().all(|b| b.is_ascii_hexdigit()),
"a reported revision is a full hex revision, got {reported:?}"
);
if let Some(stamped) = &stamp {
assert_ne!(
stamped, &revision,
"this case can only prove the override rule if the two differ"
);
}
let init = gw_env("admin init", &admin_env);
assert_eq!(code(&init), 0, "{}", String::from_utf8_lossy(&init.stdout));
let answer = json(&init);
assert_eq!(answer["type"], "admin_initialized");
assert_eq!(answer["outcome"], "initialized");
assert_eq!(answer["public_revision"], reported);
let again = gw_env("admin init", &admin_env);
assert_eq!(
code(&again),
0,
"{}",
String::from_utf8_lossy(&again.stdout)
);
assert_eq!(json(&again)["outcome"], "already_initialized");
let verified = gw_env("admin verify", &admin_env);
assert_eq!(
code(&verified),
0,
"{}",
String::from_utf8_lossy(&verified.stdout)
);
let answer = json(&verified);
assert_eq!(answer["target"], "initialized");
assert_eq!(answer["runtime_role_exists"], true);
assert_eq!(answer["violations"], serde_json::json!([]));
assert_eq!(
answer["detail"]["contract_sha256"],
answer["expected_contract_sha256"]
);
let rebuilt = gw_env(
&format!("admin rebuild-projections --scratch-database {scratch}"),
&admin_env,
);
assert_eq!(
code(&rebuilt),
0,
"{}",
String::from_utf8_lossy(&rebuilt.stdout)
);
let answer = json(&rebuilt);
assert_eq!(answer["agrees"], true, "{answer}");
assert_eq!(answer["live_hash"], answer["rebuilt_hash"]);
let mut serving = Command::new(env!("CARGO_BIN_EXE_gw"))
.arg("daemon")
.env("GWK_DATABASE_URL", &runtime_dsn)
.env("GWK_SOCKET_PATH", &socket)
.env("GWK_BLOB_ROOT", &blob_root)
.env("GWK_BLOB_KEK", &kek)
.env("GWK_BLOB_KEK_ID", "kek-test")
.env("GWK_PUBLIC_REVISION", &revision)
.spawn()
.expect("spawn the daemon");
let mut health = None;
for _ in 0..100 {
let out = gw(&socket, "kernel health");
if code(&out) == 0 {
health = Some(out);
break;
}
assert!(
serving.try_wait().expect("wait").is_none(),
"the daemon exited before it served: {}",
String::from_utf8_lossy(&out.stdout)
);
std::thread::sleep(std::time::Duration::from_millis(100));
}
let health = health.expect("the daemon never answered health");
assert_eq!(
json(&health),
serde_json::json!({"type": "health", "ready": true, "sealed": true})
);
let status = gw(&socket, "kernel status");
assert_eq!(code(&status), 0);
let answer = json(&status);
assert_eq!(answer["public_revision"], reported);
assert_eq!(answer["sealed"], true);
let payload = dir.join("hold.bin");
std::fs::write(&payload, b"a blob nothing in the log points at").expect("write a payload");
let put = gw(
&socket,
&format!(
"blob put --file {} --media-type text/plain",
payload.display()
),
);
assert_eq!(code(&put), 0, "{}", String::from_utf8_lossy(&put.stdout));
let held = json(&put)["descriptor"]["address"]
.as_str()
.expect("an address")
.to_owned();
let pinned = gw_env(&format!("admin blob pin {held} evidence-1"), &admin_env);
assert_eq!(
code(&pinned),
0,
"{}",
String::from_utf8_lossy(&pinned.stdout)
);
assert_eq!(json(&pinned)["type"], "blob_pinned");
let swept = gw_env("admin blob sweep", &admin_env);
assert_eq!(
code(&swept),
0,
"{}",
String::from_utf8_lossy(&swept.stdout)
);
let removed = json(&swept);
assert!(
!removed["removed"]
.as_array()
.expect("an array")
.iter()
.any(|address| address == &serde_json::Value::String(held.clone())),
"a pinned blob was swept: {removed}"
);
assert_eq!(code(&gw(&socket, &format!("blob stat {held}"))), 0);
let unpinned = gw_env(&format!("admin blob unpin {held} evidence-1"), &admin_env);
assert_eq!(
code(&unpinned),
0,
"{}",
String::from_utf8_lossy(&unpinned.stdout)
);
let swept = gw_env("admin blob sweep", &admin_env);
assert_eq!(code(&swept), 0);
assert!(
json(&swept)["removed"]
.as_array()
.expect("an array")
.iter()
.any(|address| address == &serde_json::Value::String(held.clone())),
"the unpinned blob survived a sweep: {}",
String::from_utf8_lossy(&swept.stdout)
);
let gone = gw(&socket, &format!("blob stat {held}"));
assert_eq!(code(&gone), 4, "{}", String::from_utf8_lossy(&gone.stdout));
assert_eq!(json(&gone)["code"], "not_found");
let doomed = dir.join("shred.bin");
std::fs::write(&doomed, b"a blob that will be cryptographically erased").expect("write");
let put = gw(
&socket,
&format!(
"blob put --file {} --media-type text/plain",
doomed.display()
),
);
assert_eq!(code(&put), 0);
let erased = json(&put)["descriptor"]["address"]
.as_str()
.expect("an address")
.to_owned();
let shredded = gw_env(&format!("admin blob shred {erased}"), &admin_env);
assert_eq!(
code(&shredded),
0,
"{}",
String::from_utf8_lossy(&shredded.stdout)
);
assert_eq!(json(&shredded)["type"], "blob_shredded");
let refused = gw(
&socket,
&format!(
"blob get {erased} --output {}",
dir.join("never.bin").display()
),
);
assert_eq!(
code(&refused),
6,
"{}",
String::from_utf8_lossy(&refused.stdout)
);
assert_eq!(json(&refused)["code"], "blob_tombstoned");
assert!(
!dir.join("never.bin").exists(),
"a refused read still wrote a file"
);
let second = gw_env(
"daemon",
&[
("GWK_DATABASE_URL", &runtime_dsn),
(
"GWK_SOCKET_PATH",
&dir.join("second.sock").to_string_lossy(),
),
("GWK_BLOB_ROOT", &blob_root.to_string_lossy()),
("GWK_BLOB_KEK", &kek),
("GWK_BLOB_KEK_ID", "kek-test"),
("GWK_PUBLIC_REVISION", &revision),
],
);
assert_eq!(
code(&second),
3,
"{}",
String::from_utf8_lossy(&second.stdout)
);
assert_eq!(json(&second)["code"], "fenced");
Command::new("kill")
.args(["-KILL", &serving.id().to_string()])
.status()
.expect("signal the daemon");
let died = serving.wait().expect("wait");
assert!(!died.success(), "SIGKILL is not a clean exit: {died:?}");
assert!(
socket.exists(),
"a killed daemon cannot have cleaned up after itself"
);
let mut released = false;
for _ in 0..100 {
let live_backends: i64 = sqlx::query_scalar(
"SELECT count(*) FROM pg_stat_activity WHERE datname = $1 AND usename = $2",
)
.bind(&live)
.bind(DAEMON_ROLE)
.fetch_one(&maintenance)
.await
.expect("count the daemon's backends");
if live_backends == 0 {
released = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
assert!(released, "the killed daemon's backends outlived it");
let mut restarted = Command::new(env!("CARGO_BIN_EXE_gw"))
.arg("daemon")
.env("GWK_DATABASE_URL", &runtime_dsn)
.env("GWK_SOCKET_PATH", &socket)
.env("GWK_BLOB_ROOT", &blob_root)
.env("GWK_BLOB_KEK", &kek)
.env("GWK_BLOB_KEK_ID", "kek-test")
.env("GWK_PUBLIC_REVISION", &revision)
.spawn()
.expect("spawn the replacement");
let mut recovered = None;
for _ in 0..100 {
let out = gw(&socket, "kernel health");
if code(&out) == 0 {
recovered = Some(out);
break;
}
if let Some(exited) = restarted.try_wait().expect("wait") {
panic!("the replacement refused to start after a crash: {exited:?}");
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
let recovered = recovered.expect("the replacement never answered health");
assert_eq!(json(&recovered)["ready"], true);
let status = gw(&socket, "kernel status");
assert_eq!(code(&status), 0);
assert_eq!(json(&status)["public_revision"], reported);
Command::new("kill")
.args(["-TERM", &restarted.id().to_string()])
.status()
.expect("signal the daemon");
let stopped = restarted.wait().expect("wait");
assert!(stopped.success(), "the daemon exited {stopped:?}");
assert!(!socket.exists(), "shutdown left the socket behind");
let _ = std::fs::remove_dir_all(&dir);
for name in [&live, &scratch] {
let _ = sqlx::raw_sql(sqlx::AssertSqlSafe(format!(
"DROP DATABASE IF EXISTS {name} WITH (FORCE);"
)))
.execute(&maintenance)
.await;
}
}