#![allow(clippy::unwrap_used)]
mod common;
use std::collections::BTreeSet;
use std::path::{Path, PathBuf};
use std::process::Output;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, mpsc};
use std::time::{Duration, Instant};
use connectrpc::server::Server;
use connectrpc::{ConnectError, RequestContext, Response, Router, handler_fn};
use mkit_attest::grant::{AcceptedSchemes, EpochStatement, OwnerScheme, SignedHeader};
use mkit_core::sign::{KeyPair, save_key};
use mkit_server::auth_v2::AuthV2Config;
use mkit_server::pipeline::{AuthMode, Hooks, Pipeline, PipelineConfig};
use mkit_server::policy::NamespacePolicy;
use mkit_server::upload::UploadLimits;
use mkit_server::{Addressing, MemoryBlobStore, MemoryKv, MultiAddressing, SystemClock};
use mkit_transport_connect::generated;
struct Party {
root: tempfile::TempDir,
repo: PathBuf,
xdg: PathBuf,
public_key: String,
}
impl Party {
fn new(seed: u8) -> Self {
let key_home = mkit_cli::config::home_dir_for_euid().unwrap();
let root = tempfile::tempdir_in(key_home).unwrap();
let repo = root.path().join("repo");
let xdg = root.path().join("xdg");
std::fs::create_dir_all(&repo).unwrap();
std::fs::create_dir_all(&xdg).unwrap();
let party = Self {
root,
repo,
xdg,
public_key: mkit_core::hash::to_hex_bytes(&KeyPair::from_seed([seed; 32]).public.0),
};
assert!(party.run(&["init"]).status.success());
let keys = party.repo.join(".mkit").join("keys");
std::fs::create_dir_all(&keys).unwrap();
save_key(&keys.join("default.key"), &KeyPair::from_seed([seed; 32])).unwrap();
party
}
fn run(&self, args: &[&str]) -> Output {
common::mkit(&self.repo, &self.xdg, args)
}
fn ok(&self, args: &[&str]) -> String {
let out = self.run(args);
assert!(
out.status.success(),
"mkit {args:?} failed: {}",
String::from_utf8_lossy(&out.stderr)
);
String::from_utf8(out.stdout).unwrap()
}
fn namespace(&self) -> String {
format!("ed25519-{}", self.public_key)
}
fn commit(&self, file: &str, text: &str, message: &str) {
std::fs::write(self.repo.join(file), text).unwrap();
self.ok(&["add", file]);
self.ok(&["commit", "-m", message]);
}
#[cfg(unix)]
fn commit_bytes(&self, file: &str, data: &[u8], message: &str) {
std::fs::write(self.repo.join(file), data).unwrap();
self.ok(&["add", file]);
self.ok(&["commit", "-m", message]);
}
fn add_remote(&self, url: &str) {
let path = self.repo.join(".mkit").join("config");
let mut config = std::fs::read_to_string(&path).unwrap_or_default();
config.push_str("\nremote.origin.url = ");
config.push_str(url);
config.push_str("\nremote.origin.type = http\n");
std::fs::write(path, config).unwrap();
}
fn connect_to(&self, url: &str) {
self.ok(&["config", "transport_auth", "envelope"]);
self.ok(&["config", "trusted_remote_endpoint", url]);
self.add_remote(url);
}
}
fn stderr(out: &Output) -> String {
String::from_utf8_lossy(&out.stderr).into_owned()
}
struct Live {
origin: String,
shutdown: Option<tokio::sync::oneshot::Sender<()>>,
thread: Option<std::thread::JoinHandle<()>>,
}
impl Live {
fn start(owner_namespace: &str) -> Self {
Self::start_capped(owner_namespace, 64 << 20)
}
fn start_capped(owner_namespace: &str, max_pack_bytes: u64) -> Self {
let owner = mkit_attest::grant::Namespace::parse(owner_namespace).unwrap();
let (addr_tx, addr_rx) = mpsc::channel();
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
let thread = std::thread::spawn(move || {
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(2)
.enable_all()
.build()
.unwrap();
rt.block_on(async move {
let bound = Server::bind("127.0.0.1:0").await.unwrap();
let origin = format!("http://{}", bound.local_addr().unwrap());
addr_tx.send(origin.clone()).unwrap();
let mut cfg = PipelineConfig::new(
Addressing::Multi(MultiAddressing::new().with_namespace_policy(
NamespacePolicy::Allowlist(BTreeSet::from([owner])),
)),
AuthMode::AuthV2(AuthV2Config::new(&origin, "").unwrap()),
UploadLimits {
max_total_bytes: max_pack_bytes,
max_chunks: 64,
},
);
cfg.grants = Some(
mkit_server::GrantConfig::new_allowing_loopback(
&origin,
AcceptedSchemes::of(&[OwnerScheme::Ed25519, OwnerScheme::Secp256k1Eip191]),
vec![],
)
.unwrap(),
);
cfg.ticket_keys = Some(
mkit_server::upload::token::TicketKeys::parse(
"dev 1111111111111111111111111111111111111111111111111111111111111111",
)
.unwrap(),
);
let clock = Arc::new(SystemClock);
let pipeline = Pipeline::new(
MemoryBlobStore::default(),
MemoryKv::with_clock(clock.clone()),
Hooks::new(),
cfg,
clock,
Arc::new(mkit_server::NoopMetrics),
)
.unwrap();
bound
.serve_with_service_and_shutdown(
mkit_server::connect::service(Arc::new(pipeline)),
async {
let _ = shutdown_rx.await;
},
)
.await
.unwrap();
});
});
Self {
origin: addr_rx.recv().unwrap(),
shutdown: Some(shutdown_tx),
thread: Some(thread),
}
}
}
impl Drop for Live {
fn drop(&mut self) {
if let Some(tx) = self.shutdown.take() {
let _ = tx.send(());
}
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
fn grant_file(party: &Party, name: &str, header: &str) -> String {
let path: &Path = &party.repo.join(name);
std::fs::write(path, header).unwrap();
path.to_str().unwrap().to_owned()
}
#[test]
fn grantee_push_revocation_and_reissue_against_a_real_server() {
let owner = Party::new(0x11);
let grantee = Party::new(0x22);
let live = Live::start(&owner.namespace());
let url = format!(
"mkit+http://{}/{}/site",
live.origin.trim_start_matches("http://"),
owner.namespace()
);
owner.connect_to(&url);
grantee.connect_to(&url);
owner.commit("a.txt", "owner\n", "first");
owner.ok(&["push", "origin"]);
assert_eq!(
owner
.ok(&["epoch", "show", "origin"])
.lines()
.last()
.unwrap(),
"epoch 0"
);
grantee.commit("base.txt", "grantee base\n", "base");
grantee.ok(&["switch", "-c", "feature"]);
grantee.commit("b.txt", "grantee 1\n", "grantee one");
let denied = grantee.run(&["push", "origin"]);
assert!(!denied.status.success(), "a push without a grant must fail");
assert!(
stderr(&denied).to_lowercase().contains("denied")
|| stderr(&denied).to_lowercase().contains("permission"),
"{}",
stderr(&denied)
);
let header = owner.ok(&[
"grant",
"create",
"--cap",
"write,read",
"--grantee",
&grantee.public_key,
"--repo",
"site",
"--refs",
"refs/heads/*=cuf",
"--ttl",
"1h",
]);
let signed = SignedHeader::parse(header.trim()).unwrap();
let statement = mkit_attest::grant::Grant::parse(&signed.statement).unwrap();
assert_eq!(statement.epoch, 0);
assert_eq!(statement.audiences, std::slice::from_ref(&live.origin));
assert_eq!(statement.capabilities.token(), "read,write");
let out = grantee.ok(&[
"grant",
"add",
&grant_file(&grantee, "g0.txt", header.trim()),
]);
assert!(out.starts_with("added grant "), "{out}");
let pushed = grantee.run(&["push", "origin"]);
assert!(pushed.status.success(), "{}", stderr(&pushed));
owner.ok(&["fetch", "origin"]);
let refs = owner.ok(&["show-ref"]);
assert!(refs.contains("refs/remotes/origin/feature"), "{refs}");
let bumped = owner.ok(&["epoch", "bump", "origin"]);
assert!(bumped.contains("epoch 1 (was 0)"), "{bumped}");
assert_eq!(
owner
.ok(&["epoch", "show", "origin"])
.lines()
.last()
.unwrap(),
"epoch 1"
);
grantee.commit("b.txt", "grantee 2\n", "grantee two");
let revoked = grantee.run(&["push", "origin"]);
assert!(!revoked.status.success(), "a revoked grant must not work");
let text = stderr(&revoked).to_lowercase();
assert!(
text.contains("denied") || text.contains("permission"),
"{text}"
);
let header = owner.ok(&[
"grant",
"create",
"--cap",
"write",
"--grantee",
&grantee.public_key,
"--repo",
"site",
"--refs",
"refs/heads/*=cuf",
"--ttl",
"1h",
"--store",
]);
let fresh =
mkit_attest::grant::Grant::parse(&SignedHeader::parse(header.trim()).unwrap().statement)
.unwrap();
assert_eq!(fresh.epoch, 1);
grantee.ok(&[
"grant",
"add",
&grant_file(&grantee, "g1.txt", header.trim()),
]);
let listed = grantee.ok(&["grant", "list", "--check"]);
assert!(listed.contains("stale epoch"), "{listed}");
assert!(listed.contains("epoch current"), "{listed}");
let pushed = grantee.run(&["push", "origin"]);
assert!(pushed.status.success(), "{}", stderr(&pushed));
let too_far = owner.run(&["epoch", "bump", "origin", "--by", "1025"]);
assert_eq!(too_far.status.code(), Some(64));
assert!(stderr(&too_far).contains("1024"), "{}", stderr(&too_far));
assert!(
owner
.ok(&["epoch", "show", "origin"])
.ends_with("epoch 1\n")
);
let revoke = owner.run(&["grant", "revoke", "origin", "--prune"]);
assert!(revoke.status.success(), "{}", stderr(&revoke));
let text = stderr(&revoke);
assert!(text.contains("revoking: 1 local grant(s)"), "{text}");
assert!(text.contains("removed 1 local grant(s)"), "{text}");
assert!(String::from_utf8_lossy(&revoke.stdout).contains("epoch 2 (was 1)"));
assert!(owner.ok(&["grant", "list"]).starts_with("no grants in "));
let listed = grantee.ok(&["grant", "list", "--check"]);
assert_eq!(listed.matches("stale epoch").count(), 2, "{listed}");
}
type Visibility = (http::HeaderMap, generated::SetRepoVisibilityRequest);
struct Stub {
port: u16,
statements: Arc<Mutex<Vec<String>>>,
visibility: Arc<Mutex<Vec<Visibility>>>,
shutdown: Option<tokio::sync::oneshot::Sender<()>>,
thread: Option<std::thread::JoinHandle<()>>,
}
impl Stub {
fn start(pending: usize) -> Self {
let statements = Arc::new(Mutex::new(Vec::new()));
let seen = statements.clone();
let visibility: Arc<Mutex<Vec<Visibility>>> = Arc::default();
let vis = visibility.clone();
let left = Arc::new(AtomicUsize::new(pending));
let (addr_tx, addr_rx) = mpsc::channel();
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
let thread = std::thread::spawn(move || {
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(1)
.enable_all()
.build()
.unwrap();
rt.block_on(async move {
let bound = Server::bind("127.0.0.1:0").await.unwrap();
addr_tx.send(bound.local_addr().unwrap().port()).unwrap();
let service = "mkit.transport.v1.TransportService";
let router = Router::new()
.route(
service,
"GetGrantEpoch",
handler_fn(
|_: RequestContext, _: generated::GetGrantEpochRequest| async {
Ok::<_, ConnectError>(Response::new(
generated::GetGrantEpochResponse {
epoch: Some(4),
..Default::default()
},
))
},
),
)
.route(
service,
"SetGrantEpoch",
handler_fn(
move |_: RequestContext, req: generated::SetGrantEpochRequest| {
seen.lock()
.unwrap()
.push(req.signed_statement.clone().unwrap_or_default());
let pending = left
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |n| {
n.checked_sub(1)
})
.is_ok();
async move {
if pending {
let mut headers = http::HeaderMap::new();
headers.insert("retry-after", "1".parse().unwrap());
return Err(ConnectError::unavailable(
"revocation in progress",
)
.with_headers(headers));
}
Ok(Response::new(generated::SetGrantEpochResponse {
epoch: Some(5),
..Default::default()
}))
}
},
),
)
.route(
service,
"SetRepoVisibility",
handler_fn(
move |ctx: RequestContext, req: generated::SetRepoVisibilityRequest| {
vis.lock().unwrap().push((ctx.headers().clone(), req));
async {
Ok::<_, ConnectError>(Response::new(
generated::SetRepoVisibilityResponse::default(),
))
}
},
),
);
bound
.serve_with_graceful_shutdown(router, async {
let _ = shutdown_rx.await;
})
.await
.unwrap();
});
});
Self {
port: addr_rx.recv().unwrap(),
statements,
visibility,
shutdown: Some(shutdown_tx),
thread: Some(thread),
}
}
}
impl Drop for Stub {
fn drop(&mut self) {
if let Some(tx) = self.shutdown.take() {
let _ = tx.send(());
}
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
#[test]
fn a_pending_bump_waits_retry_after_and_resends_the_identical_statement() {
let owner = Party::new(0x11);
let stub = Stub::start(2);
let url = format!(
"mkit+http://127.0.0.1:{}/{}/site",
stub.port,
owner.namespace()
);
owner.add_remote(&url);
let started = Instant::now();
let out = owner.run(&["epoch", "bump", "origin"]);
assert!(out.status.success(), "{}", stderr(&out));
assert!(String::from_utf8_lossy(&out.stdout).contains("epoch 5 (was 4)"));
assert!(
started.elapsed() >= Duration::from_secs(2),
"two Retry-After: 1 waits, took {:?}",
started.elapsed()
);
let sent = stub.statements.lock().unwrap().clone();
assert_eq!(sent.len(), 3);
assert!(
sent.iter().all(|s| s == &sent[0]),
"every retry is byte-identical"
);
let header = SignedHeader::parse(&sent[0]).unwrap();
let statement = EpochStatement::parse(&header.statement).unwrap();
assert_eq!(statement.new_epoch, 5);
assert_eq!(
statement.audiences,
[format!("http://127.0.0.1:{}", stub.port)]
);
}
#[test]
fn a_bump_that_stays_pending_stops_at_the_timeout_and_says_it_may_still_finish() {
let owner = Party::new(0x11);
let stub = Stub::start(usize::MAX);
let url = format!(
"mkit+http://127.0.0.1:{}/{}/site",
stub.port,
owner.namespace()
);
owner.add_remote(&url);
let out = owner.run(&["epoch", "bump", "origin", "--timeout", "2s"]);
assert_eq!(out.status.code(), Some(75), "{}", stderr(&out));
let text = stderr(&out);
assert!(text.contains("still completing"), "{text}");
assert!(text.contains("may yet finish"), "{text}");
assert_eq!(stub.statements.lock().unwrap().len(), 2);
}
#[test]
fn a_loopback_audience_needs_a_loopback_dev_remote() {
let owner = Party::new(0x11);
let out = owner.run(&[
"grant",
"create",
"--cap",
"read",
"--all",
"--grantee",
&owner.public_key,
"--audience",
"http://127.0.0.1:8080",
"--offline",
]);
assert_eq!(out.status.code(), Some(64));
assert!(stderr(&out).contains("loopback"), "{}", stderr(&out));
}
#[test]
fn visibility_set_sends_a_signed_request_or_an_unsigned_owner_statement() {
use generated::set_repo_visibility_request::Mode;
let owner = Party::new(0x11);
let stub = Stub::start(0);
let url = format!(
"mkit+http://127.0.0.1:{}/{}/site",
stub.port,
owner.namespace()
);
let repository = format!("{}/site", owner.namespace());
owner.add_remote(&url);
let out = owner.run(&["visibility", "set", "origin", "private"]);
assert_eq!(out.status.code(), Some(77), "{}", stderr(&out));
assert!(
stderr(&out).contains("transport_auth envelope"),
"{}",
stderr(&out)
);
assert!(stub.visibility.lock().unwrap().is_empty());
let out = owner.run(&[
"visibility",
"set",
"origin",
"private",
"--print-statement",
]);
assert_eq!(out.status.code(), Some(64));
assert!(stderr(&out).contains("add --statement"), "{}", stderr(&out));
owner.ok(&["config", "transport_auth", "envelope"]);
owner.ok(&["config", "trusted_remote_endpoint", &url]);
let out = owner.ok(&["visibility", "set", "origin", "private"]);
assert_eq!(out.trim(), format!("{repository} is now private"));
{
let seen = stub.visibility.lock().unwrap();
let (headers, request) = &seen[0];
assert!(headers.contains_key("x-signature") && headers.contains_key("x-public-key"));
assert_eq!(
headers.get("x-repository").unwrap().to_str().unwrap(),
repository
);
assert!(!headers.contains_key("x-write-grant"));
assert!(matches!(
&request.mode,
Some(Mode::Visibility(v)) if v.as_known() == Some(generated::RepoVisibility::REPO_VISIBILITY_PRIVATE)
));
}
let out = owner.ok(&["visibility", "set", "origin", "public", "--statement"]);
assert_eq!(out.trim(), format!("{repository} is now public"));
let seen = stub.visibility.lock().unwrap();
let (headers, request) = &seen[1];
for name in [
"x-signature",
"x-public-key",
"x-digest",
"x-envelope-version",
] {
assert!(!headers.contains_key(name), "{name}");
}
assert_eq!(
headers.get("x-repository").unwrap().to_str().unwrap(),
repository
);
let Some(Mode::SignedStatement(header)) = &request.mode else {
panic!("statement mode expected: {request:?}")
};
let audience = format!("http://127.0.0.1:{}", stub.port);
let cfg = mkit_attest::grant::VerifierConfig::new_allowing_loopback(
&audience,
AcceptedSchemes::of(&[OwnerScheme::Ed25519]),
vec![],
)
.unwrap();
let verified = mkit_attest::grant::verify_visibility_statement(
&cfg,
header,
&mkit_core::repo_identity::RepositoryIdentity::parse(&repository).unwrap(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| i64::try_from(d.as_millis()).unwrap())
.unwrap(),
)
.unwrap();
assert_eq!(
verified.statement().visibility,
mkit_attest::grant::Visibility::Public
);
}
#[test]
fn epoch_show_and_grant_list_check_pin_their_output() {
let owner = Party::new(0x11);
let stub = Stub::start(0);
let url = format!(
"mkit+http://127.0.0.1:{}/{}/site",
stub.port,
owner.namespace()
);
owner.add_remote(&url);
owner.ok(&["config", "trusted_remote_endpoint", &url]);
let filters = vec![
(r"[0-9a-f]{64}", "<hex>"),
(r"127\.0\.0\.1:\d+", "127.0.0.1:<port>"),
(r"\d{4}-\d\d-\d\d \d\d:\d\d:\d\d \+0000", "<date>"),
(r"\d{13}", "<ms>"),
];
for epoch in ["3", "4", "5"] {
owner.ok(&[
"grant",
"create",
"--cap",
"read",
"--all",
"--grantee",
&owner.public_key,
"--audience",
&format!("http://127.0.0.1:{}", stub.port),
"--epoch",
epoch,
"--store",
]);
}
insta::with_settings!({filters => filters.clone()}, {
insta::assert_snapshot!("epoch_show_human", owner.ok(&["epoch", "show", "origin"]));
insta::assert_snapshot!("epoch_show_json", owner.ok(&["epoch", "show", "origin", "--json"]));
insta::assert_snapshot!(
"grant_list_check_human",
owner.ok(&["grant", "list", "--check", "--remote", "origin"])
);
insta::assert_snapshot!(
"grant_list_check_json",
owner.ok(&["grant", "list", "--check", "--remote", "origin", "--json"])
);
insta::assert_snapshot!("epoch_bump_human", owner.ok(&["epoch", "bump", "origin"]));
insta::assert_snapshot!("epoch_bump_json", owner.ok(&["epoch", "bump", "origin", "--json"]));
});
}
#[test]
fn a_grant_for_another_audience_is_unchecked_and_survives_prune() {
let owner = Party::new(0x11);
let stub = Stub::start(0);
let url = format!(
"mkit+http://127.0.0.1:{}/{}/site",
stub.port,
owner.namespace()
);
owner.add_remote(&url);
owner.ok(&["config", "trusted_remote_endpoint", &url]);
let here = format!("http://127.0.0.1:{}", stub.port);
let create = |audiences: &[&str], epoch: &str| {
let mut args = vec![
"grant",
"create",
"--cap",
"read",
"--all",
"--grantee",
&owner.public_key,
"--epoch",
epoch,
"--store",
];
for a in audiences {
args.extend(["--audience", a]);
}
owner.ok(&args);
};
create(&["https://other.example"], "4");
create(&[&here, "https://other.example"], "4");
create(&[&here], "4");
let listed = owner.ok(&["grant", "list", "--check", "--remote", "origin"]);
assert!(
listed.contains("unchecked (audience is not the checked remote)"),
"{listed}"
);
assert!(listed.contains("epoch current"), "{listed}");
let out = owner.run(&["grant", "revoke", "origin", "--prune"]);
assert!(out.status.success(), "{}", stderr(&out));
let text = stderr(&out);
assert!(
text.contains("still valid at https://other.example"),
"{text}"
);
assert!(text.contains("removed 1 local grant(s)"), "{text}");
let left = owner.ok(&["grant", "list"]);
assert_eq!(left.matches("epoch 4").count(), 2, "{left}");
}
#[test]
fn add_warns_when_the_grant_is_for_a_future_epoch() {
let owner = Party::new(0x11);
let grantee = Party::new(0x22);
let stub = Stub::start(0);
let url = format!(
"mkit+http://127.0.0.1:{}/{}/site",
stub.port,
owner.namespace()
);
let audience = format!("http://127.0.0.1:{}", stub.port);
owner.add_remote(&url);
grantee.add_remote(&url);
let make = |epoch: &str| {
owner.ok(&[
"grant",
"create",
"--cap",
"read",
"--all",
"--grantee",
&grantee.public_key,
"--audience",
&audience,
"--epoch",
epoch,
"--remote",
"origin",
])
};
let future = grant_file(&grantee, "future.txt", make("5").trim());
let out = grantee.run(&["grant", "add", "--remote", "origin", &future]);
assert!(out.status.success(), "{}", stderr(&out));
assert!(
stderr(&out).contains("is for epoch 5 but")
&& stderr(&out).contains("is at epoch 4")
&& stderr(&out).contains("outranks your other grants"),
"{}",
stderr(&out)
);
let current = grant_file(&grantee, "current.txt", make("4").trim());
let out = grantee.run(&["grant", "add", "--remote", "origin", ¤t]);
assert!(
out.status.success() && !stderr(&out).contains("warning"),
"{}",
stderr(&out)
);
let other = grant_file(&grantee, "other.txt", make("6").trim());
let out = grantee.run(&["grant", "add", "--offline", &other]);
assert!(
out.status.success() && !stderr(&out).contains("warning"),
"{}",
stderr(&out)
);
}
#[cfg(unix)]
#[test]
fn ctrl_c_cancels_a_pending_bump_and_says_it_may_still_complete() {
use std::process::{Command, Stdio};
let owner = Party::new(0x11);
let stub = Stub::start(usize::MAX);
let url = format!(
"mkit+http://127.0.0.1:{}/{}/site",
stub.port,
owner.namespace()
);
owner.add_remote(&url);
let child = Command::new(env!("CARGO_BIN_EXE_mkit"))
.args(["epoch", "bump", "origin"])
.current_dir(&owner.repo)
.env("XDG_CONFIG_HOME", &owner.xdg)
.env("HOME", &owner.xdg)
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
let deadline = Instant::now() + Duration::from_secs(20);
while stub.statements.lock().unwrap().is_empty() {
assert!(
Instant::now() < deadline,
"the bump never reached the server"
);
std::thread::sleep(Duration::from_millis(20));
}
let killed = Command::new("kill")
.args(["-INT", &child.id().to_string()])
.status()
.unwrap();
assert!(killed.success());
let out = child.wait_with_output().unwrap();
assert_eq!(out.status.code(), Some(75), "{}", stderr(&out));
let text = stderr(&out);
assert!(text.contains("interrupted"), "{text}");
assert!(text.contains("may still complete"), "{text}");
assert!(out.stdout.is_empty());
}
#[cfg(unix)]
fn noise(seed: u64, len: usize) -> Vec<u8> {
let mut out = vec![0_u8; len];
let mut state = seed * 2 + 1;
for chunk in out.chunks_mut(8) {
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1_442_695_040_888_963_407);
chunk.copy_from_slice(&state.to_le_bytes()[..chunk.len()]);
}
out
}
#[cfg(unix)]
fn copy_dir(from: &Path, to: &Path) {
std::fs::create_dir_all(to).unwrap();
for entry in std::fs::read_dir(from).unwrap() {
let entry = entry.unwrap();
let target = to.join(entry.file_name());
if entry.file_type().unwrap().is_dir() {
copy_dir(&entry.path(), &target);
} else {
std::fs::copy(entry.path(), target).unwrap();
}
}
}
#[cfg(unix)]
fn copy_of(master: &Party, live: &Live, name: &str) -> Party {
let party = Party::new(0x11);
std::fs::remove_dir_all(party.repo.join(".mkit")).unwrap();
copy_dir(&master.repo.join(".mkit"), &party.repo.join(".mkit"));
{
use std::os::unix::fs::PermissionsExt as _;
let keys = party.repo.join(".mkit").join("keys");
std::fs::set_permissions(&keys, std::fs::Permissions::from_mode(0o700)).unwrap();
for key in std::fs::read_dir(&keys).unwrap() {
let key = key.unwrap().path();
std::fs::set_permissions(key, std::fs::Permissions::from_mode(0o600)).unwrap();
}
}
let url = format!(
"mkit+http://{}/{}/{name}",
live.origin.trim_start_matches("http://"),
party.namespace()
);
party.connect_to(&url);
party
}
#[cfg(unix)]
#[test]
fn a_split_push_reports_its_steps_and_the_published_prefix() {
use std::io::{BufRead as _, BufReader};
use std::process::{Command, Stdio};
let master = Party::new(0x11);
let live = Live::start_capped(&master.namespace(), 8192);
for i in 0..30_u64 {
master.commit_bytes(&format!("f{i:02}.bin"), &noise(i, 6000), &format!("c{i}"));
}
let piped = copy_of(&master, &live, "piped");
let out = piped.run(&["push", "origin"]);
assert!(out.status.success(), "{}", stderr(&out));
let text = stderr(&out);
let steps: Vec<&str> = text
.lines()
.filter(|line| line.starts_with("pushed step "))
.collect();
assert!(steps.len() >= 3, "{text}");
for (index, line) in steps.iter().enumerate() {
let expected = format!("pushed step {}/{}: ", index + 1, steps.len());
assert!(line.starts_with(&expected), "{line} (wanted {expected})");
}
assert!(text.contains("* [new branch]"), "{text}");
let json = copy_of(&master, &live, "json");
let out = json.run(&["push", "--format=json", "origin"]);
assert!(out.status.success(), "{}", stderr(&out));
let reported = String::from_utf8(out.stdout).unwrap();
let steps_field: u64 = reported
.split("\"steps\":")
.nth(1)
.and_then(|rest| rest.trim_end_matches(['}', '\n']).parse().ok())
.unwrap_or_else(|| panic!("no steps in {reported}"));
assert!(steps_field > 1, "{reported}");
assert_eq!(steps_field, steps.len() as u64);
let quiet = copy_of(&master, &live, "quiet");
let out = quiet.run(&["push", "--quiet", "origin"]);
assert!(out.status.success(), "{}", stderr(&out));
assert!(out.stdout.is_empty());
assert!(!stderr(&out).contains("step"), "{}", stderr(&out));
let cut = copy_of(&master, &live, "cut");
let mut child = Command::new(env!("CARGO_BIN_EXE_mkit"))
.args(["push", "origin"])
.current_dir(&cut.repo)
.env("XDG_CONFIG_HOME", &cut.xdg)
.env("HOME", &cut.xdg)
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
let mut lines = BufReader::new(child.stderr.take().unwrap());
let mut seen = String::new();
loop {
let mut next = String::new();
assert!(lines.read_line(&mut next).unwrap() > 0, "{seen}");
seen.push_str(&next);
if next.starts_with("pushed step 1/") {
break;
}
}
assert!(
Command::new("kill")
.args(["-INT", &child.id().to_string()])
.status()
.unwrap()
.success()
);
std::io::Read::read_to_string(&mut lines, &mut seen).unwrap();
let status = child.wait().unwrap();
assert_eq!(status.code(), Some(75), "{seen}");
assert!(seen.contains("advances were published"), "{seen}");
assert!(
seen.contains("the last published advance on branch"),
"{seen}"
);
assert!(seen.contains("re-run the push to resume"), "{seen}");
}
impl Party {
fn use_key_outside_the_repo(&self) {
let key = self.repo.join(".mkit").join("keys").join("default.key");
self.ok(&["config", "signing_key", key.to_str().unwrap()]);
}
}
struct Site {
live: Live,
owner: Party,
url: String,
}
impl Site {
fn new() -> Self {
let owner = Party::new(0x11);
let live = Live::start(&owner.namespace());
let url = format!(
"mkit+http://{}/{}/site",
live.origin.trim_start_matches("http://"),
owner.namespace()
);
owner.connect_to(&url);
owner.use_key_outside_the_repo();
owner.commit("a.txt", "owner\n", "first");
owner.ok(&["push", "origin"]);
Self { live, owner, url }
}
fn party(&self, seed: u8) -> Party {
let party = Party::new(seed);
party.connect_to(&self.url);
party.use_key_outside_the_repo();
party
}
fn clone_as(&self, who: &Party, dir: &str) -> Output {
let target = who.root.path().join(dir);
common::mkit(
who.root.path(),
&who.xdg,
&["clone", &self.url, target.to_str().unwrap()],
)
}
fn clone_anonymously(&self, dir: &str) -> (Party, Output) {
let anon = Party::new(0x77);
let out = self.clone_as(&anon, dir);
(anon, out)
}
fn cloned_file(who: &Party, dir: &str) -> String {
std::fs::read_to_string(who.root.path().join(dir).join("a.txt")).unwrap()
}
}
fn assert_not_found(out: &Output, ns: &str, what: &str) {
assert!(!out.status.success(), "{what}: the clone must fail");
let text = stderr(out).to_lowercase();
let expected = format!("repository `{}/site` not found at", ns.to_lowercase());
assert!(
text.contains(&expected),
"{what}: expected `{expected}`, got: {text}"
);
for hint in ["permission", "denied"] {
assert!(
!text.contains(hint),
"{what}: a private repository must not be told apart from a missing one: {text}"
);
}
}
fn grant(owner: &Party, grantee: &Party, cap: &str, extra: &[&str]) -> String {
let mut args = vec![
"grant",
"create",
"--cap",
cap,
"--grantee",
&grantee.public_key,
"--repo",
"site",
"--ttl",
"1h",
];
assert!(!extra.contains(&"--all"), "--all excludes --repo");
args.extend_from_slice(extra);
owner.ok(&args).trim().to_owned()
}
#[test]
fn visibility_set_makes_a_repository_private_and_public_in_envelope_mode() {
let site = Site::new();
let anon = Party::new(0x77);
let clone = |dir: &str| site.clone_as(&anon, dir);
let ok = clone("public-0");
assert!(ok.status.success(), "{}", stderr(&ok));
let out = site.owner.ok(&["visibility", "set", "origin", "private"]);
assert_eq!(
out.trim(),
format!("{}/site is now private", site.owner.namespace())
);
assert_not_found(
&clone("private-0"),
&site.owner.namespace(),
"anonymous, private",
);
let own = site.clone_as(&site.owner, "owner-0");
assert!(own.status.success(), "{}", stderr(&own));
let out = site.owner.ok(&["visibility", "set", "origin", "public"]);
assert_eq!(
out.trim(),
format!("{}/site is now public", site.owner.namespace())
);
let ok = clone("public-1");
assert!(ok.status.success(), "{}", stderr(&ok));
assert_eq!(Site::cloned_file(&anon, "public-1"), "owner\n");
}
#[test]
fn visibility_set_statement_mode_flips_visibility_with_an_ed25519_owner() {
let site = Site::new();
let anon = Party::new(0x77);
let out = site
.owner
.ok(&["visibility", "set", "origin", "private", "--statement"]);
assert_eq!(
out.trim(),
format!("{}/site is now private", site.owner.namespace())
);
assert_not_found(
&site.clone_as(&anon, "private-0"),
&site.owner.namespace(),
"anonymous, private",
);
std::thread::sleep(Duration::from_millis(5));
let out = site
.owner
.ok(&["visibility", "set", "origin", "public", "--statement"]);
assert_eq!(
out.trim(),
format!("{}/site is now public", site.owner.namespace())
);
let ok = site.clone_as(&anon, "public-0");
assert!(ok.status.success(), "{}", stderr(&ok));
let stranger = site.party(0x33);
let denied = stranger.run(&["visibility", "set", "origin", "private", "--statement"]);
assert!(!denied.status.success(), "{}", stderr(&denied));
let text = stderr(&denied).to_lowercase();
assert!(text.contains("the signing key owns namespace"), "{text}");
}
#[test]
fn a_private_repository_is_cloned_by_the_owner_and_a_read_grantee_only() {
let site = Site::new();
let reader = site.party(0x22);
let writer = site.party(0x33);
site.owner.ok(&["visibility", "set", "origin", "private"]);
let out = site.clone_as(&site.owner, "owner");
assert!(out.status.success(), "{}", stderr(&out));
assert_eq!(Site::cloned_file(&site.owner, "owner"), "owner\n");
assert_not_found(
&site.clone_as(&reader, "reader-none"),
&site.owner.namespace(),
"no grant",
);
let (_anon, out) = site.clone_anonymously("anon");
assert_not_found(&out, &site.owner.namespace(), "anonymous");
let header = grant(&site.owner, &reader, "read", &[]);
let out = reader.ok(&["grant", "add", &grant_file(&reader, "read.txt", &header)]);
assert!(out.starts_with("added grant "), "{out}");
let out = site.clone_as(&reader, "reader");
assert!(out.status.success(), "{}", stderr(&out));
assert_eq!(Site::cloned_file(&reader, "reader"), "owner\n");
let header = grant(
&site.owner,
&writer,
"write",
&["--refs", "refs/heads/*=cuf"],
);
let out = writer.ok(&["grant", "add", &grant_file(&writer, "write.txt", &header)]);
assert!(out.starts_with("added grant "), "{out}");
assert_not_found(
&site.clone_as(&writer, "writer"),
&site.owner.namespace(),
"write-only grant",
);
}
#[test]
fn an_epoch_bump_revokes_a_read_grant_and_a_reissue_at_the_new_epoch_works() {
let site = Site::new();
let reader = site.party(0x22);
site.owner.ok(&["visibility", "set", "origin", "private"]);
let header = grant(&site.owner, &reader, "read", &[]);
reader.ok(&["grant", "add", &grant_file(&reader, "g0.txt", &header)]);
let out = site.clone_as(&reader, "before");
assert!(out.status.success(), "{}", stderr(&out));
let bumped = site.owner.ok(&["epoch", "bump", "origin"]);
assert!(bumped.contains("epoch 1 (was 0)"), "{bumped}");
assert_not_found(
&site.clone_as(&reader, "revoked"),
&site.owner.namespace(),
"stale-epoch read grant",
);
let header = grant(&site.owner, &reader, "read", &[]);
let fresh =
mkit_attest::grant::Grant::parse(&SignedHeader::parse(&header).unwrap().statement).unwrap();
assert_eq!(fresh.epoch, 1);
reader.ok(&["grant", "add", &grant_file(&reader, "g1.txt", &header)]);
let out = site.clone_as(&reader, "after");
assert!(out.status.success(), "{}", stderr(&out));
assert_eq!(Site::cloned_file(&reader, "after"), "owner\n");
}
#[test]
fn an_older_unsubmitted_public_statement_cannot_undo_a_later_envelope_flip() {
let site = Site::new();
let anon = Party::new(0x77);
let repository = format!("{}/site", site.owner.namespace());
let seed = ed25519_dalek::SigningKey::from_bytes(&[0x11; 32]);
let created = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| i64::try_from(d.as_millis()).unwrap())
.unwrap();
let statement = mkit_attest::grant::VisibilityStatement {
repository: mkit_core::repo_identity::RepositoryIdentity::parse(&repository).unwrap(),
visibility: mkit_attest::grant::Visibility::Public,
audiences: vec![site.live.origin.clone()],
created_ms: created,
expiry_ms: created + 3_600_000,
nonce: [9; 32],
}
.encode()
.unwrap();
let header = SignedHeader {
statement: statement.clone(),
scheme: OwnerScheme::Ed25519,
blob: ed25519_dalek::Signer::sign(&seed, &mkit_core::hash::hash(&statement))
.to_bytes()
.to_vec(),
}
.encode()
.unwrap();
std::thread::sleep(Duration::from_millis(20));
site.owner.ok(&["visibility", "set", "origin", "private"]);
assert_not_found(
&site.clone_as(&anon, "private"),
&site.owner.namespace(),
"after the envelope flip",
);
let response = reqwest::blocking::Client::new()
.post(format!(
"{}/mkit.transport.v1.TransportService/SetRepoVisibility",
site.live.origin
))
.header("content-type", "application/json")
.header("connect-protocol-version", "1")
.header("x-repository", &repository)
.body(serde_json::json!({ "signedStatement": header }).to_string())
.send()
.unwrap();
assert_eq!(response.status(), 403, "permission_denied");
let body = response.text().unwrap();
assert!(
body.contains("permission_denied") && body.contains("not newer than the stored statement"),
"{body}"
);
assert_not_found(
&site.clone_as(&anon, "still-private"),
&site.owner.namespace(),
"after the stale statement",
);
}