mod common;
use common::{
TEST_KEK, address_of, blob_store, drop_database, maintenance_pool, put, put_dedup, raw_store,
read_all,
};
use gwk_domain::blob::{BLOB_CHUNK_BYTES, BlobAddress};
use gwk_domain::envelope::{Actor, EventEnvelope, Origin, PayloadRef};
use gwk_domain::ids::{
AggregateId, BlobUploadId, ByteCount, EventId, EvidenceId, ProjectId, Seq, Timestamp,
};
use gwk_domain::port::{BlobError, BlobStore, EventStore};
fn bytes(len: usize) -> Vec<u8> {
(0..len).map(|i| ((i * 31 + i / 97) % 251) as u8).collect()
}
fn referencing_event(address: &BlobAddress, size: u64) -> EventEnvelope {
EventEnvelope {
event_id: EventId::new("evt-1"),
project_id: ProjectId::new("p"),
aggregate_type: "task".into(),
aggregate_id: AggregateId::new("t1"),
aggregate_version: 1,
event_type: "artifact_recorded".into(),
schema_version: 1,
global_sequence: Seq::new(0),
occurred_at: Timestamp::new("2026-07-28T00:00:00Z"),
appended_at: Timestamp::new("2026-07-28T00:00:00Z"),
actor: Actor {
kind: "kernel".into(),
id: None,
},
origin: Origin {
system: "gw".into(),
r#ref: None,
},
causation_id: None,
correlation_id: None,
idempotency_key: None,
payload: serde_json::json!({}),
payload_ref: Some(PayloadRef {
digest: address.as_str().to_owned(),
media_type: "application/json".to_owned(),
byte_size: ByteCount::new(size),
retention_class: None,
evidence_pin: None,
}),
}
}
async fn teardown(maintenance: &sqlx::PgPool, name: &str, root: &std::path::Path) {
drop_database(maintenance, name).await;
let _ = std::fs::remove_dir_all(root);
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn a_blob_survives_chunked_upload_and_every_ranged_read_of_it() {
let maintenance = maintenance_pool().await;
let (name, store) = raw_store(&maintenance, "blobround", 8).await;
let (root, blobs) = blob_store(&store, "blobround").await;
let plaintext = bytes(BLOB_CHUNK_BYTES * 2 + 7);
let size = plaintext.len() as u64;
let address = put(&blobs, &plaintext, "application/octet-stream").await;
assert_eq!(address, address_of(&plaintext));
let descriptor = blobs
.stat(&address)
.await
.expect("stat")
.expect("committed");
assert_eq!(descriptor.byte_size.value(), size);
assert_eq!(descriptor.media_type, "application/octet-stream");
assert_eq!(descriptor.kek_id, common::TEST_KEK_ID);
assert!(!descriptor.pinned && !descriptor.tombstoned);
assert_eq!(read_all(&blobs, &address, size).await, plaintext);
let chunk = BLOB_CHUNK_BYTES as u64;
for (why, offset, length) in [
("the very start", 0u64, 10u64),
("mid first chunk", 500, 100),
("across the first boundary", chunk - 5, 10),
("exactly a chunk start", chunk, 16),
("across the second boundary", chunk * 2 - 3, 6),
("the short final chunk", chunk * 2, 7),
("the last byte", size - 1, 1),
] {
let got = blobs
.read(&address, ByteCount::new(offset), ByteCount::new(length))
.await
.unwrap_or_else(|e| panic!("{why}: {e}"));
let want = &plaintext[offset as usize..(offset + length) as usize];
assert_eq!(got, want, "{why}");
}
let big = blobs
.read(&address, ByteCount::new(0), ByteCount::new(u64::MAX))
.await
.expect("clamped read");
assert_eq!(big.len(), BLOB_CHUNK_BYTES);
assert_eq!(big, plaintext[..BLOB_CHUNK_BYTES]);
for offset in [size, size + 1, u64::MAX] {
let past = blobs
.read(&address, ByteCount::new(offset), ByteCount::new(64))
.await
.expect("past the end");
assert!(past.is_empty(), "offset {offset}");
}
let tail = blobs
.read(&address, ByteCount::new(size - 3), ByteCount::new(999))
.await
.expect("tail");
assert_eq!(tail, plaintext[(size - 3) as usize..]);
let empty = put(&blobs, b"", "text/plain").await;
assert_eq!(
blobs
.stat(&empty)
.await
.expect("stat")
.expect("committed")
.byte_size
.value(),
0
);
assert!(
blobs
.read(&empty, ByteCount::new(0), ByteCount::new(16))
.await
.expect("read empty")
.is_empty()
);
teardown(&maintenance, &name, &root).await;
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn dedup_needs_the_digest_and_the_media_type_to_agree() {
let maintenance = maintenance_pool().await;
let (name, store) = raw_store(&maintenance, "blobdedup", 8).await;
let (root, blobs) = blob_store(&store, "blobdedup").await;
let plaintext = b"the same bytes twice".to_vec();
let (address, first) = put_dedup(&blobs, &plaintext, "text/plain").await;
assert!(!first, "the first commit stores");
let (again, second) = put_dedup(&blobs, &plaintext, "text/plain").await;
assert!(second, "the second commit dedups");
assert_eq!(again, address);
let upload = blobs
.begin(
"application/json".to_owned(),
ByteCount::new(plaintext.len() as u64),
)
.await
.expect("begin");
blobs
.write_chunk(&upload, 0, &plaintext)
.await
.expect("chunk");
let err = blobs
.commit(upload, address.clone())
.await
.expect_err("a conflicting media type must be refused");
assert!(
matches!(&err, BlobError::Integrity(reason) if reason.contains("text/plain")),
"{err:?}"
);
assert_eq!(
blobs
.stat(&address)
.await
.expect("stat")
.expect("still there")
.media_type,
"text/plain"
);
teardown(&maintenance, &name, &root).await;
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn an_upload_is_bounded_by_the_order_and_the_size_it_declared() {
let maintenance = maintenance_pool().await;
let (name, store) = raw_store(&maintenance, "blobupload", 8).await;
let (root, blobs) = blob_store(&store, "blobupload").await;
let upload = blobs
.begin("text/plain".to_owned(), ByteCount::new(10))
.await
.expect("begin");
for bad in [1u32, 2, 7] {
let err = blobs
.write_chunk(&upload, bad, b"xx")
.await
.expect_err("out of order");
assert!(matches!(err, BlobError::Integrity(_)), "{bad}: {err:?}");
}
blobs
.write_chunk(&upload, 0, b"12345")
.await
.expect("chunk 0");
let err = blobs
.write_chunk(&upload, 0, b"12345")
.await
.expect_err("a repeat is not a retry once it landed");
assert!(matches!(err, BlobError::Integrity(_)), "{err:?}");
let err = blobs
.write_chunk(&upload, 1, &[0u8; 6])
.await
.expect_err("past the declared size");
assert!(
matches!(&err, BlobError::Integrity(reason) if reason.contains("declared")),
"{err:?}"
);
let err = blobs
.commit(upload.clone(), address_of(b"12345"))
.await
.expect_err("short of the declared size");
assert!(
matches!(&err, BlobError::Integrity(reason) if reason.contains("5 of the 10")),
"{err:?}"
);
blobs
.write_chunk(&upload, 1, b"67890")
.await
.expect("chunk 1");
let wrong = address_of(b"something else entirely");
let err = blobs
.commit(upload.clone(), wrong.clone())
.await
.expect_err("a claimed address that is not the digest");
match err {
BlobError::DigestMismatch { expected, actual } => {
assert_eq!(expected, wrong);
assert_eq!(actual, address_of(b"1234567890"));
}
other => panic!("expected DigestMismatch, got {other:?}"),
}
let (descriptor, deduped) = blobs
.commit(upload.clone(), address_of(b"1234567890"))
.await
.expect("commit");
assert!(!deduped);
assert_eq!(descriptor.byte_size.value(), 10);
let err = blobs.abort(upload).await.expect_err("already committed");
assert!(matches!(err, BlobError::NotFound), "{err:?}");
for forged in ["../../etc/passwd", "", &"a".repeat(31), "NOTHEX"] {
let err = blobs
.write_chunk(&BlobUploadId::new(forged), 0, b"x")
.await
.expect_err("a forged upload id");
assert!(matches!(err, BlobError::NotFound), "{forged:?}: {err:?}");
}
let doomed = blobs
.begin("text/plain".to_owned(), ByteCount::new(4))
.await
.expect("begin");
blobs.write_chunk(&doomed, 0, b"abcd").await.expect("chunk");
let staged = root.join("uploads").join(doomed.as_str());
assert!(staged.exists());
blobs.abort(doomed).await.expect("abort");
assert!(!staged.exists(), "abort must take the staging file with it");
teardown(&maintenance, &name, &root).await;
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn sweep_reclaims_what_the_log_stopped_pointing_at() {
let maintenance = maintenance_pool().await;
let (name, store) = raw_store(&maintenance, "blobsweep", 8).await;
let (root, blobs) = blob_store(&store, "blobsweep").await;
let referenced = bytes(64);
let orphan = bytes(65);
let pinned = bytes(66);
let kept = put(&blobs, &referenced, "application/json").await;
let loose = put(&blobs, &orphan, "application/json").await;
let held = put(&blobs, &pinned, "application/json").await;
store
.append(
0,
None,
vec![referencing_event(&kept, referenced.len() as u64)],
)
.await
.expect("append");
blobs
.pin(&held, &EvidenceId::new("ev-1"))
.await
.expect("pin");
let swept = blobs.sweep().await.expect("sweep");
assert_eq!(swept, vec![loose.clone()]);
assert!(blobs.stat(&loose).await.expect("stat").is_none());
assert!(
!root
.join("blobs")
.join(&loose.digest_hex()[0..2])
.join(&loose.digest_hex()[2..4])
.join(loose.digest_hex())
.exists()
);
assert_eq!(
read_all(&blobs, &kept, referenced.len() as u64).await,
referenced
);
assert_eq!(read_all(&blobs, &held, pinned.len() as u64).await, pinned);
assert!(blobs.sweep().await.expect("second sweep").is_empty());
blobs
.unpin(&held, &EvidenceId::new("ev-1"))
.await
.expect("unpin");
assert_eq!(blobs.sweep().await.expect("sweep"), vec![held.clone()]);
assert!(blobs.stat(&kept).await.expect("stat").is_some());
teardown(&maintenance, &name, &root).await;
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn a_pin_survives_every_release_but_the_last() {
let maintenance = maintenance_pool().await;
let (name, store) = raw_store(&maintenance, "blobpin", 8).await;
let (root, blobs) = blob_store(&store, "blobpin").await;
let plaintext = bytes(32);
let address = put(&blobs, &plaintext, "application/json").await;
for evidence in ["ev-a", "ev-b"] {
blobs
.pin(&address, &EvidenceId::new(evidence))
.await
.expect("pin");
}
blobs
.pin(&address, &EvidenceId::new("ev-a"))
.await
.expect("pin again");
assert!(
blobs
.stat(&address)
.await
.expect("stat")
.expect("there")
.pinned
);
assert!(blobs.sweep().await.expect("sweep").is_empty());
assert!(matches!(
blobs.shred(&address).await.expect_err("pinned"),
BlobError::Pinned
));
blobs
.unpin(&address, &EvidenceId::new("ev-a"))
.await
.expect("unpin a");
assert!(blobs.sweep().await.expect("sweep").is_empty());
blobs
.unpin(&address, &EvidenceId::new("ev-never"))
.await
.expect("unpin an absent hold");
blobs
.unpin(&address, &EvidenceId::new("ev-b"))
.await
.expect("unpin b");
assert!(
!blobs
.stat(&address)
.await
.expect("stat")
.expect("there")
.pinned
);
blobs.shred(&address).await.expect("shred");
teardown(&maintenance, &name, &root).await;
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn crypto_shred_is_permanent_even_against_someone_holding_the_plaintext() {
let maintenance = maintenance_pool().await;
let (name, store) = raw_store(&maintenance, "blobshred", 8).await;
let (root, blobs) = blob_store(&store, "blobshred").await;
let plaintext = bytes(4096);
let address = put(&blobs, &plaintext, "application/json").await;
let path = root
.join("blobs")
.join(&address.digest_hex()[0..2])
.join(&address.digest_hex()[2..4])
.join(address.digest_hex());
assert!(path.exists());
blobs.shred(&address).await.expect("shred");
assert!(!path.exists(), "shred removes the ciphertext after the key");
for err in [
blobs
.read(&address, ByteCount::new(0), ByteCount::new(16))
.await
.expect_err("read a shredded blob"),
blobs
.stat(&address)
.await
.expect_err("stat a shredded blob"),
] {
assert!(matches!(err, BlobError::Tombstoned), "{err:?}");
}
blobs.shred(&address).await.expect("shred again");
let upload = blobs
.begin(
"application/json".to_owned(),
ByteCount::new(plaintext.len() as u64),
)
.await
.expect("begin");
blobs
.write_chunk(&upload, 0, &plaintext)
.await
.expect("chunk");
let err = blobs
.commit(upload, address.clone())
.await
.expect_err("a shredded address stays shredded");
assert!(matches!(err, BlobError::Tombstoned), "{err:?}");
assert!(blobs.sweep().await.expect("sweep").is_empty());
assert!(matches!(
blobs.stat(&address).await.expect_err("still tombstoned"),
BlobError::Tombstoned
));
assert!(
blobs
.stat(&address_of(b"never uploaded"))
.await
.expect("stat")
.is_none()
);
assert!(matches!(
blobs
.shred(&address_of(b"never uploaded"))
.await
.expect_err("no such blob"),
BlobError::NotFound
));
teardown(&maintenance, &name, &root).await;
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn a_tampered_container_never_returns_bytes() {
let maintenance = maintenance_pool().await;
let (name, store) = raw_store(&maintenance, "blobtamper", 8).await;
let (root, blobs) = blob_store(&store, "blobtamper").await;
let plaintext = bytes(BLOB_CHUNK_BYTES + 128);
let address = put(&blobs, &plaintext, "application/json").await;
let path = root
.join("blobs")
.join(&address.digest_hex()[0..2])
.join(&address.digest_hex()[2..4])
.join(address.digest_hex());
let original = std::fs::read(&path).expect("read container");
let last_chunk = BLOB_CHUNK_BYTES as u64;
for (why, at, offset) in [
("the header", 16usize, 0u64),
("the first chunk", original.len() / 4, 0),
("the final chunk's tag", original.len() - 1, last_chunk),
] {
let mut edited = original.clone();
edited[at] ^= 0xff;
std::fs::write(&path, &edited).expect("write");
let err = blobs
.read(&address, ByteCount::new(offset), ByteCount::new(64))
.await
.err()
.unwrap_or_else(|| panic!("{why}: a tampered container must not read"));
assert!(matches!(err, BlobError::Integrity(_)), "{why}: {err:?}");
}
let mut edited = original.clone();
let end = edited.len() - 1;
edited[end] ^= 0xff;
std::fs::write(&path, &edited).expect("write");
assert_eq!(
blobs
.read(&address, ByteCount::new(0), ByteCount::new(64))
.await
.expect("an untouched chunk still reads"),
plaintext[..64]
);
std::fs::write(&path, &original[..original.len() - 4096]).expect("truncate");
let err = blobs
.read(&address, ByteCount::new(last_chunk), ByteCount::new(64))
.await
.expect_err("a truncated container must not read");
assert!(matches!(err, BlobError::Integrity(_)), "{err:?}");
std::fs::write(&path, &original).expect("restore");
assert_eq!(
read_all(&blobs, &address, plaintext.len() as u64).await,
plaintext
);
teardown(&maintenance, &name, &root).await;
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn rotation_replaces_every_wrapped_key_and_no_ciphertext() {
let maintenance = maintenance_pool().await;
let (name, store) = raw_store(&maintenance, "blobrotate", 8).await;
let (root, blobs) = blob_store(&store, "blobrotate").await;
let blobs_in = [bytes(100), bytes(BLOB_CHUNK_BYTES + 1), Vec::new()];
let mut addresses = Vec::new();
for (n, plaintext) in blobs_in.iter().enumerate() {
addresses.push(put(&blobs, plaintext, &format!("application/x-{n}")).await);
}
let containers: Vec<Vec<u8>> = addresses
.iter()
.map(|a| {
std::fs::read(
root.join("blobs")
.join(&a.digest_hex()[0..2])
.join(&a.digest_hex()[2..4])
.join(a.digest_hex()),
)
.expect("read container")
})
.collect();
let doomed = put(&blobs, &bytes(7), "application/json").await;
blobs.shred(&doomed).await.expect("shred");
let mut new_kek = TEST_KEK;
new_kek[0] ^= 0xff;
assert_eq!(blobs.rewrap_all(&new_kek).await.expect("rewrap"), 3);
for (address, before) in addresses.iter().zip(&containers) {
let after = std::fs::read(
root.join("blobs")
.join(&address.digest_hex()[0..2])
.join(&address.digest_hex()[2..4])
.join(address.digest_hex()),
)
.expect("read container");
assert_eq!(&after, before, "{address} ciphertext changed");
}
for (address, plaintext) in addresses.iter().zip(&blobs_in) {
if plaintext.is_empty() {
continue;
}
assert!(
blobs
.read(address, ByteCount::new(0), ByteCount::new(8))
.await
.is_err(),
"{address} still opens under the old key"
);
}
let rotated = common::blob_store_with(&store, &root, new_kek).await;
for (address, plaintext) in addresses.iter().zip(&blobs_in) {
assert_eq!(
read_all(&rotated, address, plaintext.len() as u64).await,
*plaintext,
"{address}"
);
}
assert!(matches!(
rotated.stat(&doomed).await.expect_err("still shredded"),
BlobError::Tombstoned
));
teardown(&maintenance, &name, &root).await;
}