#![allow(clippy::unnecessary_wraps)]
use std::collections::VecDeque;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use bytes::Bytes;
use futures_executor::block_on;
use mkit_core::hash::{Hash, hash};
use mkit_core::protocol::PackKey;
use mkit_core::refs::RefWriteCondition;
use mkit_core::repo_identity::Namespace;
use mkit_rpc::mkit::common::v1::RefExpectation;
use mkit_rpc::mkit::rpc::v1::ssh::{
DownloadPack, DownloadPackHeader, Hello, ListRefs, PackChunk, PackExists, ReadRef, SshFrame,
UpdateRef, UploadPack, ssh_frame,
};
use mkit_rpc::mkit::rpc::v1::{Error as RpcError, ErrorCode, ProtocolVersion};
use super::verbs::Verbs;
use super::*;
use crate::error::{Redacted, ServerError};
use crate::op::Procedure;
use crate::pipeline::{
Admission, AdmissionDecision, AdmissionInput, AuthMode, D34Shards, HookSet, Hooks,
OpenAuthorizer, Pipeline, PipelineConfig, ShardMap, Sharding,
};
use crate::policy::{NamespacePolicy, WritePolicy};
use crate::principal::Principal;
use crate::repo::{Addressing, MultiAddressing, NamespaceKey, RepoId, RepoName};
use crate::rt::ManualClock;
use crate::store::{
Batch, BatchOutcome, BlobKey, BlobStore, Cursor, Key, MultipartBlobStore, NamespaceStore,
Partition, PartitionStats, ScanPage, StoreCapabilities, StoreError, UnsupportedPartSink, Value,
Write, codec, keys,
};
use crate::telemetry::NoopMetrics;
use crate::upload::UploadError;
use crate::upload::token::TicketKeys;
use crate::{MemoryBlobStore, MemoryFault, MemoryKv};
type Body = ssh_frame::Body;
const SERVER_ID: &str = "mkit serve/0.4.2";
const T0: i64 = 1_700_000_000_000;
struct VecSource(VecDeque<Result<SshFrame, FrameIoError>>);
impl FrameSource for VecSource {
async fn next_frame(&mut self) -> Result<SshFrame, FrameIoError> {
self.0.pop_front().unwrap_or(Err(FrameIoError::Eof))
}
}
#[derive(Default)]
struct VecSink {
frames: Vec<SshFrame>,
fail_after: Option<usize>,
}
impl FrameSink for VecSink {
async fn send(&mut self, frame: &SshFrame) -> Result<(), FrameIoError> {
if self.fail_after.is_some_and(|n| self.frames.len() >= n) {
return Err(FrameIoError::Io(Redacted::new("peer closed")));
}
self.frames.push(frame.clone());
Ok(())
}
}
fn frame(body: Option<Body>) -> SshFrame {
SshFrame {
body,
..Default::default()
}
}
fn hello() -> Option<Body> {
let hello = Hello::default().with_proto(ProtocolVersion::ProtocolVersion1);
Some(Body::Hello(Box::new(hello)))
}
fn script(bodies: impl IntoIterator<Item = Option<Body>>) -> VecSource {
let frames = core::iter::once(hello()).chain(bodies);
VecSource(frames.map(|b| Ok(frame(b))).collect())
}
fn repo() -> RepoId {
RepoId {
namespace: NamespaceKey::deployment_default(),
name: RepoName::new("room-a").unwrap(),
}
}
fn cfg(auth: AuthMode) -> PipelineConfig {
PipelineConfig::new(Addressing::Single { repo: repo() }, auth, upload_limits())
}
fn pipeline<B: MultipartBlobStore, N: NamespaceStore>(
blobs: B,
meta: N,
auth: AuthMode,
) -> Pipeline<B, N> {
let clock = Arc::new(ManualClock::new(T0));
let metrics = Arc::new(NoopMetrics);
Pipeline::new(blobs, meta, Hooks::new(), cfg(auth), clock, metrics).unwrap()
}
fn mem() -> (Pipeline<MemoryBlobStore, MemoryKv>, MemoryBlobStore) {
let blobs = MemoryBlobStore::default();
let clock = Arc::new(ManualClock::new(T0));
let meta = MemoryKv::with_clock(clock);
let pipe = pipeline(blobs.clone(), meta, AuthMode::TransportIdentity);
(pipe, blobs)
}
fn principal() -> Principal {
Principal::SshForcedCommand { key: None }
}
fn serve<B: MultipartBlobStore, N: NamespaceStore, H: HookSet>(
pipe: &Pipeline<B, N, H>,
mut src: VecSource,
sink: &mut VecSink,
cfg: &SessionConfig,
) -> SessionEnd {
block_on(serve_session(pipe, principal(), &mut src, sink, cfg))
}
fn run<B: MultipartBlobStore, N: NamespaceStore, H: HookSet>(
pipe: &Pipeline<B, N, H>,
bodies: impl IntoIterator<Item = Option<Body>>,
) -> (SessionEnd, Vec<SshFrame>) {
let mut sink = VecSink::default();
let end = serve(
pipe,
script(bodies),
&mut sink,
&SessionConfig::new(SERVER_ID),
);
let mut frames = sink.frames.into_iter();
let first = frames.next().expect("a HelloResponse");
let Some(Body::HelloResponse(resp)) = first.body else {
panic!("expected HelloResponse, got {:?}", first.body);
};
assert_eq!(resp.server_id.as_deref(), Some(SERVER_ID));
(end, frames.collect())
}
#[track_caller]
fn error(f: &SshFrame) -> &RpcError {
match &f.body {
Some(Body::Error(e)) => e,
other => panic!("expected an Error frame, got {other:?}"),
}
}
#[track_caller]
fn assert_error(f: &SshFrame, code: ErrorCode, message: &str) {
let e = error(f);
assert!(e.code.is_some_and(|c| c == code), "code of {e:?}");
assert_eq!(e.message.as_deref(), Some(message));
assert_eq!(e.details.as_deref().map_or(0, <[u8]>::len), 0);
}
fn upload_header(id: &[u8], total: Option<u64>) -> Option<Body> {
Some(Body::UploadPack(Box::new(UploadPack {
pack_id: Some(id.to_vec()),
total_bytes: total,
..Default::default()
})))
}
fn chunk(id: &[u8], offset: Option<u64>, data: &[u8], last: bool) -> Option<Body> {
Some(Body::PackChunk(Box::new(PackChunk {
pack_id: Some(id.to_vec()),
offset,
data: Some(data.to_vec()),
last: Some(last),
..Default::default()
})))
}
fn update(
name: &str,
new: &[u8],
expectation: Option<RefExpectation>,
expected: Option<&[u8]>,
) -> Option<Body> {
let mut req = UpdateRef::default()
.with_name(name)
.with_new_id(new.to_vec());
if let Some(e) = expectation {
req = req.with_expectation(e);
}
if let Some(e) = expected {
req = req.with_expected_id(e.to_vec());
}
Some(Body::UpdateRef(Box::new(req)))
}
fn read_ref(name: &str) -> Option<Body> {
Some(Body::ReadRef(Box::new(ReadRef::default().with_name(name))))
}
fn list_refs(prefix: Option<&str>) -> Option<Body> {
let mut req = ListRefs::default();
if let Some(p) = prefix {
req = req.with_prefix(p);
}
Some(Body::ListRefs(Box::new(req)))
}
fn exists(id: &[u8]) -> Option<Body> {
let req = PackExists::default().with_pack_id(id.to_vec());
Some(Body::PackExists(Box::new(req)))
}
fn download(id: Option<&[u8]>) -> Option<Body> {
let mut req = DownloadPack::default();
if let Some(id) = id {
req = req.with_pack_id(id.to_vec());
}
Some(Body::DownloadPack(Box::new(req)))
}
fn close() -> Option<Body> {
Some(Body::Close(Box::default()))
}
fn pack_bytes(len: usize, seed: u8) -> Vec<u8> {
#[allow(clippy::cast_possible_truncation)]
(0..len)
.map(|i| (i as u8).wrapping_mul(31).wrapping_add(seed))
.collect()
}
fn seed<B: MultipartBlobStore, N: NamespaceStore, H: HookSet>(
pipe: &Pipeline<B, N, H>,
refs: &[(&str, Hash)],
packs: &[&[u8]],
) {
let mut bodies = Vec::new();
for (name, id) in refs {
bodies.push(update(name, id, Some(RefExpectation::Any), None));
}
for pack in packs {
let id = hash(pack);
bodies.push(upload_header(&id, Some(pack.len() as u64)));
bodies.push(chunk(&id, Some(0), pack, true));
}
let (end, frames) = run(pipe, bodies);
assert_eq!(end, SessionEnd::Clean);
let ok = |f: &SshFrame| {
matches!(
f.body,
Some(Body::UpdateRefResponse(_) | Body::UploadPackResponse(_))
)
};
assert!(frames.iter().all(ok), "seeding failed: {frames:?}");
}
fn blob_present(blobs: &impl BlobStore, id: Hash) -> bool {
block_on(blobs.head(&BlobKey::pack(id))).unwrap().is_some()
}
fn valid_pack() -> (Vec<u8>, Hash) {
let bytes = b"valid pack bytes".to_vec();
let id = hash(&bytes);
(bytes, id)
}
#[test]
fn handshake_rejects_non_hello_first_frame() {
let (pipe, _) = mem();
let src = VecSource(VecDeque::from([Ok(frame(list_refs(None)))]));
let mut sink = VecSink::default();
let end = serve(&pipe, src, &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::ProtocolError);
assert_eq!(sink.frames.len(), 1);
assert_error(
&sink.frames[0],
ErrorCode::InvalidRequest,
"first frame must be Hello",
);
}
#[test]
fn handshake_rejects_unsupported_proto() {
let (pipe, _) = mem();
let unspecified = Some(Body::Hello(Box::default()));
let src = VecSource(VecDeque::from([Ok(frame(unspecified))]));
let mut sink = VecSink::default();
let end = serve(&pipe, src, &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::ProtocolError);
assert_eq!(sink.frames.len(), 1);
assert_error(
&sink.frames[0],
ErrorCode::InvalidRequest,
"unsupported proto_version 0",
);
}
#[test]
fn handshake_read_failure_is_silent_protocol_error() {
let (pipe, _) = mem();
for err in [FrameIoError::Eof, FrameIoError::Malformed] {
let mut sink = VecSink::default();
let src = VecSource(VecDeque::from([Err(err)]));
let end = serve(&pipe, src, &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::ProtocolError);
assert!(sink.frames.is_empty());
}
let mut sink = VecSink {
fail_after: Some(0),
..VecSink::default()
};
let end = serve(&pipe, script([]), &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::ProtocolError, "HelloResponse write failed");
}
#[test]
fn stop_after_hello_ends_clean() {
let (pipe, _) = mem();
let mut cfg = SessionConfig::new(SERVER_ID);
cfg.stop_after_hello = true;
let mut sink = VecSink::default();
let end = serve(&pipe, script([list_refs(None)]), &mut sink, &cfg);
assert_eq!(end, SessionEnd::Clean);
assert_eq!(sink.frames.len(), 1, "only the HelloResponse");
assert!(matches!(sink.frames[0].body, Some(Body::HelloResponse(_))));
}
#[test]
fn close_and_eof_end_clean() {
let (pipe, _) = mem();
let (end, frames) = run(&pipe, [close(), list_refs(None)]);
assert_eq!(end, SessionEnd::Clean);
assert!(frames.is_empty(), "nothing after Close is answered");
let (end, frames) = run(&pipe, []);
assert_eq!(end, SessionEnd::Clean);
assert!(frames.is_empty());
}
#[test]
fn pipeline_outside_transport_identity_is_refused() {
let blobs = MemoryBlobStore::default();
let pipe = pipeline(blobs, MemoryKv::default(), AuthMode::Open);
let mut src = script([]);
let mut sink = VecSink::default();
let cfg = SessionConfig::new(SERVER_ID);
let end = block_on(serve_session(&pipe, principal(), &mut src, &mut sink, &cfg));
assert_eq!(end, SessionEnd::ProtocolError);
assert!(sink.frames.is_empty(), "refused before the handshake");
assert_eq!(src.0.len(), 1, "the Hello is still unread");
}
#[test]
fn malformed_frame_emits_parse_error_and_ends() {
let (pipe, _) = mem();
for err in [
FrameIoError::Malformed,
FrameIoError::Io(Redacted::new("reset")),
] {
let mut src = script([]);
src.0.push_back(Err(err));
src.0.push_back(Ok(frame(list_refs(None))));
let mut sink = VecSink::default();
let end = serve(&pipe, src, &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::ProtocolError);
assert_eq!(sink.frames.len(), 2);
assert_error(
&sink.frames[1],
ErrorCode::InvalidRequest,
"frame parse error",
);
}
}
#[test]
fn timeout_from_source_ends_timeout() {
let (pipe, blobs) = mem();
let mut sink = VecSink::default();
let src = VecSource(VecDeque::from([Err(FrameIoError::Timeout)]));
let end = serve(&pipe, src, &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::Timeout);
assert!(sink.frames.is_empty());
let mut src = script([list_refs(None)]);
src.0.push_back(Err(FrameIoError::Timeout));
let mut sink = VecSink::default();
let end = serve(&pipe, src, &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::Timeout);
assert_eq!(sink.frames.len(), 2, "HelloResponse and ListRefsResponse");
let (bytes, id) = valid_pack();
let mut src = script([
upload_header(&id, Some(bytes.len() as u64)),
chunk(&id, Some(0), &bytes[..4], false),
]);
src.0.push_back(Err(FrameIoError::Timeout));
let mut sink = VecSink::default();
let end = serve(&pipe, src, &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::Timeout);
assert_eq!(sink.frames.len(), 1, "only the HelloResponse");
assert!(!blob_present(&blobs, id));
}
#[test]
fn sink_failure_after_handshake_ends_io_error() {
let (pipe, _) = mem();
let mut sink = VecSink {
fail_after: Some(1),
..VecSink::default()
};
let src = script([list_refs(None), list_refs(None)]);
let end = serve(&pipe, src, &mut sink, &SessionConfig::new(SERVER_ID));
assert_eq!(end, SessionEnd::IoError);
assert_eq!(sink.frames.len(), 1);
}
#[test]
fn frame_budget_exceeded_ends_protocol_error() {
let (pipe, _) = mem();
let n = MAX_FRAMES_PER_CONN as usize;
let (end, frames) = run(&pipe, core::iter::repeat_n(None, n + 2));
assert_eq!(end, SessionEnd::ProtocolError);
assert_eq!(frames.len(), n + 1);
assert_error(&frames[n - 1], ErrorCode::InvalidRequest, "empty frame");
assert_error(
&frames[n],
ErrorCode::InvalidRequest,
"per-connection frame budget exceeded",
);
}
#[test]
fn byte_budget_exceeded_ends_protocol_error() {
let (pipe, _) = mem();
let header = |total| {
Some(Body::DownloadPackHeader(Box::new(DownloadPackHeader {
total_bytes: Some(total),
..Default::default()
})))
};
let half = MAX_BYTES_PER_CONN / 2 + 1;
let (end, frames) = run(&pipe, [header(half), header(half), list_refs(None)]);
assert_eq!(end, SessionEnd::ProtocolError);
assert_eq!(frames.len(), 2);
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
"unexpected request frame",
);
assert_error(
&frames[1],
ErrorCode::InvalidRequest,
"per-connection byte budget exceeded",
);
let (_, id) = valid_pack();
let (end, frames) = run(&pipe, [upload_header(&id, Some(u64::MAX))]);
assert_eq!(end, SessionEnd::ProtocolError);
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
"per-connection byte budget exceeded",
);
}
#[test]
fn frame_byte_estimate_matches_mkit_serve() {
let est = |body| frame_byte_estimate(&frame(body));
assert_eq!(est(chunk(&[0; 32], Some(0), &[7; 100], false)), 100);
assert_eq!(est(upload_header(&[0; 32], Some(12_345))), 12_345);
assert_eq!(est(upload_header(&[0; 32], None)), 0);
assert_eq!(est(list_refs(None)), 64);
assert_eq!(est(None), 64);
assert_eq!(
upload_limits(),
crate::upload::UploadLimits {
max_total_bytes: 1024 * 1024 * 1024,
max_chunks: 10_000,
}
);
}
#[test]
fn misplaced_frames_get_mkit_serve_errors() {
let (pipe, _) = mem();
let (end, frames) = run(
&pipe,
[
chunk(&[1; 32], Some(0), b"x", true),
hello(),
None,
Some(Body::HelloResponse(Box::default())),
Some(Body::Error(Box::default())),
],
);
assert_eq!(end, SessionEnd::Clean);
let want = [
"PackChunk arrived without UploadPack header",
"Hello after handshake",
"empty frame",
"unexpected request frame",
"unexpected request frame",
];
assert_eq!(frames.len(), want.len());
for (f, message) in frames.iter().zip(want) {
assert_error(f, ErrorCode::InvalidRequest, message);
}
}
#[test]
fn pack_key_from_id_rejects_bad_length_as_invalid_request() {
let (pipe, _) = mem();
let missing = Some(Body::PackExists(Box::default()));
let (end, frames) = run(
&pipe,
[
exists(&[0; 16]),
missing,
download(Some(&[0; 16])),
download(None),
exists(&[7; 32]),
],
);
assert_eq!(end, SessionEnd::Clean);
let bad = [
"pack_id must be 32 bytes",
"pack_id missing",
"pack_id must be 32 bytes",
"pack_id missing",
];
for (f, message) in frames.iter().zip(bad) {
assert_error(f, ErrorCode::InvalidRequest, message);
}
let Some(Body::PackExistsResponse(resp)) = &frames[4].body else {
panic!("a correct 32-byte id still decodes");
};
assert_eq!(resp.exists, Some(false));
}
#[test]
fn update_ref_field_errors_keep_ssh_messages() {
let (pipe, _) = mem();
let main = "refs/heads/main";
let (end, frames) = run(
&pipe,
[
update(main, &[1; 5], Some(RefExpectation::Any), None),
update(main, &[1; 32], None, None),
update(main, &[1; 32], Some(RefExpectation::Match), None),
update(main, &[1; 32], Some(RefExpectation::Match), Some(&[2; 31])),
update(main, &[], None, None),
update(
"refs/heads/.hidden",
&[1; 32],
Some(RefExpectation::Any),
None,
),
],
);
assert_eq!(end, SessionEnd::Clean);
let want = [
"new_id must be 32 bytes",
"UpdateRef.expectation is required",
"MATCH expectation requires a 32-byte expected_id",
"MATCH expectation requires a 32-byte expected_id",
"new_id must be 32 bytes",
"update ref failed",
];
assert_eq!(frames.len(), want.len());
for (f, message) in frames.iter().zip(want) {
assert_error(f, ErrorCode::InvalidRequest, message);
}
}
#[test]
fn any_and_missing_ignore_expected_id() {
let (pipe, _) = mem();
let (end, frames) = run(
&pipe,
[
update(
"refs/heads/a",
&[1; 32],
Some(RefExpectation::Any),
Some(&[9; 3]),
),
update(
"refs/heads/b",
&[1; 32],
Some(RefExpectation::Missing),
Some(&[9; 32]),
),
read_ref("refs/heads/b"),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UpdateRefResponse(_))));
assert!(matches!(frames[1].body, Some(Body::UpdateRefResponse(_))));
let Some(Body::ReadRefResponse(resp)) = &frames[2].body else {
panic!("expected ReadRefResponse");
};
assert_eq!(resp.object_id.as_deref(), Some(&[1; 32][..]));
}
#[test]
fn read_and_list_refs() {
let (pipe, _) = mem();
seed(
&pipe,
&[("refs/heads/main", [1; 32]), ("refs/tags/v1", [2; 32])],
&[],
);
let (end, frames) = run(
&pipe,
[
read_ref("refs/heads/main"),
read_ref("refs/heads/absent"),
read_ref("refs/heads/.hidden"),
list_refs(Some("refs/heads/")),
list_refs(None),
list_refs(Some("/bad")),
],
);
assert_eq!(end, SessionEnd::Clean);
let read = |f: &SshFrame| match &f.body {
Some(Body::ReadRefResponse(r)) => r.object_id.clone(),
other => panic!("expected ReadRefResponse, got {other:?}"),
};
assert_eq!(read(&frames[0]), Some(vec![1; 32]));
assert_eq!(read(&frames[1]), Some(Vec::new()), "absent is empty");
assert_error(&frames[2], ErrorCode::Internal, "read ref failed");
let listed = |f: &SshFrame| match &f.body {
Some(Body::ListRefsResponse(r)) => r
.refs
.iter()
.map(|e| (e.name.clone().unwrap(), e.object_id.clone().unwrap()))
.collect::<Vec<_>>(),
other => panic!("expected ListRefsResponse, got {other:?}"),
};
assert_eq!(listed(&frames[3]), [("main".to_owned(), vec![1; 32])]);
assert_eq!(
listed(&frames[4]),
[
("refs/heads/main".to_owned(), vec![1; 32]),
("refs/tags/v1".to_owned(), vec![2; 32]),
]
);
assert_error(&frames[5], ErrorCode::Internal, "list refs failed");
}
#[test]
fn serve_loop_cas_conflict_carries_current_id_in_details() {
let (pipe, _) = mem();
let (winner, loser) = ([0xA1u8; 32], [0xB2u8; 32]);
let main = "refs/heads/main";
let (end, frames) = run(
&pipe,
[
update(main, &winner, Some(RefExpectation::Missing), None),
update(main, &loser, Some(RefExpectation::Missing), None),
read_ref(main),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UpdateRefResponse(_))));
let err = error(&frames[1]).clone();
assert!(err.code.is_some_and(|c| c == ErrorCode::InvalidRequest));
assert_eq!(err.details.as_deref(), Some(&winner[..]));
assert!(matches!(
mkit_rpc::map_update_ref_error(err, RefWriteCondition::Missing, "ssh"),
mkit_core::protocol::TransportError::RefConflict
));
let Some(Body::ReadRefResponse(resp)) = &frames[2].body else {
panic!("expected ReadRefResponse");
};
assert_eq!(
resp.object_id.as_deref(),
Some(&winner[..]),
"loser clobbered nothing"
);
}
#[test]
fn serve_loop_match_conflict_reports_current_value() {
let (pipe, _) = mem();
let (current, stale, next) = ([0x11u8; 32], [0x22u8; 32], [0x33u8; 32]);
let main = "refs/heads/main";
seed(&pipe, &[(main, current)], &[]);
let (_, frames) = run(
&pipe,
[
update(main, &next, Some(RefExpectation::Match), Some(&stale)),
read_ref(main),
],
);
let err = error(&frames[0]);
assert!(err.code.is_some_and(|c| c == ErrorCode::InvalidRequest));
assert_eq!(err.details.as_deref(), Some(¤t[..]));
assert_eq!(
err.message.as_deref(),
Some("ref update conflict: expectation does not match current ref value")
);
let Some(Body::ReadRefResponse(resp)) = &frames[1].body else {
panic!("expected ReadRefResponse");
};
assert_eq!(resp.object_id.as_deref(), Some(¤t[..]));
}
#[test]
fn serve_loop_match_conflict_on_absent_ref_has_empty_details() {
let (pipe, _) = mem();
let ghost = "refs/heads/ghost";
let (_, frames) = run(
&pipe,
[
update(
ghost,
&[0x44; 32],
Some(RefExpectation::Match),
Some(&[0x55; 32]),
),
read_ref(ghost),
],
);
let err = error(&frames[0]).clone();
assert!(err.code.is_some_and(|c| c == ErrorCode::InvalidRequest));
assert_eq!(err.details.as_deref().map_or(0, <[u8]>::len), 0);
let mapped = mkit_rpc::map_update_ref_error(err, RefWriteCondition::Match([0x55; 32]), "ssh");
match mapped {
mkit_core::protocol::TransportError::RemoteError(msg) => {
assert!(msg.contains("absent"), "message should say absent: {msg}");
}
other => panic!("expected RemoteError, got {other:?}"),
}
let Some(Body::ReadRefResponse(resp)) = &frames[1].body else {
panic!("expected ReadRefResponse");
};
assert_eq!(resp.object_id.as_deref(), Some(&[][..]));
}
#[test]
fn cas_conflict_body_matches_mkit_serve() {
let Body::Error(e) = cas_conflict_body(Some([7; 32])) else {
panic!("an Error body");
};
assert!(e.code.is_some_and(|c| c == ErrorCode::InvalidRequest));
assert_eq!(e.details.as_deref(), Some(&[7; 32][..]));
let Body::Error(e) = cas_conflict_body(None) else {
panic!("an Error body");
};
assert_eq!(
e.message.as_deref(),
Some("ref update conflict: expectation not met and ref is currently absent")
);
assert_eq!(e.details.as_deref(), Some(&[][..]));
}
#[test]
fn upload_drain_accepts_valid_chunks() {
let (pipe, blobs) = mem();
let (bytes, id) = valid_pack();
let (end, frames) = run(
&pipe,
[
upload_header(&id, Some(bytes.len() as u64)),
chunk(&id, Some(0), &bytes[..5], false),
chunk(&id, Some(5), &bytes[5..], true),
exists(&id),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UploadPackResponse(_))));
let Some(Body::PackExistsResponse(resp)) = &frames[1].body else {
panic!("expected PackExistsResponse");
};
assert_eq!(resp.exists, Some(true));
assert!(blob_present(&blobs, id));
}
#[test]
fn upload_drain_rejects_malformed_streams() {
let (bytes, id) = valid_pack();
let len = bytes.len() as u64;
let wrong = b"wrong pack bytes";
let cases: Vec<(Vec<Option<Body>>, &str)> = vec![
(
vec![upload_header(&id, None)],
"UploadPack.total_bytes is required",
),
(
vec![upload_header(&[1; 7], Some(len))],
"pack_id must be 32 bytes",
),
(
vec![
upload_header(&id, Some(len)),
chunk(&id, Some(1), &bytes, true),
],
"PackChunk.offset is not the expected next offset",
),
(
vec![
upload_header(&id, Some(len)),
chunk(&id, None, &bytes, true),
],
"PackChunk.offset is required",
),
(
vec![
upload_header(&id, Some(len)),
chunk(&[0xAA; 32], Some(0), &bytes, true),
],
"PackChunk.pack_id does not match UploadPack",
),
(
vec![
upload_header(&id, Some(len - 1)),
chunk(&id, Some(0), &bytes, true),
],
"PackChunk data exceeds declared total_bytes",
),
(
vec![
upload_header(&id, Some(len)),
chunk(&id, Some(0), &bytes[..bytes.len() - 1], true),
],
"PackChunk stream ended before declared total_bytes",
),
(
vec![
upload_header(&id, Some(wrong.len() as u64)),
chunk(&id, Some(0), wrong, true),
],
"uploaded pack bytes do not match UploadPack.pack_id",
),
(
vec![upload_header(&id, Some(len)), list_refs(None)],
"expected PackChunk after UploadPack",
),
(
vec![upload_header(&id, Some(len))],
"pack chunk read failed",
),
];
for (bodies, message) in cases {
let (pipe, blobs) = mem();
let mut bodies = bodies;
let eof = bodies.len() == 1 && message == "pack chunk read failed";
if !eof {
bodies.push(exists(&id));
}
let (end, frames) = run(&pipe, bodies);
assert_eq!(end, SessionEnd::Clean, "{message}");
assert_error(&frames[0], ErrorCode::InvalidRequest, message);
assert!(!blob_present(&blobs, id), "{message}: nothing stored");
if !eof {
assert!(
matches!(frames[1].body, Some(Body::PackExistsResponse(_))),
"{message}"
);
}
}
}
#[test]
fn declared_size_over_the_cap() {
let (pipe, _) = mem();
let (_, id) = valid_pack();
let (end, frames) = run(&pipe, [upload_header(&id, Some(MAX_BYTES_PER_CONN + 1))]);
assert_eq!(end, SessionEnd::ProtocolError);
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
"per-connection byte budget exceeded",
);
let clock = Arc::new(ManualClock::new(T0));
let mut small = cfg(AuthMode::TransportIdentity);
small.upload_limits.max_total_bytes = 8;
let meta = MemoryKv::with_clock(clock.clone());
let pipe = Pipeline::new(
MemoryBlobStore::default(),
meta,
Hooks::new(),
small,
clock,
Arc::new(NoopMetrics),
)
.unwrap();
let (end, frames) = run(&pipe, [upload_header(&id, Some(9)), exists(&id)]);
assert_eq!(end, SessionEnd::Clean);
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
"UploadPack.total_bytes exceeds server cap",
);
assert!(matches!(frames[1].body, Some(Body::PackExistsResponse(_))));
}
#[test]
fn too_many_chunks_is_rejected() {
let (pipe, blobs) = mem();
let n = MAX_FRAMES_PER_CONN as usize;
let data = pack_bytes(n + 1, 3);
let id = hash(&data);
let mut bodies = vec![upload_header(&id, Some(data.len() as u64))];
for i in 0..=n {
bodies.push(chunk(&id, Some(i as u64), &data[i..=i], i == n));
}
let (end, frames) = run(&pipe, bodies);
assert_eq!(end, SessionEnd::Clean);
assert_eq!(frames.len(), 1);
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
"too many PackChunk frames before last=true",
);
assert!(!blob_present(&blobs, id));
}
#[test]
fn chunk_cap_holds_with_an_uncapped_chunk_pipeline() {
let clock = Arc::new(ManualClock::new(T0));
let mut uncapped = cfg(AuthMode::TransportIdentity);
uncapped.upload_limits.max_chunks = u32::MAX;
let blobs = MemoryBlobStore::default();
let meta = MemoryKv::with_clock(clock.clone());
let metrics = Arc::new(NoopMetrics);
let pipe = Pipeline::new(blobs.clone(), meta, Hooks::new(), uncapped, clock, metrics).unwrap();
let n = MAX_FRAMES_PER_CONN as usize;
let data = pack_bytes(n + 1, 3);
let id = hash(&data);
let mut bodies = vec![upload_header(&id, Some(data.len() as u64))];
for i in 0..=n {
bodies.push(chunk(&id, Some(i as u64), &data[i..=i], i == n));
}
bodies.push(exists(&id));
let (end, frames) = run(&pipe, bodies);
assert_eq!(end, SessionEnd::Clean);
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
"too many PackChunk frames before last=true",
);
assert!(!blob_present(&blobs, id));
let data = pack_bytes(n, 4);
let id = hash(&data);
let mut bodies = vec![upload_header(&id, Some(data.len() as u64))];
for i in 0..n {
bodies.push(chunk(&id, Some(i as u64), &data[i..=i], i + 1 == n));
}
let (_, frames) = run(&pipe, bodies);
assert!(matches!(frames[0].body, Some(Body::UploadPackResponse(_))));
}
#[test]
fn over_long_ref_names_are_refused_by_name() {
let (pipe, _) = mem();
let longest = format!(
"refs/heads/{}",
"a".repeat(crate::refs::MAX_REF_NAME_BYTES - 11)
);
let over = format!("{longest}a");
let (end, frames) = run(
&pipe,
[
update(&over, &[1; 32], Some(RefExpectation::Any), None),
read_ref(&over),
update(&over, &[1; 5], Some(RefExpectation::Any), None),
update(&longest, &[1; 32], Some(RefExpectation::Any), None),
read_ref(&longest),
list_refs(Some(&over)),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_error(&frames[0], ErrorCode::InvalidRequest, "ref name too long");
assert_error(&frames[1], ErrorCode::InvalidRequest, "ref name too long");
assert_error(
&frames[2],
ErrorCode::InvalidRequest,
"new_id must be 32 bytes",
);
assert!(matches!(frames[3].body, Some(Body::UpdateRefResponse(_))));
let Some(Body::ReadRefResponse(resp)) = &frames[4].body else {
panic!("expected ReadRefResponse");
};
assert_eq!(resp.object_id.as_deref(), Some(&[1; 32][..]));
assert_error(&frames[5], ErrorCode::Internal, "list refs failed");
}
#[test]
fn ref_names_outside_refs_are_refused_by_name() {
let (pipe, _) = mem();
let (end, frames) = run(
&pipe,
[
update("main", &[1; 32], Some(RefExpectation::Any), None),
read_ref("main"),
update("packs/x", &[1; 32], Some(RefExpectation::Missing), None),
read_ref("heads/main"),
update("main", &[1; 5], Some(RefExpectation::Any), None),
read_ref(".main"),
update(".main", &[1; 32], Some(RefExpectation::Any), None),
list_refs(Some("heads/")),
],
);
assert_eq!(end, SessionEnd::Clean);
let outside = crate::refs::REF_NAME_OUTSIDE_REFS;
for f in &frames[..4] {
assert_error(f, ErrorCode::InvalidRequest, outside);
}
assert_error(
&frames[4],
ErrorCode::InvalidRequest,
"new_id must be 32 bytes",
);
assert_error(&frames[5], ErrorCode::Internal, "read ref failed");
assert_error(&frames[6], ErrorCode::InvalidRequest, "update ref failed");
let Some(Body::ListRefsResponse(listed)) = &frames[7].body else {
panic!("expected ListRefsResponse, got {:?}", frames[7].body);
};
assert!(listed.refs.is_empty());
}
#[test]
fn list_refs_prefixes_match_at_component_boundaries() {
let (pipe, _) = mem();
let refs = [
("refs/heads/feat/x", [1; 32]),
("refs/heads/featx", [2; 32]),
("refs/heads/main", [3; 32]),
];
seed(&pipe, &refs, &[]);
let prefixes = [
"refs/heads",
"refs/heads/",
"refs/heads/feat",
"refs/heads/ma",
"refs/heads/main",
"refs//",
];
let (_, frames) = run(&pipe, prefixes.map(|p| list_refs(Some(p))));
let names: Vec<Vec<String>> = frames
.iter()
.map(|f| match &f.body {
Some(Body::ListRefsResponse(r)) => {
r.refs.iter().map(|e| e.name.clone().unwrap()).collect()
}
other => panic!("expected ListRefsResponse, got {other:?}"),
})
.collect();
let heads = ["feat/x", "featx", "main"].map(String::from).to_vec();
assert_eq!(names[0], heads);
assert_eq!(names[1], heads);
assert_eq!(names[2], ["x"]);
assert!(names[3].is_empty());
assert!(names[4].is_empty());
assert_eq!(names[5], ["heads/feat/x", "heads/featx", "heads/main"]);
}
#[test]
fn storage_failure_drains_then_answers_upload_failed() {
let data = pack_bytes(3_000, 5);
let id = hash(&data);
let bodies = || {
vec![
upload_header(&id, Some(data.len() as u64)),
chunk(&id, Some(0), &data[..1_000], false),
chunk(&id, Some(1_000), &data[1_000..2_000], false),
chunk(&id, Some(2_000), &data[2_000..], true),
exists(&id),
]
};
for fault in [MemoryFault::BlobWrite(0), MemoryFault::BlobCommit] {
let (pipe, blobs) = mem();
let _armed = blobs.clone().with_fault(fault);
let (end, frames) = run(&pipe, bodies());
assert_eq!(end, SessionEnd::Clean);
assert_eq!(frames.len(), 2, "{fault:?}");
assert_error(&frames[0], ErrorCode::Internal, "upload failed");
let Some(Body::PackExistsResponse(resp)) = &frames[1].body else {
panic!("expected PackExistsResponse");
};
assert_eq!(resp.exists, Some(false));
}
}
#[test]
fn serve_loop_rejects_invalid_upload_before_storage() {
let (pipe, blobs) = mem();
let bogus = [0x77; 32];
let (end, frames) = run(
&pipe,
[
upload_header(&bogus, Some(5)),
chunk(&bogus, Some(0), b"wrong", true),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(!blob_present(&blobs, bogus));
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
UploadError::DigestMismatch.ssh_message(),
);
}
#[test]
fn serve_loop_rejected_upload_does_not_overwrite_existing_pack() {
let (pipe, _) = mem();
let (bytes, id) = valid_pack();
seed(&pipe, &[], &[&bytes]);
let (end, frames) = run(
&pipe,
[
upload_header(&id, Some(5)),
chunk(&id, Some(0), b"wrong", true),
download(Some(&id)),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::Error(_))));
assert_eq!(downloaded(&frames[1..], &id), bytes);
}
#[test]
fn no_replay_or_quota_rows_are_written() {
let (pipe, _) = mem();
let (bytes, _) = valid_pack();
seed(
&pipe,
&[("refs/heads/main", [1; 32])],
&[&bytes, b"another"],
);
let p = Partition::Namespace(NamespaceKey::deployment_default());
let (start, end) = (Key::new(vec![0u8]), Key::new(vec![0xffu8]));
let page = block_on(pipe.meta_store().scan(&p, &start, &end, None, 100)).unwrap();
let tags: Vec<u8> = page.entries.iter().map(|(k, _)| k.as_bytes()[0]).collect();
assert!(!tags.is_empty());
assert!(
page.entries.iter().all(|(key, _)| matches!(
keys::parse(key),
Some(
keys::ParsedKey::Ref { .. }
| keys::ParsedKey::LayoutVersion
| keys::ParsedKey::Publication { .. }
| keys::ParsedKey::Advance { .. }
| keys::ParsedKey::PublishedRef { .. }
)
)),
"only ref, publication and layout-version rows, got tags {tags:?}"
);
}
#[track_caller]
fn downloaded(frames: &[SshFrame], id: &[u8]) -> Vec<u8> {
let Some(Body::DownloadPackHeader(h)) = &frames[0].body else {
panic!("expected DownloadPackHeader, got {:?}", frames[0].body);
};
let total = h.total_bytes.unwrap();
let mut out = Vec::new();
for (i, f) in frames[1..].iter().enumerate() {
let Some(Body::PackChunk(c)) = &f.body else {
panic!("expected PackChunk, got {:?}", f.body);
};
assert_eq!(c.pack_id.as_deref(), Some(id));
assert_eq!(c.offset, Some(out.len() as u64));
out.extend_from_slice(c.data.as_deref().unwrap());
if c.last == Some(true) {
assert_eq!(out.len() as u64, total);
assert_eq!(i + 2, frames.len(), "nothing after the last chunk");
return out;
}
}
panic!("no last chunk");
}
#[test]
fn download_emits_header_then_contiguous_chunks_last() {
let (pipe, _) = mem();
let max = crate::download::DOWNLOAD_CHUNK_MAX;
let big = pack_bytes(2 * max + 123, 11);
let id = hash(&big);
seed(&pipe, &[], &[&big, b""]);
let (end, frames) = run(&pipe, [download(Some(&id))]);
assert_eq!(end, SessionEnd::Clean);
assert_eq!(frames.len(), 4, "a header and three chunks");
let chunks: Vec<(u64, usize)> = frames[1..]
.iter()
.map(|f| match &f.body {
Some(Body::PackChunk(c)) => (c.offset.unwrap(), c.data.as_ref().unwrap().len()),
_ => panic!("chunk"),
})
.collect();
let max64 = max as u64;
assert_eq!(chunks, [(0, max), (max64, max), (2 * max64, 123)]);
assert_eq!(downloaded(&frames, &id), big);
let empty = hash(b"");
let (_, frames) = run(&pipe, [download(Some(&empty))]);
assert_eq!(frames.len(), 2);
assert!(downloaded(&frames, &empty).is_empty());
let (_, frames) = run(&pipe, [download(Some(&[9; 32]))]);
assert_error(&frames[0], ErrorCode::KeyNotFound, "pack not found");
}
#[test]
fn download_read_failure_after_header_is_internal() {
use crate::store::{BlobBody, BlobMeta, ByteRange, StoreError};
struct Failing(MemoryBlobStore);
impl BlobStore for Failing {
type Sink = <MemoryBlobStore as BlobStore>::Sink;
async fn begin(&self, key: BlobKey, len: u64) -> Result<Self::Sink, StoreError> {
self.0.begin(key, len).await
}
async fn get(
&self,
_: &BlobKey,
_: Option<ByteRange>,
) -> Result<Option<BlobBody>, StoreError> {
let piece: Result<Bytes, StoreError> = Ok(Bytes::from_static(b"abc"));
let pieces = [
piece,
Err(StoreError::unavailable(std::io::Error::other("gone"))),
];
let stream = futures_stream(pieces);
Ok(Some(BlobBody::Stream { len: 10, stream }))
}
async fn head(&self, key: &BlobKey) -> Result<Option<BlobMeta>, StoreError> {
self.0.head(key).await
}
async fn probe(&self) -> Result<(), StoreError> {
Ok(())
}
async fn delete(&self, key: &BlobKey) -> Result<bool, StoreError> {
self.0.delete(key).await
}
}
impl MultipartBlobStore for Failing {
type PartSink = UnsupportedPartSink;
const MAX_PARTS: u32 = u32::MAX;
}
let clock = Arc::new(ManualClock::new(T0));
let meta = MemoryKv::with_clock(clock);
let pipe = pipeline(
Failing(MemoryBlobStore::default()),
meta,
AuthMode::TransportIdentity,
);
let (end, frames) = run(&pipe, [download(Some(&[3; 32])), list_refs(None)]);
assert_eq!(end, SessionEnd::Clean);
let Some(Body::DownloadPackHeader(h)) = &frames[0].body else {
panic!("expected DownloadPackHeader");
};
assert_eq!(h.total_bytes, Some(10));
assert_error(&frames[1], ErrorCode::Internal, "pack read failed");
assert!(matches!(frames[2].body, Some(Body::ListRefsResponse(_))));
}
fn futures_stream<T: Send + Unpin + 'static>(
items: impl IntoIterator<Item = T>,
) -> crate::rt::BoxStream<'static, T> {
struct Iter<I>(I);
impl<I: Iterator + Unpin> futures_core::Stream for Iter<I> {
type Item = I::Item;
fn poll_next(
mut self: core::pin::Pin<&mut Self>,
_: &mut core::task::Context<'_>,
) -> core::task::Poll<Option<I::Item>> {
core::task::Poll::Ready(self.0.next())
}
}
let items: Vec<T> = items.into_iter().collect();
Box::pin(Iter(items.into_iter()))
}
#[test]
fn session_future_is_send() {
fn assert_send<T: Send>(_: &T) {}
let (pipe, _) = mem();
let mut src = script([]);
let mut sink = VecSink::default();
let cfg = SessionConfig::new(SERVER_ID);
let fut = serve_session(&pipe, principal(), &mut src, &mut sink, &cfg);
assert_send(&fut);
assert_eq!(block_on(fut), SessionEnd::Clean);
}
#[test]
fn read_frames_maps_framing_errors() {
let next = |bytes: Vec<u8>| block_on(ReadFrames(std::io::Cursor::new(bytes)).next_frame());
assert!(matches!(next(Vec::new()), Err(FrameIoError::Eof)));
assert!(matches!(next(vec![1, 0]), Err(FrameIoError::Eof)));
let too_long = (mkit_rpc::MAX_FRAME_BYTES + 1).to_le_bytes().to_vec();
assert!(matches!(next(too_long), Err(FrameIoError::Malformed)));
assert!(matches!(
next(vec![5, 0, 0, 0, 1]),
Err(FrameIoError::Malformed)
));
assert!(matches!(
next(vec![1, 0, 0, 0, 0xff]),
Err(FrameIoError::Malformed)
));
let (pipe, _) = mem();
let mut input = Vec::new();
mkit_rpc::write_frame(&mut input, &frame(hello())).unwrap();
input.extend_from_slice(&[1, 0, 0, 0, 0xff]);
let mut src = ReadFrames(std::io::Cursor::new(input));
let mut sink = VecSink::default();
let cfg = SessionConfig::new(SERVER_ID);
let end = block_on(serve_session(&pipe, principal(), &mut src, &mut sink, &cfg));
assert_eq!(end, SessionEnd::ProtocolError);
assert_error(
&sink.frames[1],
ErrorCode::InvalidRequest,
"frame parse error",
);
}
struct Golden {
seeded: Vec<u8>,
refs: [(&'static str, Hash); 3],
input: Vec<u8>,
}
#[allow(clippy::too_many_lines)]
fn golden() -> Golden {
let seeded = pack_bytes(1000, 7);
let seeded_id = hash(&seeded);
let fresh = pack_bytes(3000, 9);
let fresh_id = hash(&fresh);
let empty_id = hash(b"");
let bogus = pack_bytes(40, 1);
let bodies = vec![
hello(),
list_refs(Some("refs/heads/")),
list_refs(None),
list_refs(Some("/bad")),
read_ref("refs/heads/main"),
read_ref("refs/heads/absent"),
read_ref("refs/heads/.hidden"),
exists(&seeded_id),
exists(&[0x99; 32]),
exists(&[0x99; 16]),
update(
"refs/heads/main",
&[0x44; 32],
Some(RefExpectation::Match),
Some(&[0x11; 32]),
),
update(
"refs/heads/main",
&[0x55; 32],
Some(RefExpectation::Match),
Some(&[0x11; 32]),
),
update(
"refs/heads/ghost",
&[0x55; 32],
Some(RefExpectation::Match),
Some(&[0x66; 32]),
),
update(
"refs/heads/dev",
&[0x55; 32],
Some(RefExpectation::Missing),
None,
),
update(
"refs/heads/main",
&[0x55; 5],
Some(RefExpectation::Any),
None,
),
update("refs/heads/main", &[0x55; 32], None, None),
update(
"refs/heads/main",
&[0x55; 32],
Some(RefExpectation::Match),
None,
),
update(
"refs/heads/any",
&[0x77; 32],
Some(RefExpectation::Any),
Some(&[0x01; 32]),
),
update(
"refs/heads/.hidden",
&[0x77; 32],
Some(RefExpectation::Any),
None,
),
download(Some(&seeded_id)),
download(Some(&[0x99; 32])),
download(None),
upload_header(&fresh_id, Some(fresh.len() as u64)),
chunk(&fresh_id, Some(0), &fresh[..1200], false),
chunk(&fresh_id, Some(1200), &fresh[1200..], true),
upload_header(&empty_id, Some(0)),
chunk(&empty_id, Some(0), &[], true),
download(Some(&empty_id)),
upload_header(&fresh_id, None),
upload_header(&[0x01; 7], Some(3)),
upload_header(&[0x77; 32], Some(bogus.len() as u64)),
chunk(&[0x77; 32], Some(0), &bogus, true),
upload_header(&fresh_id, Some(fresh.len() as u64)),
chunk(&fresh_id, Some(5), &fresh, true),
upload_header(&fresh_id, Some(fresh.len() as u64)),
read_ref("refs/heads/main"),
chunk(&fresh_id, Some(0), &fresh[..10], false),
hello(),
None,
Some(Body::HelloResponse(Box::default())),
list_refs(Some("refs/")),
exists(&fresh_id),
download(Some(&fresh_id)),
close(),
read_ref("refs/heads/main"),
];
let mut input = Vec::new();
for body in bodies {
mkit_rpc::write_frame(&mut input, &frame(body)).unwrap();
}
Golden {
seeded,
refs: [
("refs/heads/main", [0x11; 32]),
("refs/heads/dev", [0x22; 32]),
("refs/tags/v1", [0x33; 32]),
],
input,
}
}
const GOLDEN_IN: &[u8] = include_bytes!("../../../../tests/golden/ssh-serve/session-1.in.bin");
const GOLDEN_OUT: &[u8] = include_bytes!("../../../../tests/golden/ssh-serve/session-1.bin");
fn replay_golden<B: MultipartBlobStore, N: NamespaceStore>(pipe: &Pipeline<B, N>) -> Vec<u8> {
let g = golden();
assert_eq!(
g.input, GOLDEN_IN,
"the builder must produce the captured script"
);
seed(pipe, &g.refs, &[&g.seeded]);
let mut src = ReadFrames(std::io::Cursor::new(g.input));
let mut sink = WriteFrames(Vec::new());
let cfg = SessionConfig::new(SERVER_ID);
let end = block_on(serve_session(pipe, principal(), &mut src, &mut sink, &cfg));
assert_eq!(end, SessionEnd::Clean, "mkit serve exited OK");
sink.0
}
fn assert_same_frames(got: &[u8], want: &[u8]) {
let decode = |bytes: &[u8]| {
let mut r = std::io::Cursor::new(bytes);
let mut out = Vec::new();
while r.position() < bytes.len() as u64 {
out.push(mkit_rpc::read_frame::<_, SshFrame>(&mut r).unwrap());
}
out
};
let (got, want) = (decode(got), decode(want));
for (i, (g, w)) in got.iter().zip(&want).enumerate() {
assert_eq!(g, w, "frame {i}");
}
assert_eq!(got.len(), want.len(), "frame count");
}
#[test]
fn golden_session_matches_mkit_serve_over_memory_stores() {
let (pipe, _) = mem();
let out = replay_golden(&pipe);
assert_same_frames(&out, GOLDEN_OUT);
assert_eq!(out, GOLDEN_OUT);
}
struct Golden2 {
refs: [(&'static str, Hash); 4],
big: Vec<u8>,
input: Vec<u8>,
}
fn golden2() -> Golden2 {
let big = pack_bytes(800 * 1024 + 1000, 5);
let big_id = hash(&big);
let prefixes = [
None,
Some(""),
Some("refs/heads"),
Some("refs/heads/"),
Some("refs/heads//"),
Some("refs/heads/feat"),
Some("refs/heads/feat/"),
Some("refs/heads/ma"),
Some("refs"),
Some("refs/"),
Some("refs//"),
Some("refs/tags"),
Some("nope/"),
];
let mut bodies = vec![hello()];
bodies.extend(prefixes.map(list_refs));
bodies.extend([download(Some(&big_id)), exists(&big_id), close()]);
let mut input = Vec::new();
for body in bodies {
mkit_rpc::write_frame(&mut input, &frame(body)).unwrap();
}
Golden2 {
refs: [
("refs/heads/main", [0x11; 32]),
("refs/heads/feat/x", [0x12; 32]),
("refs/heads/featx", [0x13; 32]),
("refs/tags/v1", [0x33; 32]),
],
big,
input,
}
}
const GOLDEN_2_IN: &[u8] = include_bytes!("../../../../tests/golden/ssh-serve/session-2.in.bin");
const GOLDEN_2_OUT: &[u8] = include_bytes!("../../../../tests/golden/ssh-serve/session-2.bin");
fn replay_golden_2<B: MultipartBlobStore, N: NamespaceStore>(pipe: &Pipeline<B, N>) -> Vec<u8> {
let input = golden2().input;
assert_eq!(
input, GOLDEN_2_IN,
"the builder must produce the captured script"
);
let mut src = ReadFrames(std::io::Cursor::new(input));
let mut sink = WriteFrames(Vec::new());
let cfg = SessionConfig::new(SERVER_ID);
let end = block_on(serve_session(pipe, principal(), &mut src, &mut sink, &cfg));
assert_eq!(end, SessionEnd::Clean, "mkit serve exited OK");
sink.0
}
#[test]
fn golden_session_2_matches_mkit_serve_over_memory_stores() {
let (pipe, _) = mem();
let g = golden2();
seed(&pipe, &g.refs, &[&g.big]);
let out = replay_golden_2(&pipe);
assert_same_frames(&out, GOLDEN_2_OUT);
assert_eq!(out, GOLDEN_2_OUT);
}
struct Golden3 {
ref_id: Hash,
input: Vec<u8>,
}
fn golden3() -> Golden3 {
let ref_id = [0x44; 32];
let pack = pack_bytes(64, 3);
let pack_id = hash(&pack);
let bodies = vec![
hello(),
upload_header(&pack_id, Some(pack.len() as u64)),
chunk(&pack_id, Some(0), &pack, true),
update(
BRANCH,
&[0x55; 32],
Some(RefExpectation::Match),
Some(&ref_id),
),
read_ref(BRANCH),
close(),
];
let mut input = Vec::new();
for body in bodies {
mkit_rpc::write_frame(&mut input, &frame(body)).unwrap();
}
Golden3 { ref_id, input }
}
const GOLDEN_3_IN: &[u8] = include_bytes!("../../../../tests/golden/ssh-serve/session-3.in.bin");
const GOLDEN_3_OUT: &[u8] = include_bytes!("../../../../tests/golden/ssh-serve/session-3.bin");
fn run_as<B: MultipartBlobStore, N: NamespaceStore, H: HookSet>(
pipe: &Pipeline<B, N, H>,
principal: Principal,
bodies: impl IntoIterator<Item = Option<Body>>,
) -> (SessionEnd, Vec<SshFrame>) {
let mut sink = VecSink::default();
let end = block_on(serve_session(
pipe,
principal,
&mut script(bodies),
&mut sink,
&SessionConfig::new(SERVER_ID),
));
let first = sink.frames.first().expect("a HelloResponse");
assert!(
matches!(first.body, Some(Body::HelloResponse(_))),
"got {:?}",
first.body
);
(end, sink.frames[1..].to_vec())
}
#[test]
fn golden_session_3_denies_a_non_owner_like_serve_root() {
let mut c = cfg(AuthMode::TransportIdentity);
c.addressing = Addressing::Single {
repo: peer_repo_id(OWNER, "room-a"),
};
c.write_policy = WritePolicy::Owner;
let clock = Arc::new(ManualClock::new(T0));
let pipe = Pipeline::new(
MemoryBlobStore::default(),
MemoryKv::with_clock(clock.clone()),
Hooks::new(),
c,
clock,
Arc::new(NoopMetrics),
)
.unwrap();
let g = golden3();
assert_eq!(g.input, GOLDEN_3_IN, "the builder must produce the file");
let owner = Principal::SshForcedCommand { key: Some(OWNER) };
let (end, frames) = run_as(
&pipe,
owner,
[update(BRANCH, &g.ref_id, Some(RefExpectation::Any), None)],
);
assert_eq!(end, SessionEnd::Clean, "seeding failed: {frames:?}");
let mut src = ReadFrames(std::io::Cursor::new(g.input));
let mut sink = WriteFrames(Vec::new());
let end = block_on(serve_session(
&pipe,
Principal::SshForcedCommand { key: Some(OTHER) },
&mut src,
&mut sink,
&SessionConfig::new(SERVER_ID),
));
assert_eq!(end, SessionEnd::Clean, "mkit serve exited OK");
assert_same_frames(&sink.0, GOLDEN_3_OUT);
assert_eq!(sink.0, GOLDEN_3_OUT);
}
#[test]
fn implicit_tickets_require_owner_policy_on_single() {
let build = |policy: WritePolicy| {
let mut c = cfg(AuthMode::TransportIdentity);
c.addressing = Addressing::Single {
repo: peer_repo_id(OWNER, "room-a"),
};
c.write_policy = policy;
let clock = Arc::new(ManualClock::new(T0));
Pipeline::new(
MemoryBlobStore::default(),
MemoryKv::with_clock(clock.clone()),
Hooks::new(),
c,
clock,
Arc::new(NoopMetrics),
)
.unwrap()
};
assert!(build(WritePolicy::Owner).implicit_tickets());
assert!(!build(WritePolicy::Open).implicit_tickets());
let (plain, _) = mem();
assert!(!plain.implicit_tickets());
}
#[test]
fn transport_identity_write_never_carries_a_grant() {
let (pipe, _, _) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let verbs = Verbs::new(
&pipe,
Principal::TransportPeer { ed25519: OWNER },
Some(peer_repo(OWNER, "room-a")),
);
for procedure in [
Procedure::UpdateRef,
Procedure::AdvanceRefs,
Procedure::BeginUpload,
Procedure::UploadPack,
] {
let a = verbs.auth(procedure).unwrap();
assert!(a.write_grant.is_none(), "{procedure:?} carried a grant");
}
}
const OWNER: [u8; 32] = [0x0a; 32];
const OTHER: [u8; 32] = [0x0b; 32];
const BRANCH: &str = "refs/heads/main";
const PACKMAP: &str = "refs/mkit/packmap/main";
fn peer_repo(key: [u8; 32], name: &str) -> String {
format!("{}/{name}", Namespace::Ed25519(key))
}
fn peer_repo_id(key: [u8; 32], name: &str) -> RepoId {
RepoId {
namespace: NamespaceKey::from_namespace(&Namespace::Ed25519(key)),
name: RepoName::new(name).unwrap(),
}
}
fn mkpl(packs: &[Hash]) -> (Vec<u8>, Hash) {
let bytes = mkit_core::transfer::encode_packlist(None, packs).unwrap();
let id = hash(&bytes);
(bytes, id)
}
fn mkpl_prev(prev: Hash, packs: &[Hash]) -> (Vec<u8>, Hash) {
let bytes = mkit_core::transfer::encode_packlist(Some(prev), packs).unwrap();
let id = hash(&bytes);
(bytes, id)
}
#[derive(Clone, Default)]
struct AdmissionSpy {
seen: Arc<Mutex<Vec<Procedure>>>,
reservation: Option<String>,
challenge: bool,
}
impl Admission for AdmissionSpy {
fn is_default(&self) -> bool {
true
}
async fn admit(&self, input: &AdmissionInput<'_>) -> Result<AdmissionDecision, ServerError> {
self.seen.lock().unwrap().push(input.op.procedure());
if self.challenge {
return Ok(AdmissionDecision::challenge(
vec![crate::pipeline::Challenge {
scheme: "payment".into(),
value: "id=\"c1\"".into(),
}],
"pay to write",
));
}
Ok(AdmissionDecision::Allow {
charges: Vec::new(),
reservation: self.reservation.clone(),
response_headers: Vec::new(),
external_ref: None,
})
}
}
fn spy_hooks(admission: AdmissionSpy) -> Hooks<OpenAuthorizer, AdmissionSpy> {
let defaults = Hooks::new();
Hooks {
authorizer: defaults.authorizer,
admission,
pre_receive: defaults.pre_receive,
receipts: defaults.receipts,
outcomes: defaults.outcomes,
}
}
#[derive(Clone)]
struct SharedKv(Arc<MemoryKv>);
impl SharedKv {
fn read(&self, p: &Partition, key: &Key) -> Option<Value> {
block_on(NamespaceStore::get(self, p, key)).unwrap()
}
}
impl NamespaceStore for SharedKv {
fn capabilities(&self) -> StoreCapabilities {
self.0.capabilities()
}
async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
self.0.get(p, key).await
}
async fn get_many(
&self,
p: &Partition,
keys: &[Key],
) -> Result<Vec<Option<Value>>, StoreError> {
self.0.get_many(p, keys).await
}
async fn scan(
&self,
p: &Partition,
start: &Key,
end: &Key,
after: Option<&Cursor>,
limit: u32,
) -> Result<ScanPage, StoreError> {
self.0.scan(p, start, end, after, limit).await
}
async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
self.0.apply(p, batch).await
}
async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
self.0.stats(p).await
}
async fn probe(&self) -> Result<(), StoreError> {
self.0.probe().await
}
}
#[derive(Clone)]
struct BumpingKv {
inner: SharedKv,
key: Key,
value: Value,
armed: Arc<AtomicBool>,
}
impl BumpingKv {
fn packmap(inner: SharedKv, repo: &RepoId, bumped: Hash) -> Self {
Self {
inner,
key: keys::ref_key(&repo.name, PACKMAP),
value: codec::encode_ref_id(&bumped),
armed: Arc::new(AtomicBool::new(false)),
}
}
fn arm(&self) {
self.armed.store(true, Ordering::SeqCst);
}
}
impl NamespaceStore for BumpingKv {
fn capabilities(&self) -> StoreCapabilities {
self.inner.capabilities()
}
async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
self.inner.get(p, key).await
}
async fn get_many(
&self,
p: &Partition,
keys: &[Key],
) -> Result<Vec<Option<Value>>, StoreError> {
self.inner.get_many(p, keys).await
}
async fn scan(
&self,
p: &Partition,
start: &Key,
end: &Key,
after: Option<&Cursor>,
limit: u32,
) -> Result<ScanPage, StoreError> {
self.inner.scan(p, start, end, after, limit).await
}
async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
if self.armed.load(Ordering::SeqCst)
&& batch
.writes
.iter()
.any(|w| matches!(w, Write::Put(k, _) if *k == self.key))
{
self.armed.store(false, Ordering::SeqCst);
self.inner
.apply(p, Batch::new().put(self.key.clone(), self.value.clone()))
.await?;
}
self.inner.apply(p, batch).await
}
async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
self.inner.stats(p).await
}
async fn probe(&self) -> Result<(), StoreError> {
self.inner.probe().await
}
}
#[derive(Clone)]
struct CountingKv {
inner: SharedKv,
reads: Arc<AtomicUsize>,
}
impl CountingKv {
fn wrap(inner: SharedKv) -> Self {
Self {
inner,
reads: Arc::new(AtomicUsize::new(0)),
}
}
fn member_reads(&self) -> usize {
self.reads.load(Ordering::SeqCst)
}
fn count(&self, key: &Key) {
if key.as_bytes().starts_with(b"m\0") {
self.reads.fetch_add(1, Ordering::SeqCst);
}
}
}
impl NamespaceStore for CountingKv {
fn capabilities(&self) -> StoreCapabilities {
self.inner.capabilities()
}
async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
self.count(key);
self.inner.get(p, key).await
}
async fn get_many(
&self,
p: &Partition,
keys: &[Key],
) -> Result<Vec<Option<Value>>, StoreError> {
for key in keys {
self.count(key);
}
self.inner.get_many(p, keys).await
}
async fn scan(
&self,
p: &Partition,
start: &Key,
end: &Key,
after: Option<&Cursor>,
limit: u32,
) -> Result<ScanPage, StoreError> {
self.count(start);
self.inner.scan(p, start, end, after, limit).await
}
async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
self.inner.apply(p, batch).await
}
async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
self.inner.stats(p).await
}
async fn probe(&self) -> Result<(), StoreError> {
self.inner.probe().await
}
}
#[derive(Clone)]
struct RecordingKv {
inner: SharedKv,
batches: Arc<Mutex<Vec<Batch>>>,
}
impl RecordingKv {
fn wrap(inner: SharedKv) -> Self {
Self {
inner,
batches: Arc::new(Mutex::new(Vec::new())),
}
}
fn take(&self) -> Vec<Batch> {
std::mem::take(&mut *self.batches.lock().unwrap())
}
}
impl NamespaceStore for RecordingKv {
fn capabilities(&self) -> StoreCapabilities {
self.inner.capabilities()
}
async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
self.inner.get(p, key).await
}
async fn get_many(
&self,
p: &Partition,
keys: &[Key],
) -> Result<Vec<Option<Value>>, StoreError> {
self.inner.get_many(p, keys).await
}
async fn scan(
&self,
p: &Partition,
start: &Key,
end: &Key,
after: Option<&Cursor>,
limit: u32,
) -> Result<ScanPage, StoreError> {
self.inner.scan(p, start, end, after, limit).await
}
async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
self.batches.lock().unwrap().push(batch.clone());
self.inner.apply(p, batch).await
}
async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
self.inner.stats(p).await
}
async fn probe(&self) -> Result<(), StoreError> {
self.inner.probe().await
}
}
fn multi_pipeline<H: HookSet>(
allowlist: impl IntoIterator<Item = Namespace>,
sharding: Sharding,
hooks: H,
) -> (
Pipeline<MemoryBlobStore, SharedKv, H>,
MemoryBlobStore,
SharedKv,
) {
let clock = Arc::new(ManualClock::new(T0));
let kv = SharedKv(Arc::new(MemoryKv::with_clock(clock.clone())));
let (pipe, blobs) = multi_pipeline_on(allowlist, sharding, hooks, kv.clone(), clock);
(pipe, blobs, kv)
}
fn multi_pipeline_on<H: HookSet, N: NamespaceStore>(
allowlist: impl IntoIterator<Item = Namespace>,
sharding: Sharding,
hooks: H,
kv: N,
clock: Arc<ManualClock>,
) -> (Pipeline<MemoryBlobStore, N, H>, MemoryBlobStore) {
let mut cfg = PipelineConfig::new(
Addressing::Multi(
MultiAddressing::new()
.with_namespace_policy(NamespacePolicy::Allowlist(allowlist.into_iter().collect())),
),
AuthMode::TransportIdentity,
upload_limits(),
);
cfg.sharding = sharding;
cfg.ticket_keys = Some(TicketKeys::new(vec![("test".into(), [7; 32])]).unwrap());
let blobs = MemoryBlobStore::default();
let pipe = Pipeline::new(blobs.clone(), kv, hooks, cfg, clock, Arc::new(NoopMetrics)).unwrap();
(pipe, blobs)
}
fn peer_run<B: MultipartBlobStore, N: NamespaceStore, H: HookSet>(
pipe: &Pipeline<B, N, H>,
peer: [u8; 32],
repository: &str,
bodies: impl IntoIterator<Item = Option<Body>>,
) -> (SessionEnd, Vec<SshFrame>) {
let mut cfg = SessionConfig::new(SERVER_ID);
cfg.repository = Some(repository.to_owned());
let mut src = script(bodies);
let mut sink = VecSink::default();
let end = block_on(serve_session(
pipe,
Principal::TransportPeer { ed25519: peer },
&mut src,
&mut sink,
&cfg,
));
let first = sink.frames.first().expect("a HelloResponse");
assert!(
matches!(first.body, Some(Body::HelloResponse(_))),
"got {:?}",
first.body
);
(end, sink.frames[1..].to_vec())
}
fn root(ns: &NamespaceKey) -> Partition {
Partition::Namespace(ns.clone())
}
fn member<N: NamespaceStore>(kv: &N, partition: &Partition, name: &RepoName, pack: &Hash) -> bool {
block_on(kv.get(partition, &keys::membership(name, pack)))
.unwrap()
.is_some()
}
fn o_rows<N: NamespaceStore>(kv: &N, partition: &Partition) -> Vec<(Key, Value)> {
block_on(kv.scan(
partition,
&Key::new(b"o".to_vec()),
&Key::new(b"p".to_vec()),
None,
64,
))
.unwrap()
.entries
}
#[test]
fn implicit_packmap_consumes_pending_into_membership() {
let spy = AdmissionSpy::default();
let seen = spy.seen.clone();
let (pipe, _, kv) = multi_pipeline(
[Namespace::Ed25519(OWNER)],
Sharding::Single,
spy_hooks(spy),
);
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (mkpl_bytes, mkpl_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&mkpl_id, Some(mkpl_bytes.len() as u64)),
chunk(&mkpl_id, Some(0), &mkpl_bytes, true),
exists(&data_id),
update(PACKMAP, &mkpl_id, Some(RefExpectation::Missing), None),
update(BRANCH, &[7; 32], Some(RefExpectation::Missing), None),
exists(&data_id),
download(Some(&data_id)),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UploadPackResponse(_))));
assert!(matches!(frames[1].body, Some(Body::UploadPackResponse(_))));
let Some(Body::PackExistsResponse(resp)) = &frames[2].body else {
panic!("expected PackExistsResponse, got {:?}", frames[2].body);
};
assert_eq!(resp.exists, Some(false));
assert!(matches!(frames[3].body, Some(Body::UpdateRefResponse(_))));
assert!(matches!(frames[4].body, Some(Body::UpdateRefResponse(_))));
let Some(Body::PackExistsResponse(resp)) = &frames[5].body else {
panic!("expected PackExistsResponse, got {:?}", frames[5].body);
};
assert_eq!(resp.exists, Some(true));
assert_eq!(downloaded(&frames[6..], &data_id), data);
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&data_id
));
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&mkpl_id
));
assert!(o_rows(&kv, &root(&repo_id.namespace)).is_empty());
assert_eq!(
*seen.lock().unwrap(),
[
Procedure::UploadPack,
Procedure::UploadPack,
Procedure::UpdateRef
]
);
}
#[test]
fn implicit_packmap_queues_relay_rows_under_d34() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::D34, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (mkpl_bytes, mkpl_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&mkpl_id, Some(mkpl_bytes.len() as u64)),
chunk(&mkpl_id, Some(0), &mkpl_bytes, true),
update(PACKMAP, &mkpl_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
for f in &frames {
assert!(
matches!(
f.body,
Some(Body::UploadPackResponse(_) | Body::UpdateRefResponse(_))
),
"got {f:?}"
);
}
let source = D34Shards.ref_shard(&repo_id, PACKMAP);
assert!(member(&kv, &source, &repo_id.name, &data_id));
assert!(member(&kv, &source, &repo_id.name, &mkpl_id));
let seq = kv.read(&source, &keys::outbox_sequence());
let seq = seq.map(|v| codec::decode_u64(&v).unwrap());
assert!(seq.is_some_and(|s| s >= 1));
assert!(kv.read(&source, &keys::relay(1)).is_some());
let mut leftovers = o_rows(&kv, &source);
leftovers.retain(|(k, _)| !k.as_bytes().starts_with(b"or") && !k.as_bytes().starts_with(b"os"));
assert!(leftovers.is_empty(), "{leftovers:?}");
}
#[test]
fn implicit_cas_conflict_keeps_pending() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (mkpl_bytes, mkpl_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&mkpl_id, Some(mkpl_bytes.len() as u64)),
chunk(&mkpl_id, Some(0), &mkpl_bytes, true),
update(
PACKMAP,
&mkpl_id,
Some(RefExpectation::Match),
Some(&[9; 32]),
),
update(PACKMAP, &mkpl_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UploadPackResponse(_))));
assert!(matches!(frames[1].body, Some(Body::UploadPackResponse(_))));
let conflict = error(&frames[2]);
assert!(
conflict
.code
.is_some_and(|c| c == ErrorCode::InvalidRequest)
);
assert!(matches!(frames[3].body, Some(Body::UpdateRefResponse(_))));
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&data_id
));
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&mkpl_id
));
}
#[test]
fn head_write_consumes_nothing() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (mkpl_bytes, mkpl_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&mkpl_id, Some(mkpl_bytes.len() as u64)),
chunk(&mkpl_id, Some(0), &mkpl_bytes, true),
update(BRANCH, &[7; 32], Some(RefExpectation::Missing), None),
exists(&data_id),
update(PACKMAP, &mkpl_id, Some(RefExpectation::Missing), None),
exists(&data_id),
],
);
assert_eq!(end, SessionEnd::Clean);
for f in &frames[..3] {
assert!(
matches!(
f.body,
Some(Body::UploadPackResponse(_) | Body::UpdateRefResponse(_))
),
"got {f:?}"
);
}
let Some(Body::PackExistsResponse(resp)) = &frames[3].body else {
panic!("expected PackExistsResponse, got {:?}", frames[3].body);
};
assert_eq!(resp.exists, Some(false), "a head write consumes nothing");
assert!(matches!(frames[4].body, Some(Body::UpdateRefResponse(_))));
let Some(Body::PackExistsResponse(resp)) = &frames[5].body else {
panic!("expected PackExistsResponse, got {:?}", frames[5].body);
};
assert_eq!(resp.exists, Some(true));
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&data_id
));
}
#[test]
fn eighth_pending_upload_is_refused_in_frame_sync() {
let (pipe, blobs, _) =
multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let mut bodies = Vec::new();
for i in 0..7u8 {
let bytes = pack_bytes(64, i + 1);
let id = hash(&bytes);
bodies.push(upload_header(&id, Some(bytes.len() as u64)));
bodies.push(chunk(&id, Some(0), &bytes, true));
}
let eighth = pack_bytes(64, 9);
let eighth_id = hash(&eighth);
bodies.push(upload_header(&eighth_id, Some(eighth.len() as u64)));
bodies.push(chunk(&eighth_id, Some(0), &eighth, true));
bodies.push(exists(&eighth_id));
let (end, frames) = peer_run(&pipe, OWNER, &repo, bodies);
assert_eq!(end, SessionEnd::Clean);
for f in &frames[..7] {
assert!(
matches!(f.body, Some(Body::UploadPackResponse(_))),
"got {f:?}"
);
}
assert_error(
&frames[7],
ErrorCode::InvalidRequest,
"too many packs uploaded before a packmap update",
);
let Some(Body::PackExistsResponse(resp)) = &frames[8].body else {
panic!("expected PackExistsResponse, got {:?}", frames[8].body);
};
assert_eq!(resp.exists, Some(false));
assert!(!blob_present(&blobs, eighth_id));
}
#[test]
fn reserving_admission_fails_the_upload_closed() {
let spy = AdmissionSpy {
reservation: Some("r:1".into()),
..AdmissionSpy::default()
};
let (pipe, blobs, kv) = multi_pipeline(
[Namespace::Ed25519(OWNER)],
Sharding::Single,
spy_hooks(spy),
);
let repo = peer_repo(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
exists(&data_id),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
"payment required: use mkit+https",
);
let Some(Body::PackExistsResponse(resp)) = &frames[1].body else {
panic!("expected PackExistsResponse, got {:?}", frames[1].body);
};
assert_eq!(resp.exists, Some(false));
assert!(!blob_present(&blobs, data_id));
let ns = NamespaceKey::from_namespace(&Namespace::Ed25519(OWNER));
let rows: Vec<_> = o_rows(&kv, &root(&ns))
.into_iter()
.filter(|(key, _)| key.as_bytes().starts_with(b"o\0"))
.collect();
assert_eq!(rows.len(), 1, "{rows:?}");
assert!(matches!(
crate::store::codec::decode_reservation(&rows[0].1).unwrap(),
crate::store::codec::ReservationV1::Aborted { .. }
));
}
fn challenging_hooks() -> Hooks<OpenAuthorizer, AdmissionSpy> {
spy_hooks(AdmissionSpy {
challenge: true,
..AdmissionSpy::default()
})
}
#[test]
fn a_payment_challenge_answers_use_https() {
use mkit_core::protocol::{RefWriteCondition, TransportError, is_retryable};
const FRAME: &str = "payment required: use mkit+https";
let (pipe, blobs, kv) = multi_pipeline(
[Namespace::Ed25519(OWNER)],
Sharding::Single,
challenging_hooks(),
);
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
update(BRANCH, &[7; 32], Some(RefExpectation::Missing), None),
update(BRANCH, &[7; 32], Some(RefExpectation::Any), None),
update(
BRANCH,
&[8; 32],
Some(RefExpectation::Match),
Some(&[7; 32]),
),
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
exists(&data_id),
],
);
assert_eq!(end, SessionEnd::Clean);
for f in &frames[..4] {
assert_error(f, ErrorCode::InvalidRequest, FRAME);
}
let Some(Body::PackExistsResponse(resp)) = &frames[4].body else {
panic!("expected PackExistsResponse, got {:?}", frames[4].body);
};
assert_eq!(resp.exists, Some(false));
assert!(!blob_present(&blobs, data_id));
assert!(
kv.read(
&root(&repo_id.namespace),
&keys::ref_key(&repo_id.name, BRANCH)
)
.is_none()
);
for condition in [RefWriteCondition::Match([7; 32]), RefWriteCondition::Any] {
let err = mkit_rpc::map_update_ref_error(error(&frames[0]).clone(), condition, "ssh");
assert!(
matches!(&err, TransportError::RemoteError(m) if m.contains(FRAME)),
"{err:?}"
);
assert!(!is_retryable(&err));
}
}
#[test]
fn packmap_in_a_new_session_names_nothing() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (mkpl_bytes, mkpl_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&mkpl_id, Some(mkpl_bytes.len() as u64)),
chunk(&mkpl_id, Some(0), &mkpl_bytes, true),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_eq!(frames.len(), 2, "two uploads, two responses");
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
update(PACKMAP, &mkpl_id, Some(RefExpectation::Missing), None),
exists(&mkpl_id),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_error(
&frames[0],
ErrorCode::InvalidRequest,
"packmap names packs not uploaded to this repository",
);
let Some(Body::PackExistsResponse(resp)) = &frames[1].body else {
panic!("expected PackExistsResponse, got {:?}", frames[1].body);
};
assert_eq!(resp.exists, Some(false));
assert!(
kv.read(
&root(&repo_id.namespace),
&keys::ref_key(&repo_id.name, PACKMAP)
)
.is_none()
);
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&data_id
));
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&mkpl_id
));
assert!(o_rows(&kv, &root(&repo_id.namespace)).is_empty());
}
#[test]
fn packmap_listing_an_unknown_pack_is_refused() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let ghost = [0x99; 32];
let (mkpl_bytes, mkpl_id) = mkpl(&[ghost]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&mkpl_id, Some(mkpl_bytes.len() as u64)),
chunk(&mkpl_id, Some(0), &mkpl_bytes, true),
update(PACKMAP, &mkpl_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UploadPackResponse(_))));
assert_error(
&frames[1],
ErrorCode::InvalidRequest,
"packmap names packs not uploaded to this repository",
);
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&mkpl_id
));
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&ghost
));
}
#[test]
fn private_repository_reads_are_not_found_over_transport_identity() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[update(
"refs/heads/main",
&[0x55; 32],
Some(RefExpectation::Missing),
None,
)],
);
assert_eq!(end, SessionEnd::Clean);
assert!(
matches!(frames[0].body, Some(Body::UpdateRefResponse(_))),
"{:?}",
frames[0].body
);
let (_, public) = peer_run(&pipe, OTHER, &repo, [read_ref("refs/heads/main")]);
assert!(matches!(public[0].body, Some(Body::ReadRefResponse(_))));
let batch = Batch::new().put(
keys::repo_visibility(&repo_id.name),
codec::encode_repo_visibility(&codec::RepoVisibilityV1 {
visibility: codec::StoredVisibility::Private,
last_created_ms: 0,
last_statement_id: None,
changed_ms: None,
}),
);
let coordinator = root(&repo_id.namespace);
assert_eq!(
block_on(kv.apply(&coordinator, batch)).unwrap(),
BatchOutcome::Committed
);
let missing = peer_repo(OWNER, "never-created");
let (_, absent) = peer_run(&pipe, OTHER, &missing, [read_ref("refs/heads/main")]);
assert!(matches!(absent[0].body, Some(Body::Error(_))));
let (_, stranger) = peer_run(&pipe, OTHER, &repo, [read_ref("refs/heads/main")]);
assert_eq!(stranger, absent, "private reads like a missing repository");
let (_, owner) = peer_run(&pipe, OWNER, &repo, [read_ref("refs/heads/main")]);
assert_eq!(owner, absent, "owner private reads fail closed over enc");
}
#[test]
fn packmap_prev_must_equal_the_replaced_value() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (m_bytes, m_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&m_id, Some(m_bytes.len() as u64)),
chunk(&m_id, Some(0), &m_bytes, true),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_eq!(frames.len(), 2, "two uploads, two responses");
let second = b"second pack bytes".to_vec();
let second_id = hash(&second);
let (n_bytes, n_id) = mkpl_prev(m_id, &[second_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&second_id, Some(second.len() as u64)),
chunk(&second_id, Some(0), &second, true),
upload_header(&n_id, Some(n_bytes.len() as u64)),
chunk(&n_id, Some(0), &n_bytes, true),
update(PACKMAP, &n_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UploadPackResponse(_))));
assert!(matches!(frames[1].body, Some(Body::UploadPackResponse(_))));
assert_error(
&frames[2],
ErrorCode::InvalidRequest,
"packmap names packs not uploaded to this repository",
);
assert!(
kv.read(
&root(&repo_id.namespace),
&keys::ref_key(&repo_id.name, PACKMAP)
)
.is_none()
);
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&second_id
));
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&n_id
));
assert!(o_rows(&kv, &root(&repo_id.namespace)).is_empty());
}
#[test]
fn packmap_prev_member_but_not_current_is_refused() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (n0_bytes, n0_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&n0_id, Some(n0_bytes.len() as u64)),
chunk(&n0_id, Some(0), &n0_bytes, true),
update(PACKMAP, &n0_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[2].body, Some(Body::UpdateRefResponse(_))));
let n1_id = hash(b"a member that is not the current packmap");
block_on(kv.apply(
&root(&repo_id.namespace),
Batch::new().put(keys::membership(&repo_id.name, &n1_id), Value::default()),
))
.unwrap();
let second = b"second pack bytes".to_vec();
let second_id = hash(&second);
let (n3_bytes, n3_id) = mkpl_prev(n1_id, &[second_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&second_id, Some(second.len() as u64)),
chunk(&second_id, Some(0), &second, true),
upload_header(&n3_id, Some(n3_bytes.len() as u64)),
chunk(&n3_id, Some(0), &n3_bytes, true),
update(
PACKMAP,
&n3_id,
Some(RefExpectation::Match),
Some(&n0_id[..]),
),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_error(
&frames[2],
ErrorCode::InvalidRequest,
"packmap names packs not uploaded to this repository",
);
let stored = kv
.read(
&root(&repo_id.namespace),
&keys::ref_key(&repo_id.name, PACKMAP),
)
.unwrap();
assert_eq!(codec::decode_ref_id(&stored).unwrap(), n0_id);
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&second_id
));
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&n3_id
));
}
#[test]
fn packmap_prev_equal_to_the_replaced_value_is_accepted() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (l2_bytes, l2_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&l2_id, Some(l2_bytes.len() as u64)),
chunk(&l2_id, Some(0), &l2_bytes, true),
update(PACKMAP, &l2_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[2].body, Some(Body::UpdateRefResponse(_))));
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&l2_id
));
let second = b"second pack bytes".to_vec();
let second_id = hash(&second);
let (m2_bytes, m2_id) = mkpl_prev(l2_id, &[second_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&second_id, Some(second.len() as u64)),
chunk(&second_id, Some(0), &second, true),
upload_header(&m2_id, Some(m2_bytes.len() as u64)),
chunk(&m2_id, Some(0), &m2_bytes, true),
update(
PACKMAP,
&m2_id,
Some(RefExpectation::Match),
Some(&l2_id[..]),
),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[2].body, Some(Body::UpdateRefResponse(_))));
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&m2_id
));
let third = b"third pack bytes".to_vec();
let third_id = hash(&third);
let (m3_bytes, m3_id) = mkpl_prev(m2_id, &[third_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&third_id, Some(third.len() as u64)),
chunk(&third_id, Some(0), &third, true),
upload_header(&m3_id, Some(m3_bytes.len() as u64)),
chunk(&m3_id, Some(0), &m3_bytes, true),
update(PACKMAP, &m3_id, Some(RefExpectation::Any), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[2].body, Some(Body::UpdateRefResponse(_))));
let stored = kv
.read(
&root(&repo_id.namespace),
&keys::ref_key(&repo_id.name, PACKMAP),
)
.unwrap();
assert_eq!(codec::decode_ref_id(&stored).unwrap(), m3_id);
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&m3_id
));
}
#[test]
fn packmap_any_guard_conflict_keeps_pending() {
let clock = Arc::new(ManualClock::new(T0));
let kv = SharedKv(Arc::new(MemoryKv::with_clock(clock.clone())));
let (pipe, _) = multi_pipeline_on(
[Namespace::Ed25519(OWNER)],
Sharding::Single,
Hooks::new(),
kv.clone(),
clock,
);
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (l2_bytes, l2_id) = mkpl(&[data_id]);
let (end, _) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&l2_id, Some(l2_bytes.len() as u64)),
chunk(&l2_id, Some(0), &l2_bytes, true),
update(PACKMAP, &l2_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
let bumped = hash(b"concurrent packmap");
let bumping = BumpingKv::packmap(kv.clone(), &repo_id, bumped);
let (pipe2, _) = multi_pipeline_on(
[Namespace::Ed25519(OWNER)],
Sharding::Single,
Hooks::new(),
bumping.clone(),
Arc::new(ManualClock::new(T0)),
);
bumping.arm();
let second = b"second pack bytes".to_vec();
let second_id = hash(&second);
let (m2_bytes, m2_id) = mkpl_prev(l2_id, &[second_id]);
let (m3_bytes, m3_id) = mkpl_prev(bumped, &[second_id]);
let (end, frames) = peer_run(
&pipe2,
OWNER,
&repo,
[
upload_header(&second_id, Some(second.len() as u64)),
chunk(&second_id, Some(0), &second, true),
upload_header(&m2_id, Some(m2_bytes.len() as u64)),
chunk(&m2_id, Some(0), &m2_bytes, true),
update(PACKMAP, &m2_id, Some(RefExpectation::Any), None),
upload_header(&m3_id, Some(m3_bytes.len() as u64)),
chunk(&m3_id, Some(0), &m3_bytes, true),
update(
PACKMAP,
&m3_id,
Some(RefExpectation::Match),
Some(&bumped[..]),
),
],
);
assert_eq!(end, SessionEnd::Clean);
let e = error(&frames[2]);
assert!(e.code.is_some_and(|c| c == ErrorCode::InvalidRequest));
assert_eq!(
e.message.as_deref(),
Some("ref update conflict: expectation does not match current ref value")
);
assert!(matches!(frames[4].body, Some(Body::UpdateRefResponse(_))));
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&second_id
));
assert!(member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&m3_id
));
let stored = kv
.read(
&root(&repo_id.namespace),
&keys::ref_key(&repo_id.name, PACKMAP),
)
.unwrap();
assert_eq!(codec::decode_ref_id(&stored).unwrap(), m3_id);
}
#[test]
fn packmap_listing_a_pending_packlist_is_refused() {
let (pipe, _, kv) = multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let ghost = hash(b"never uploaded");
let (n1_bytes, n1_id) = mkpl(&[ghost]);
let (n2_bytes, n2_id) = mkpl(&[n1_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&n1_id, Some(n1_bytes.len() as u64)),
chunk(&n1_id, Some(0), &n1_bytes, true),
upload_header(&n2_id, Some(n2_bytes.len() as u64)),
chunk(&n2_id, Some(0), &n2_bytes, true),
update(PACKMAP, &n2_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UploadPackResponse(_))));
assert!(matches!(frames[1].body, Some(Body::UploadPackResponse(_))));
assert_error(
&frames[2],
ErrorCode::InvalidRequest,
"packmap names packs not uploaded to this repository",
);
assert!(
kv.read(
&root(&repo_id.namespace),
&keys::ref_key(&repo_id.name, PACKMAP)
)
.is_none()
);
assert!(!member(
&kv,
&root(&repo_id.namespace),
&repo_id.name,
&n2_id
));
}
#[test]
fn packmap_listing_too_many_packs_is_refused_without_membership_reads() {
let clock = Arc::new(ManualClock::new(T0));
let kv = SharedKv(Arc::new(MemoryKv::with_clock(clock.clone())));
let counting = CountingKv::wrap(kv.clone());
let (pipe, _) = multi_pipeline_on(
[Namespace::Ed25519(OWNER)],
Sharding::Single,
Hooks::new(),
counting.clone(),
clock,
);
let repo = peer_repo(OWNER, "room-a");
let repo_id = peer_repo_id(OWNER, "room-a");
let packs: Vec<Hash> = (0..1025u16).map(|i| hash(&i.to_be_bytes())).collect();
let (node_bytes, node_id) = mkpl(&packs);
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&node_id, Some(node_bytes.len() as u64)),
chunk(&node_id, Some(0), &node_bytes, true),
update(PACKMAP, &node_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_error(
&frames[1],
ErrorCode::InvalidRequest,
"packmap names packs not uploaded to this repository",
);
assert_eq!(
counting.member_reads(),
0,
"the cap refuses before any membership read"
);
assert!(
kv.read(
&root(&repo_id.namespace),
&keys::ref_key(&repo_id.name, PACKMAP)
)
.is_none()
);
}
#[test]
fn empty_pending_consume_plans_no_membership_work() {
fn shape(batch: &Batch) -> Vec<u8> {
let mut tags: Vec<u8> = batch
.writes
.iter()
.map(|w| match w {
Write::Put(k, _) | Write::Delete(k) => k.as_bytes()[0],
})
.collect();
tags.sort_unstable();
tags
}
let clock = Arc::new(ManualClock::new(T0));
let kv = SharedKv(Arc::new(MemoryKv::with_clock(clock.clone())));
let recording = RecordingKv::wrap(kv.clone());
let (pipe, _) = multi_pipeline_on(
[Namespace::Ed25519(OWNER)],
Sharding::D34,
Hooks::new(),
recording.clone(),
clock,
);
let repo = peer_repo(OWNER, "room-a");
let (data, data_id) = valid_pack();
let (m1_bytes, m1_id) = mkpl(&[data_id]);
let (end, _) = peer_run(
&pipe,
OWNER,
&repo,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&m1_id, Some(m1_bytes.len() as u64)),
chunk(&m1_id, Some(0), &m1_bytes, true),
update(PACKMAP, &m1_id, Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
recording.take();
let (end, frames) = peer_run(
&pipe,
OWNER,
&repo,
[
update(
PACKMAP,
&m1_id,
Some(RefExpectation::Match),
Some(&m1_id[..]),
),
update(BRANCH, &data_id, Some(RefExpectation::Any), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(
frames
.iter()
.all(|f| matches!(f.body, Some(Body::UpdateRefResponse(_))))
);
let batches = recording.take();
assert!(batches.len() >= 2, "one batch per ref write: {batches:?}");
let (packmap_batch, head_batch) = (&batches[batches.len() - 2], &batches[batches.len() - 1]);
assert_eq!(
shape(packmap_batch),
shape(head_batch),
"an empty consume plans exactly the ref write's own work"
);
assert!(
!packmap_batch.writes.iter().any(|w| match w {
Write::Put(k, _) | Write::Delete(k) => k.as_bytes().starts_with(b"m\0"),
}),
"no membership puts: {batches:?}"
);
}
#[test]
fn foreign_and_unlisted_namespaces_write_nothing() {
let foreign = peer_repo(OTHER, "room-b");
let (pipe, blobs, kv) =
multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let (data, data_id) = valid_pack();
let (end, frames) = peer_run(
&pipe,
OWNER,
&foreign,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
update(BRANCH, &[7; 32], Some(RefExpectation::Missing), None),
exists(&data_id),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_error(&frames[0], ErrorCode::InvalidRequest, "write not permitted");
assert_error(&frames[1], ErrorCode::InvalidRequest, "write not permitted");
let Some(Body::PackExistsResponse(resp)) = &frames[2].body else {
panic!("expected PackExistsResponse, got {:?}", frames[2].body);
};
assert_eq!(resp.exists, Some(false));
let room_b = RepoName::new("room-b").unwrap();
let other_ns = peer_repo_id(OTHER, "room-b").namespace;
assert!(!blob_present(&blobs, data_id));
assert!(
kv.read(&root(&other_ns), &keys::repo_record(&room_b))
.is_none(),
"the denied write created no repo record"
);
assert!(
kv.read(&root(&other_ns), &keys::repo_known(&room_b))
.is_none(),
"the denied write created no repo-known marker"
);
let own = peer_repo(OWNER, "room-a");
let (pipe, blobs, kv) =
multi_pipeline([Namespace::Ed25519(OTHER)], Sharding::Single, Hooks::new());
let (end, frames) = peer_run(
&pipe,
OWNER,
&own,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
update(BRANCH, &[7; 32], Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_error(&frames[0], ErrorCode::InvalidRequest, "write not permitted");
assert_error(&frames[1], ErrorCode::InvalidRequest, "write not permitted");
assert!(!blob_present(&blobs, data_id));
let room_a = RepoName::new("room-a").unwrap();
let own_ns = peer_repo_id(OWNER, "room-a").namespace;
assert!(
kv.read(&root(&own_ns), &keys::repo_record(&room_a))
.is_none()
);
}
#[test]
fn repositories_see_only_their_own_packs() {
let (pipe, blobs, kv) =
multi_pipeline([Namespace::Ed25519(OWNER)], Sharding::Single, Hooks::new());
let room_a = peer_repo(OWNER, "room-a");
let room_b = peer_repo(OWNER, "room-b");
let (data, data_id) = valid_pack();
let (mkpl_bytes, mkpl_id) = mkpl(&[data_id]);
let (end, frames) = peer_run(
&pipe,
OWNER,
&room_a,
[
upload_header(&data_id, Some(data.len() as u64)),
chunk(&data_id, Some(0), &data, true),
upload_header(&mkpl_id, Some(mkpl_bytes.len() as u64)),
chunk(&mkpl_id, Some(0), &mkpl_bytes, true),
update(PACKMAP, &mkpl_id, Some(RefExpectation::Missing), None),
update(BRANCH, &[7; 32], Some(RefExpectation::Missing), None),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_eq!(frames.len(), 4);
let (end, frames) = peer_run(
&pipe,
OWNER,
&room_b,
[
update(BRANCH, &[8; 32], Some(RefExpectation::Missing), None),
exists(&data_id),
download(Some(&data_id)),
read_ref(BRANCH),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::UpdateRefResponse(_))));
let Some(Body::PackExistsResponse(resp)) = &frames[1].body else {
panic!("expected PackExistsResponse, got {:?}", frames[1].body);
};
assert_eq!(resp.exists, Some(false), "A's pack must not leak to B");
assert_error(&frames[2], ErrorCode::KeyNotFound, "pack not found");
let Some(Body::ReadRefResponse(resp)) = &frames[3].body else {
panic!("expected ReadRefResponse, got {:?}", frames[3].body);
};
assert_eq!(resp.object_id.as_deref(), Some(&[8; 32][..]));
assert!(blob_present(&blobs, data_id));
let room_b_id = peer_repo_id(OWNER, "room-b");
assert!(!member(
&kv,
&root(&room_b_id.namespace),
&room_b_id.name,
&data_id
));
}
#[cfg(feature = "fs")]
mod fs {
use std::path::{Path, PathBuf};
use mkit_core::protocol::Transport as _;
use mkit_transport_file::FileTransport;
use super::*;
use crate::fs::{FsBlobStore, FsLayoutStore};
fn fs_pipe(root: &Path) -> Pipeline<FsBlobStore, FsLayoutStore> {
let clock = Arc::new(ManualClock::new(T0));
let meta = FsLayoutStore::new(root, &repo()).with_clock(clock);
pipeline(FsBlobStore::new(root), meta, AuthMode::TransportIdentity)
}
fn files(root: &Path) -> Vec<PathBuf> {
fn walk(root: &Path, dir: &Path, out: &mut Vec<PathBuf>) {
for entry in std::fs::read_dir(dir).unwrap() {
let path = entry.unwrap().path();
if path.is_dir() {
walk(root, &path, out);
} else {
out.push(path.strip_prefix(root).unwrap().to_owned());
}
}
}
let mut out = Vec::new();
walk(root, root, &mut out);
out.sort();
out
}
#[test]
fn golden_session_matches_mkit_serve_over_fs_stores() {
let td = tempfile::tempdir().unwrap();
let pipe = fs_pipe(td.path());
let out = replay_golden(&pipe);
assert_same_frames(&out, GOLDEN_OUT);
assert_eq!(out, GOLDEN_OUT);
let tx = FileTransport::new(td.path());
assert_eq!(tx.read_ref("refs/heads/main").unwrap(), Some([0x44; 32]));
assert_eq!(tx.read_ref("refs/heads/any").unwrap(), Some([0x77; 32]));
}
#[test]
fn golden_session_2_matches_mkit_serve_over_fs_stores() {
let td = tempfile::tempdir().unwrap();
let tx = FileTransport::new(td.path());
let g = golden2();
for (name, id) in &g.refs {
tx.update_ref(name, RefWriteCondition::Any, id).unwrap();
}
tx.upload_pack(&g.big, &PackKey::new(hash(&g.big))).unwrap();
std::fs::write(td.path().join("refs/heads/README"), b"not a ref\n").unwrap();
let out = replay_golden_2(&fs_pipe(td.path()));
assert_same_frames(&out, GOLDEN_2_OUT);
assert_eq!(out, GOLDEN_2_OUT);
}
#[test]
fn stray_file_under_refs_is_skipped_by_listings() {
let td = tempfile::tempdir().unwrap();
let tx = FileTransport::new(td.path());
tx.update_ref("refs/heads/main", RefWriteCondition::Any, &[0x11; 32])
.unwrap();
std::fs::write(td.path().join("refs/heads/README"), b"not a ref\n").unwrap();
let pipe = fs_pipe(td.path());
let (end, frames) = run(
&pipe,
[
list_refs(Some("refs/heads/")),
list_refs(Some("refs/tags/")),
list_refs(None),
read_ref("refs/heads/README"),
update(
"refs/heads/main",
&[0x12; 32],
Some(RefExpectation::Match),
Some(&[0x11; 32]),
),
],
);
assert_eq!(end, SessionEnd::Clean);
let listed = |f: &SshFrame| match &f.body {
Some(Body::ListRefsResponse(r)) => r
.refs
.iter()
.map(|e| e.name.clone().unwrap())
.collect::<Vec<_>>(),
other => panic!("expected ListRefsResponse, got {other:?}"),
};
assert_eq!(listed(&frames[0]), ["main"]);
assert!(listed(&frames[1]).is_empty());
assert_eq!(listed(&frames[2]), ["refs/heads/main"]);
assert_error(&frames[3], ErrorCode::Internal, "read ref failed");
assert!(matches!(frames[4].body, Some(Body::UpdateRefResponse(_))));
}
#[test]
fn over_long_legacy_ref_file_is_skipped_and_refused() {
let td = tempfile::tempdir().unwrap();
let tx = FileTransport::new(td.path());
let long = format!("refs/heads/{}", vec!["a".repeat(200); 3].join("/"));
let path = long
.split('/')
.fold(td.path().to_path_buf(), |p, s| p.join(s));
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, mkit_core::refs::encode_ref_wire(&[9; 32])).unwrap();
tx.update_ref("refs/heads/main", RefWriteCondition::Any, &[1; 32])
.unwrap();
let pipe = fs_pipe(td.path());
let (_, frames) = run(
&pipe,
[
list_refs(Some("refs/heads/")),
read_ref(&long),
update(&long, &[8; 32], Some(RefExpectation::Match), Some(&[9; 32])),
],
);
let Some(Body::ListRefsResponse(r)) = &frames[0].body else {
panic!("expected ListRefsResponse");
};
let names: Vec<_> = r.refs.iter().map(|e| e.name.clone().unwrap()).collect();
assert_eq!(names, ["main"]);
assert_error(&frames[1], ErrorCode::InvalidRequest, "ref name too long");
assert_error(&frames[2], ErrorCode::InvalidRequest, "ref name too long");
assert_eq!(
std::fs::read(&path).unwrap(),
mkit_core::refs::encode_ref_wire(&[9; 32])
);
}
#[test]
fn serve_loop_rejects_invalid_upload_before_storage() {
let td = tempfile::tempdir().unwrap();
let pipe = fs_pipe(td.path());
let bogus = [0x77; 32];
let (end, frames) = run(
&pipe,
[
upload_header(&bogus, Some(5)),
chunk(&bogus, Some(0), b"wrong", true),
],
);
assert_eq!(end, SessionEnd::Clean);
assert!(matches!(frames[0].body, Some(Body::Error(_))));
let tx = FileTransport::new(td.path());
assert!(!tx.pack_exists(&PackKey::new(bogus)).unwrap());
assert!(
files(td.path()).iter().all(|p| !p.starts_with("packs")),
"no pack or temp file left: {:?}",
files(td.path())
);
}
#[test]
fn serve_loop_rejected_upload_does_not_overwrite_existing_pack() {
let td = tempfile::tempdir().unwrap();
let tx = FileTransport::new(td.path());
let (bytes, id) = valid_pack();
tx.upload_pack(&bytes, &PackKey::new(id)).unwrap();
let pipe = fs_pipe(td.path());
let (end, _) = run(
&pipe,
[
upload_header(&id, Some(5)),
chunk(&id, Some(0), b"wrong", true),
],
);
assert_eq!(end, SessionEnd::Clean);
assert_eq!(tx.download_pack(&PackKey::new(id)).unwrap(), bytes);
}
#[test]
fn serve_loop_cas_conflict_carries_current_id_in_details() {
let td = tempfile::tempdir().unwrap();
let pipe = fs_pipe(td.path());
let (winner, loser) = ([0xA1u8; 32], [0xB2u8; 32]);
let main = "refs/heads/main";
let (_, frames) = run(
&pipe,
[
update(main, &winner, Some(RefExpectation::Missing), None),
update(main, &loser, Some(RefExpectation::Missing), None),
],
);
let err = error(&frames[1]).clone();
assert_eq!(err.details.as_deref(), Some(&winner[..]));
assert!(matches!(
mkit_rpc::map_update_ref_error(err, RefWriteCondition::Missing, "ssh"),
mkit_core::protocol::TransportError::RefConflict
));
let tx = FileTransport::new(td.path());
assert_eq!(tx.read_ref(main).unwrap(), Some(winner));
}
#[test]
fn serve_loop_match_conflicts_report_current_value_or_empty() {
let td = tempfile::tempdir().unwrap();
let tx = FileTransport::new(td.path());
let (current, stale, next) = ([0x11u8; 32], [0x22u8; 32], [0x33u8; 32]);
tx.update_ref("refs/heads/main", RefWriteCondition::Any, ¤t)
.unwrap();
let pipe = fs_pipe(td.path());
let (_, frames) = run(
&pipe,
[
update(
"refs/heads/main",
&next,
Some(RefExpectation::Match),
Some(&stale),
),
update(
"refs/heads/ghost",
&next,
Some(RefExpectation::Match),
Some(&stale),
),
],
);
assert_eq!(error(&frames[0]).details.as_deref(), Some(¤t[..]));
assert_eq!(error(&frames[1]).details.as_deref(), Some(&[][..]));
assert_eq!(tx.read_ref("refs/heads/main").unwrap(), Some(current));
assert_eq!(tx.read_ref("refs/heads/ghost").unwrap(), None);
}
#[test]
fn ssh_session_writes_only_packs_and_refs() {
let td = tempfile::tempdir().unwrap();
let pipe = fs_pipe(td.path());
let (bytes, id) = valid_pack();
seed(&pipe, &[("refs/heads/main", [1; 32])], &[&bytes]);
let (_, frames) = run(
&pipe,
[
upload_header(&id, Some(5)),
chunk(&id, Some(0), b"wrong", true),
download(Some(&id)),
],
);
assert_eq!(downloaded(&frames[1..], &id), bytes);
let hex = mkit_core::hash::to_hex(&id);
assert_eq!(
files(td.path()),
[
PathBuf::from(".mkit/refs/.lock"),
PathBuf::from("packs").join(hex),
PathBuf::from("refs/heads/main"),
],
);
}
}