mod common;
use base64::prelude::{BASE64_STANDARD, Engine as _};
use common::*;
use gwk_domain::blob::BLOB_CHUNK_BYTES;
use gwk_domain::ids::Seq;
use gwk_domain::port::EventStore;
use gwk_domain::protocol::{
CONNECTION_EGRESS_BYTES_PER_WINDOW, CONNECTION_INGRESS_BYTES_PER_WINDOW, CONTRACT_VERSION,
FRAME_BODY_MAX_BYTES, KernelErrorCode, KernelResult, MAX_SUBSCRIPTIONS_PER_CONNECTION,
ProjectionKind, ProjectionRecord, ProtocolVersion, SLOW_CONSUMER_TIMEOUT_SECS,
SUBSCRIPTION_POLL_SECS, ServerControl,
};
use gwk_kernel::store::connect_pool;
use gwk_kernel::wire::frame::{Budget, Incoming, read_frame};
use gwk_kernel::wire::listen::Listener;
use gwk_kernel::wire::serve::serve_stream;
use std::sync::Arc;
use tokio::net::UnixStream;
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_sealed_daemon_answers_the_whole_surface_it_promises() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_sealed_store(&maintenance, "wire_sealed", 8).await;
let dir = runtime_dir("sealed");
let path = dir.join("gwk.sock");
let listener = Listener::bind(&path).await.expect("bind");
let (daemon, blobs) = daemon_for(store, "sealed").await;
let daemon = Arc::new(daemon);
let serving = tokio::spawn({
let daemon = Arc::clone(&daemon);
async move {
let (stream, _) = listener.accept().await.expect("accept");
let _ = serve_stream(&daemon, stream).await;
listener.remove();
}
});
let (mut client, ack) = Client::connect(&path).await;
match ack {
ServerControl::HelloAck {
protocol_major,
sealed,
capabilities,
..
} => {
assert_eq!(protocol_major, ProtocolVersion::V1);
assert!(sealed);
assert!(capabilities.is_empty());
}
other => panic!("{other:?}"),
}
assert_eq!(
client.ask("r-health", r#"{"type":"health"}"#).await,
KernelResult::Health {
ready: true,
sealed: true
}
);
match client.ask("r-status", r#"{"type":"status"}"#).await {
KernelResult::Status {
sealed,
contract_version,
public_revision,
watermark,
..
} => {
assert!(sealed);
assert_eq!(contract_version, CONTRACT_VERSION);
assert_eq!(public_revision, TEST_REVISION);
assert!(watermark.is_some(), "genesis is in the log");
}
other => panic!("{other:?}"),
}
let watermark = match client.ask("r-wm", r#"{"type":"watermark"}"#).await {
KernelResult::Watermark { watermark } => watermark.expect("genesis is in the log"),
other => panic!("{other:?}"),
};
match client.ask("r-sealed", r#"{"type":"verify_sealed"}"#).await {
KernelResult::SealedVerification {
sealed,
genesis_watermark,
event_count,
..
} => {
assert!(sealed);
assert_eq!(event_count.value(), 1);
assert_eq!(genesis_watermark, watermark);
}
other => panic!("{other:?}"),
}
drop(client);
serving.await.expect("join");
let _ = std::fs::remove_dir_all(&dir);
let _ = std::fs::remove_dir_all(&blobs);
drop_database(&maintenance, &name).await;
}
fn record_key(record: &ProjectionRecord) -> String {
let json = serde_json::to_value(record).expect("serialize");
let body = &json[record.kind().as_str()];
body.get("id")
.or_else(|| body.get("orchestrator_id"))
.and_then(serde_json::Value::as_str)
.expect("every projection record carries the key it is paged by")
.to_owned()
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn every_projection_a_client_can_name_comes_back_through_the_wire() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_projections", 8).await;
populate(&store).await;
let mut served = Served::open(store, "projections").await;
for kind in ProjectionKind::ALL {
let tag = kind.as_str();
match served
.client
.ask(
&format!("r-{tag}"),
&format!(r#"{{"type":"list_projection","projection":"{tag}"}}"#),
)
.await
{
KernelResult::ProjectionPage { records, .. } => {
assert!(!records.is_empty(), "{tag} came back empty");
for record in &records {
assert_eq!(record.kind(), *kind, "{tag} answered with another table");
}
}
other => panic!("{tag}: {other:?}"),
}
}
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_page_at_a_time_walks_a_projection_exactly_once() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_paging", 8).await;
let ids = ["t-a", "t-A", "t_a", "t-a-1", "t-a1", "tA", "ta"];
for id in ids {
apply(&store, id, task(id)).await;
}
let mut served = Served::open(store, "paging").await;
let mut seen: Vec<String> = Vec::new();
let mut cursor: Option<String> = None;
for round in 0..ids.len() + 2 {
let request = match &cursor {
Some(c) => format!(
r#"{{"type":"list_projection","projection":"task","cursor":"{c}","limit":1}}"#
),
None => r#"{"type":"list_projection","projection":"task","limit":1}"#.to_owned(),
};
match served.client.ask(&format!("r-{round}"), &request).await {
KernelResult::ProjectionPage {
records,
next_cursor,
} => {
seen.extend(records.iter().map(record_key));
match next_cursor {
Some(next) => cursor = Some(next),
None => break,
}
}
other => panic!("{other:?}"),
}
}
let mut expected: Vec<String> = ids.iter().map(|s| (*s).to_owned()).collect();
expected.sort_unstable();
let mut got = seen.clone();
got.sort_unstable();
assert_eq!(
got, expected,
"the walk did not see every task exactly once"
);
assert_eq!(seen.len(), ids.len(), "a row was delivered twice");
assert_eq!(seen, expected, "pages did not arrive in key order");
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn one_record_by_id_is_the_same_record_the_page_delivered() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_get", 8).await;
apply(&store, "t-1", task("t-1")).await;
let mut served = Served::open(store, "get").await;
let from_page = match served
.client
.ask(
"r-list",
r#"{"type":"list_projection","projection":"task"}"#,
)
.await
{
KernelResult::ProjectionPage { mut records, .. } => records.pop().expect("one task"),
other => panic!("{other:?}"),
};
match served
.client
.ask(
"r-get",
r#"{"type":"get_projection","projection":"task","id":"t-1"}"#,
)
.await
{
KernelResult::Projection { record } => assert_eq!(record, from_page),
other => panic!("{other:?}"),
}
match served
.client
.ask(
"r-missing",
r#"{"type":"get_projection","projection":"task","id":"t-nope"}"#,
)
.await
{
KernelResult::Error { code, .. } => assert_eq!(code, KernelErrorCode::NotFound),
other => panic!("{other:?}"),
}
assert!(matches!(
served.client.ask("r-after", r#"{"type":"health"}"#).await,
KernelResult::Health { .. }
));
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn the_log_reads_back_from_a_cursor_in_the_order_it_was_written() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_events", 8).await;
for i in 0..6 {
apply(&store, &format!("t-{i}"), task(&format!("t-{i}"))).await;
}
let mut served = Served::open(store, "events").await;
let watermark = match served.client.ask("r-wm", r#"{"type":"watermark"}"#).await {
KernelResult::Watermark { watermark } => watermark.expect("events exist"),
other => panic!("{other:?}"),
};
let mut collected: Vec<u64> = Vec::new();
let mut cursor: Option<u64> = None;
for round in 0..10 {
let request = match cursor {
Some(c) => format!(r#"{{"type":"read_events","cursor":"{c}","limit":2}}"#),
None => r#"{"type":"read_events","limit":2}"#.to_owned(),
};
match served.client.ask(&format!("r-e{round}"), &request).await {
KernelResult::Events {
events,
cursor: last,
..
} => {
if events.is_empty() {
assert!(last.is_none());
break;
}
collected.extend(events.iter().map(|e| e.global_sequence.value()));
assert_eq!(
last.expect("a non-empty page reports its last sequence"),
events.last().expect("non-empty").global_sequence
);
cursor = Some(last.expect("checked").value());
}
other => panic!("{other:?}"),
}
}
assert!(!collected.is_empty());
assert_eq!(*collected.last().expect("non-empty"), watermark.value());
assert!(
collected.windows(2).all(|w| w[0] < w[1]),
"the log came back out of order or repeated: {collected:?}"
);
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_command_submitted_over_the_wire_lands_once_however_often_it_is_sent() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_submit", 8).await;
let mut served = Served::open(store, "submit").await;
let envelope = serde_json::to_string(&envelope("k-1", &task("t-1"))).expect("serialize");
let first = match served
.client
.ask(
"r-submit",
&format!(r#"{{"type":"submit_command","envelope":{envelope}}}"#),
)
.await
{
KernelResult::CommandApplied { events, .. } => events,
other => panic!("{other:?}"),
};
assert_eq!(first.len(), 1);
let again = match served
.client
.ask(
"r-submit-2",
&format!(r#"{{"type":"submit_command","envelope":{envelope}}}"#),
)
.await
{
KernelResult::CommandApplied { events, .. } => events,
other => panic!("{other:?}"),
};
assert_eq!(again, first, "a retry appended a second event");
match served
.client
.ask(
"r-check",
r#"{"type":"list_projection","projection":"task"}"#,
)
.await
{
KernelResult::ProjectionPage { records, .. } => assert_eq!(records.len(), 1),
other => panic!("{other:?}"),
}
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_subscription_delivers_what_the_log_gains_after_it_started() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_subscribe", 8).await;
let served = Running::open(store, "subscribe").await;
let mut watcher = served.client().await;
let mut appender = served.client().await;
let watermark = watermark_of(&mut watcher, "r-wm").await;
match watcher.ask("r-sub", &subscribe_from(watermark)).await {
KernelResult::Subscribed { cursor } => assert_eq!(cursor, Some(watermark)),
other => panic!("{other:?}"),
}
let envelope = serde_json::to_string(&envelope("k-1", &task("t-1"))).expect("serialize");
match appender
.ask(
"r-submit",
&format!(r#"{{"type":"submit_command","envelope":{envelope}}}"#),
)
.await
{
KernelResult::CommandApplied { .. } => {}
other => panic!("{other:?}"),
}
let batch = tokio::time::timeout(std::time::Duration::from_secs(3), watcher.recv())
.await
.expect("a notified subscription delivers without waiting for the poll")
.expect("the connection stayed open");
match batch {
ServerControl::EventBatch {
request_id,
events,
cursor,
} => {
assert_eq!(request_id.as_str(), "r-sub");
assert!(!events.is_empty(), "a batch with no events");
assert!(
events.iter().all(|e| e.global_sequence > watermark),
"the stream replayed what the cursor had already covered"
);
assert_eq!(
cursor,
events.last().expect("non-empty").global_sequence,
"the batch cursor is not its last event"
);
}
other => panic!("{other:?}"),
}
drop(watcher);
drop(appender);
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_subscription_that_never_hears_a_notification_still_catches_up() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_poll", 8).await;
let mut served = Served::open(store, "poll").await;
let watermark = watermark_of(&mut served.client, "r-wm").await;
match served.client.ask("r-sub", &subscribe_from(watermark)).await {
KernelResult::Subscribed { .. } => {}
other => panic!("{other:?}"),
}
let envelope = serde_json::to_string(&envelope("k-1", &task("t-1"))).expect("serialize");
match served
.client
.ask(
"r-submit",
&format!(r#"{{"type":"submit_command","envelope":{envelope}}}"#),
)
.await
{
KernelResult::CommandApplied { .. } => {}
other => panic!("{other:?}"),
}
let batch = tokio::time::timeout(
std::time::Duration::from_secs(SUBSCRIPTION_POLL_SECS + 5),
served.client.recv(),
)
.await
.expect("the poll delivered what the lost notification did not")
.expect("the connection stayed open");
match batch {
ServerControl::EventBatch { events, cursor, .. } => {
assert!(events.iter().all(|e| e.global_sequence > watermark));
assert_eq!(cursor, events.last().expect("non-empty").global_sequence);
}
other => panic!("{other:?}"),
}
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_frame_the_codec_refuses_takes_down_its_own_connection_and_no_other() {
use tokio::io::AsyncWriteExt as _;
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_strict", 8).await;
let served = Running::open(store, "strict").await;
let mut bystander = served.client().await;
assert!(matches!(
bystander.ask("r-before", r#"{"type":"health"}"#).await,
KernelResult::Health { .. }
));
let mut liar = UnixStream::connect(&served.path).await.expect("connect");
liar.write_all(&(FRAME_BODY_MAX_BYTES + 1).to_be_bytes())
.await
.expect("write a length");
liar.write_all(&[1u8]).await.expect("write a kind");
let mut budget = Budget::new(
CONNECTION_INGRESS_BYTES_PER_WINDOW,
CONNECTION_EGRESS_BYTES_PER_WINDOW,
);
match read_frame(&mut liar, FRAME_BODY_MAX_BYTES, &mut budget)
.await
.expect("read the refusal")
{
Incoming::Frame(frame) => {
match serde_json::from_slice(&frame.body).expect("decode the refusal") {
ServerControl::HelloRefusal { code, .. } => {
assert_eq!(code, KernelErrorCode::FrameSize)
}
other => panic!("{other:?}"),
}
}
Incoming::Closed => panic!("the daemon hung up without saying why"),
}
match read_frame(&mut liar, FRAME_BODY_MAX_BYTES, &mut budget).await {
Ok(Incoming::Closed) => {}
Err(error) => assert!(
error.fatal,
"a non-fatal error on a dead connection: {error}"
),
Ok(Incoming::Frame(frame)) => {
panic!("the connection outlived the frame that broke it: {frame:?}")
}
}
let (mut sneak, _) = Client::connect(&served.path).await;
sneak
.send(r#"{"type":"request","request_id":"r-x","request":{"type":"health","extra":1}}"#)
.await;
assert!(
sneak.recv().await.is_none(),
"an unknown field was tolerated"
);
assert!(matches!(
bystander.ask("r-after", r#"{"type":"health"}"#).await,
KernelResult::Health { .. }
));
drop(bystander);
drop(sneak);
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_stream_survives_losing_the_listener_it_was_being_notified_through() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_deaf", 8).await;
let admin = connect_pool(&secret(&name), 2).await.expect("connect");
let served = Running::open(store, "deaf").await;
let mut watcher = served.client().await;
let watermark = watermark_of(&mut watcher, "r-wm").await;
match watcher.ask("r-sub", &subscribe_from(watermark)).await {
KernelResult::Subscribed { .. } => {}
other => panic!("{other:?}"),
}
let first = seed_events(&admin, 3).await;
let batch = tokio::time::timeout(
std::time::Duration::from_secs(SUBSCRIPTION_POLL_SECS + 5),
watcher.recv(),
)
.await
.expect("only the poll could have delivered this, and it had to")
.expect("the connection stayed open");
match batch {
ServerControl::EventBatch { events, cursor, .. } => {
assert!(events.iter().all(|e| e.global_sequence > watermark));
assert_eq!(cursor, first, "the batch did not reach the seeded tail");
}
other => panic!("{other:?}"),
}
let killed: Vec<bool> = sqlx::query_scalar(
"SELECT pg_terminate_backend(pid) FROM pg_stat_activity \
WHERE datname = current_database() AND pid <> pg_backend_pid() \
AND query LIKE 'LISTEN %'",
)
.fetch_all(&admin)
.await
.expect("terminate the listener");
assert!(
killed.iter().any(|ok| *ok),
"no LISTEN backend was found to kill: this half proved nothing"
);
let second = seed_events(&admin, 3).await;
let batch = tokio::time::timeout(
std::time::Duration::from_secs(SUBSCRIPTION_POLL_SECS + 5),
watcher.recv(),
)
.await
.expect("a stream whose listener died still delivers")
.expect("the connection stayed open");
match batch {
ServerControl::EventBatch { cursor, .. } => {
assert_eq!(cursor, second, "the stream stalled at the break")
}
other => panic!("{other:?}"),
}
drop(watcher);
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL, and waits out SLOW_CONSUMER_TIMEOUT_SECS on purpose"]
async fn a_consumer_that_stops_reading_is_cut_off_at_the_last_batch_it_actually_got() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_slow", 8).await;
let before = store
.watermark()
.await
.expect("watermark")
.expect("genesis");
let batches_needed = MAX_SUBSCRIPTIONS_PER_CONNECTION as u64 + 6;
seed_events(store.pool(), batches_needed * 256).await;
let (daemon, blobs) = daemon_for(store, "slow").await;
let dir = runtime_dir("slow");
let (mine, theirs) = tokio::io::duplex(4096);
let serving = tokio::spawn(async move {
let (mut reader, mut writer) = tokio::io::split(theirs);
let _ = gwk_kernel::wire::serve::serve_connection(&daemon, &mut reader, &mut writer).await;
});
let (mut client, _) = Client::greet(mine).await;
match client.ask("r-sub", &subscribe_from(before)).await {
KernelResult::Subscribed { .. } => {}
other => panic!("{other:?}"),
}
let mut received: Vec<Seq> = Vec::new();
for _ in 0..2 {
match client.recv().await.expect("the connection stayed open") {
ServerControl::EventBatch { cursor, .. } => received.push(cursor),
other => panic!("{other:?}"),
}
}
tokio::time::sleep(std::time::Duration::from_secs(
SLOW_CONSUMER_TIMEOUT_SECS + 3,
))
.await;
let closed = loop {
match client.recv().await.expect("the connection stayed open") {
ServerControl::EventBatch { cursor, .. } => received.push(cursor),
ServerControl::StreamClosed {
request_id,
code,
last_cursor,
} => {
assert_eq!(request_id.as_str(), "r-sub");
assert_eq!(code, KernelErrorCode::SlowConsumer);
break last_cursor.expect("the consumer had received batches");
}
other => panic!("{other:?}"),
}
};
let position = received
.iter()
.position(|cursor| *cursor == closed)
.unwrap_or_else(|| {
panic!("the stream closed at {closed:?}, which was never sent: {received:?}")
});
assert!(
received.len() - 1 - position <= 1,
"closed at {closed:?}, {} batches behind the {} that were sent",
received.len() - 1 - position,
received.len()
);
assert!(
closed.value() < before.value() + batches_needed * 256,
"the kernel reported its own read position, not what it delivered"
);
drop(client);
serving.await.expect("join");
let _ = std::fs::remove_dir_all(&dir);
let _ = std::fs::remove_dir_all(&blobs);
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_connection_may_not_hold_more_subscriptions_than_its_cap() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_subcap", 8).await;
let mut served = Served::open(store, "subcap").await;
let watermark = watermark_of(&mut served.client, "r-wm").await;
for i in 0..MAX_SUBSCRIPTIONS_PER_CONNECTION {
match served
.client
.ask(&format!("r-s{i}"), &subscribe_from(watermark))
.await
{
KernelResult::Subscribed { .. } => {}
other => panic!("subscription {i}: {other:?}"),
}
}
match served
.client
.ask("r-over", &subscribe_from(watermark))
.await
{
KernelResult::Error { code, message, .. } => {
assert_eq!(code, KernelErrorCode::Overloaded);
assert!(
message.contains(&MAX_SUBSCRIPTIONS_PER_CONNECTION.to_string()),
"the refusal does not say what the cap is: {message}"
);
}
other => panic!("{other:?}"),
}
assert!(matches!(
served.client.ask("r-after", r#"{"type":"health"}"#).await,
KernelResult::Health { .. }
));
served.close().await;
drop_database(&maintenance, &name).await;
}
async fn begin_upload(client: &mut Client, id: &str, plaintext: &[u8]) -> String {
match client
.ask(
id,
&format!(
r#"{{"type":"blob_begin","media_type":"application/octet-stream","byte_size":"{}"}}"#,
plaintext.len()
),
)
.await
{
KernelResult::BlobBegun { upload_id } => upload_id.as_str().to_owned(),
other => panic!("{other:?}"),
}
}
fn chunk_request(upload_id: &str, sequence: u32, chunk: &[u8]) -> String {
format!(
r#"{{"type":"blob_chunk","upload_id":"{upload_id}","sequence":{sequence},"data_base64":"{}"}}"#,
BASE64_STANDARD.encode(chunk)
)
}
async fn read_blob(client: &mut Client, address: &str, size: usize) -> Vec<u8> {
let mut out: Vec<u8> = Vec::with_capacity(size);
while out.len() < size {
let request = format!(
r#"{{"type":"blob_read","address":"{address}","offset":"{}","length":"{}"}}"#,
out.len(),
size - out.len()
);
match client.ask(&format!("r-read{}", out.len()), &request).await {
KernelResult::BlobBytes {
offset,
data_base64,
..
} => {
assert_eq!(offset.value(), out.len() as u64);
let part = BASE64_STANDARD.decode(&data_base64).expect("base64");
assert!(!part.is_empty(), "the read stalled at {}", out.len());
out.extend_from_slice(&part);
}
other => panic!("{other:?}"),
}
}
out
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_blob_goes_up_in_chunks_and_comes_back_byte_for_byte() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_blob", 8).await;
let mut served = Served::open(store, "blob").await;
let plaintext: Vec<u8> = (0..9_001u32).map(|i| (i % 251) as u8).collect();
let address = address_of(&plaintext);
let upload = begin_upload(&mut served.client, "r-begin", &plaintext).await;
for (sequence, chunk) in [&plaintext[..1], &plaintext[1..5_000], &plaintext[5_000..]]
.into_iter()
.enumerate()
{
let sequence = sequence as u32;
match served
.client
.ask(
&format!("r-chunk{sequence}"),
&chunk_request(&upload, sequence, chunk),
)
.await
{
KernelResult::BlobChunkAccepted {
upload_id,
sequence: acked,
} => {
assert_eq!(upload_id.as_str(), upload);
assert_eq!(acked, sequence);
}
other => panic!("chunk {sequence}: {other:?}"),
}
}
let descriptor = match served
.client
.ask(
"r-commit",
&format!(
r#"{{"type":"blob_commit","upload_id":"{upload}","address":"{}"}}"#,
address.as_str()
),
)
.await
{
KernelResult::BlobCommitted {
descriptor,
deduplicated,
} => {
assert!(!deduplicated, "nothing was there to deduplicate against");
assert_eq!(descriptor.address, address);
assert_eq!(descriptor.byte_size.value(), plaintext.len() as u64);
assert!(!descriptor.tombstoned);
descriptor
}
other => panic!("{other:?}"),
};
match served
.client
.ask(
"r-stat",
&format!(r#"{{"type":"blob_stat","address":"{}"}}"#, address.as_str()),
)
.await
{
KernelResult::BlobStat { descriptor: stat } => assert_eq!(stat, descriptor),
other => panic!("{other:?}"),
}
let back = read_blob(&mut served.client, address.as_str(), plaintext.len()).await;
assert_eq!(back, plaintext, "the blob did not come back as it went in");
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn the_same_bytes_uploaded_twice_are_stored_once() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_blob_dedup", 8).await;
let mut served = Served::open(store, "blobdedup").await;
let plaintext = b"the same bytes, twice".to_vec();
let address = address_of(&plaintext);
let commit = format!(
r#"{{"type":"blob_commit","upload_id":"{{upload}}","address":"{}"}}"#,
address.as_str()
);
for round in 0..2 {
let upload = begin_upload(&mut served.client, &format!("r-begin{round}"), &plaintext).await;
match served
.client
.ask(
&format!("r-chunk{round}"),
&chunk_request(&upload, 0, &plaintext),
)
.await
{
KernelResult::BlobChunkAccepted { .. } => {}
other => panic!("{other:?}"),
}
match served
.client
.ask(
&format!("r-commit{round}"),
&commit.replace("{upload}", &upload),
)
.await
{
KernelResult::BlobCommitted { deduplicated, .. } => assert_eq!(
deduplicated,
round == 1,
"round {round} reported the wrong dedup verdict"
),
other => panic!("{other:?}"),
}
}
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn an_upload_that_does_not_add_up_is_refused_and_can_be_abandoned() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_blob_bad", 8).await;
let mut served = Served::open(store, "blobbad").await;
let plaintext = b"eleven bytes".to_vec();
let upload = begin_upload(&mut served.client, "r-begin", &plaintext).await;
match served
.client
.ask(
"r-garbage",
&format!(
r#"{{"type":"blob_chunk","upload_id":"{upload}","sequence":0,"data_base64":"not base64 at all!"}}"#
),
)
.await
{
KernelResult::Error { code, .. } => assert_eq!(code, KernelErrorCode::Validation),
other => panic!("{other:?}"),
}
match served
.client
.ask("r-gap", &chunk_request(&upload, 3, &plaintext))
.await
{
KernelResult::Error { code, .. } => assert_eq!(code, KernelErrorCode::BlobIntegrity),
other => panic!("{other:?}"),
}
match served
.client
.ask(
"r-oversized",
&chunk_request(&upload, 0, &vec![0u8; BLOB_CHUNK_BYTES + 1]),
)
.await
{
KernelResult::Error { code, message, .. } => {
assert_eq!(code, KernelErrorCode::Validation);
assert!(message.contains("chunk"), "{message}");
}
other => panic!("{other:?}"),
}
match served
.client
.ask("r-chunk", &chunk_request(&upload, 0, &plaintext))
.await
{
KernelResult::BlobChunkAccepted { .. } => {}
other => panic!("{other:?}"),
}
let wrong = address_of(b"different bytes entirely");
match served
.client
.ask(
"r-liar",
&format!(
r#"{{"type":"blob_commit","upload_id":"{upload}","address":"{}"}}"#,
wrong.as_str()
),
)
.await
{
KernelResult::Error { code, .. } => assert_eq!(code, KernelErrorCode::BlobIntegrity),
other => panic!("{other:?}"),
}
match served
.client
.ask(
"r-abort",
&format!(r#"{{"type":"blob_abort","upload_id":"{upload}"}}"#),
)
.await
{
KernelResult::BlobAborted { upload_id } => assert_eq!(upload_id.as_str(), upload),
other => panic!("{other:?}"),
}
match served
.client
.ask("r-after-abort", &chunk_request(&upload, 1, &plaintext))
.await
{
KernelResult::Error { code, .. } => assert_eq!(code, KernelErrorCode::NotFound),
other => panic!("{other:?}"),
}
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_read_larger_than_a_frame_is_clamped_rather_than_refused() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "wire_blob_clamp", 8).await;
let mut served = Served::open(store, "blobclamp").await;
let plaintext = b"a few bytes".to_vec();
let address = address_of(&plaintext);
let upload = begin_upload(&mut served.client, "r-begin", &plaintext).await;
match served
.client
.ask("r-chunk", &chunk_request(&upload, 0, &plaintext))
.await
{
KernelResult::BlobChunkAccepted { .. } => {}
other => panic!("{other:?}"),
}
match served
.client
.ask(
"r-commit",
&format!(
r#"{{"type":"blob_commit","upload_id":"{upload}","address":"{}"}}"#,
address.as_str()
),
)
.await
{
KernelResult::BlobCommitted { .. } => {}
other => panic!("{other:?}"),
}
match served
.client
.ask(
"r-huge",
&format!(
r#"{{"type":"blob_read","address":"{}","offset":"0","length":"{}"}}"#,
address.as_str(),
u64::from(u32::MAX)
),
)
.await
{
KernelResult::BlobBytes { data_base64, .. } => {
let bytes = BASE64_STANDARD.decode(&data_base64).expect("base64");
assert_eq!(bytes, plaintext);
}
other => panic!("{other:?}"),
}
served.close().await;
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn a_blob_that_was_never_written_is_absent_rather_than_an_error() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_sealed_store(&maintenance, "wire_unwritten", 8).await;
let dir = runtime_dir("unwritten");
let path = dir.join("gwk.sock");
let listener = Listener::bind(&path).await.expect("bind");
let (daemon, blob_root) = daemon_for(store, "unwritten").await;
let daemon = Arc::new(daemon);
let serving = tokio::spawn({
let daemon = Arc::clone(&daemon);
async move {
let (stream, _) = listener.accept().await.expect("accept");
let _ = serve_stream(&daemon, stream).await;
listener.remove();
}
});
let (mut client, _) = Client::connect(&path).await;
let address = format!("sha256:{}", "0".repeat(64));
match client
.ask(
"r-stat",
&format!(r#"{{"type":"blob_stat","address":"{address}"}}"#),
)
.await
{
KernelResult::Error { code, .. } => assert_eq!(code, KernelErrorCode::NotFound),
other => panic!("{other:?}"),
}
assert!(matches!(
client.ask("r-after", r#"{"type":"health"}"#).await,
KernelResult::Health { .. }
));
drop(client);
serving.await.expect("join");
let _ = std::fs::remove_dir_all(&dir);
let _ = std::fs::remove_dir_all(&blob_root);
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn the_accept_loop_stops_taking_work_and_takes_its_socket_with_it() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_sealed_store(&maintenance, "wire_drain", 8).await;
let dir = runtime_dir("drain");
let path = dir.join("gwk.sock");
let listener = Listener::bind(&path).await.expect("bind");
let (daemon, blob_root) = daemon_for(store, "drain").await;
let daemon = Arc::new(daemon);
let (stop, stopped) = tokio::sync::oneshot::channel::<()>();
let running = tokio::spawn(gwk_kernel::wire::serve::run(listener, daemon, async move {
let _ = stopped.await;
}));
let (mut client, _) = Client::connect(&path).await;
assert!(matches!(
client.ask("r-1", r#"{"type":"health"}"#).await,
KernelResult::Health { .. }
));
stop.send(()).expect("signal shutdown");
drop(client);
running.await.expect("join").expect("run");
assert!(!path.exists(), "shutdown left the socket behind");
assert!(UnixStream::connect(&path).await.is_err());
let _ = std::fs::remove_dir_all(&dir);
let _ = std::fs::remove_dir_all(&blob_root);
drop_database(&maintenance, &name).await;
}