mod common;
use std::sync::Arc;
use common::{
TEST_KEK, TEST_REVISION, apply, blob_store_with, checkpointing_store, drop_database,
fresh_store, maintenance_pool, populate, read_all, runtime_dir, secret, task,
};
use gwk_domain::blob::BlobAddress;
use gwk_domain::port::{BlobStore, EventStore};
use gwk_kernel::checkpoint::{self, RECORDS_MEDIA_TYPE};
use gwk_kernel::recover::Verdict;
use gwk_kernel::store::{PgEventStore, connect_pool};
use gwk_kernel::wire::listen::Listener;
use gwk_kernel::wire::serve::{self, Daemon};
fn tags(records: &[u8]) -> Vec<String> {
let text = String::from_utf8(records.to_vec()).expect("canonical records are utf-8");
text.lines()
.map(|line| {
let value: serde_json::Value = serde_json::from_str(line).expect("a record per line");
value["projection_type"]
.as_str()
.expect("every record is tagged")
.to_owned()
})
.collect()
}
async fn canonical(store: &PgEventStore) -> Vec<u8> {
let mut conn = store.pool().acquire().await.expect("connection");
checkpoint::canonical_records(&mut conn)
.await
.expect("canonical records")
}
async fn derived(store: &PgEventStore) -> Vec<u8> {
let mut conn = store.pool().acquire().await.expect("connection");
checkpoint::derived_records(&mut conn)
.await
.expect("derived records")
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn every_projection_round_trips_through_its_contract_type() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "cpparity", 8).await;
populate(&store).await;
let records = canonical(&store).await;
let present = tags(&records);
for expected in [
"attempt",
"attention_item",
"authority_grant",
"command",
"dispatch_node",
"engine_session",
"evidence",
"gate",
"lease",
"message",
"orchestrator_checkpoint",
"receipt",
"task",
"worktree",
] {
assert!(
present.iter().any(|tag| tag == expected),
"no {expected} row reached the canonical form: {present:?}"
);
}
let mut sorted = present.clone();
sorted.sort();
assert_eq!(present, sorted, "the visit order is not stable");
let text = String::from_utf8(records.clone()).expect("utf-8");
assert!(
text.contains("\"byte_size\":\"18446744073709551615\""),
"evidence byte_size lost its decimal-string form"
);
assert!(
text.contains("\"seq\":\"18446744073709551615\""),
"orchestrator checkpoint seq lost its decimal-string form"
);
assert_eq!(canonical(&store).await, records);
drop_database(&maintenance, &name).await;
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn the_barrier_fires_on_either_bound_and_stores_what_it_hashed() {
let maintenance = maintenance_pool().await;
let (name, root, store) = checkpointing_store(&maintenance, "cpbarrier").await;
populate(&store).await;
assert!(checkpoints(&store).await.is_empty());
sqlx::query("UPDATE gwk_internal.writer SET checkpoint_at = now() - interval '1 hour'")
.execute(store.pool())
.await
.expect("backdate the barrier");
apply(&store, "after-time", task("t-2")).await;
let taken = checkpoints(&store).await;
assert_eq!(taken.len(), 1, "the time bound must trip exactly once");
let first = &taken[0];
assert_eq!(first.schema_version, 1);
assert_eq!(first.records_ref.media_type, RECORDS_MEDIA_TYPE);
let records = derived(&store).await;
assert_eq!(first.projection_hash, checkpoint::projection_hash(&records));
let blobs = blob_store_with(&store, &root, common::TEST_KEK).await;
let address = BlobAddress::parse(&first.records_ref.digest).expect("a legal address");
assert_eq!(address.digest_hex(), first.projection_hash);
let stored = read_all(&blobs, &address, first.records_ref.byte_size.value()).await;
assert_eq!(stored, records);
apply(&store, "quiet", task("t-3")).await;
assert_eq!(checkpoints(&store).await.len(), 1);
sqlx::query("UPDATE gwk_internal.writer SET next_seq = 20000")
.execute(store.pool())
.await
.expect("jump the sequence");
apply(&store, "after-count", task("t-4")).await;
let taken = checkpoints(&store).await;
assert_eq!(taken.len(), 2, "the event bound must trip");
assert!(taken[0].through_sequence.value() > taken[1].through_sequence.value());
let swept = blobs.sweep().await.expect("sweep");
for checkpoint in &taken {
let address = BlobAddress::parse(&checkpoint.records_ref.digest).expect("address");
assert!(
!swept.contains(&address),
"swept a live checkpoint's records"
);
assert!(
blobs.stat(&address).await.expect("stat").is_some(),
"a live checkpoint's records were removed"
);
}
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_store_with_no_blob_home_takes_no_snapshots() {
let maintenance = maintenance_pool().await;
let (name, store) = fresh_store(&maintenance, "cpnoblobs", 8).await;
sqlx::query("UPDATE gwk_internal.writer SET checkpoint_at = now() - interval '1 hour'")
.execute(store.pool())
.await
.expect("backdate the barrier");
apply(&store, "overdue", task("t-1")).await;
assert!(store.blobs().is_none());
assert!(checkpoints(&store).await.is_empty());
drop_database(&maintenance, &name).await;
}
async fn checkpoints(store: &PgEventStore) -> Vec<gwk_domain::checkpoint::Checkpoint> {
let mut conn = store.pool().acquire().await.expect("connection");
checkpoint::checkpoints(&mut conn)
.await
.expect("read checkpoints")
}
#[tokio::test]
#[ignore = "needs a PostgreSQL; see tests/common/mod.rs"]
async fn a_clean_stop_leaves_a_checkpoint_the_next_start_can_verify() {
let maintenance = maintenance_pool().await;
let (name, root, store) = checkpointing_store(&maintenance, "cpshutdown").await;
populate(&store).await;
let watermark = store
.watermark()
.await
.expect("watermark")
.expect("a populated log");
assert!(
checkpoints(&store).await.is_empty(),
"the barrier fired on its own — this case would prove nothing"
);
let dir = runtime_dir("cpshutdown");
let path = dir.join("gwk.sock");
let listener = Listener::bind(&path).await.expect("bind");
let daemon = Arc::new(Daemon::new(store, TEST_REVISION.to_owned()).expect("daemon"));
let (stop, stopped) = tokio::sync::oneshot::channel::<()>();
let serving = tokio::spawn(async move {
serve::run(listener, daemon, async move {
let _ = stopped.await;
})
.await
});
let _ = stop.send(());
let report = serving.await.expect("join").expect("a clean stop");
assert_eq!(
report.checkpoint,
Some(watermark),
"a clean stop must snapshot AT the watermark"
);
assert_eq!(report.checkpoint_error, None);
let pool = connect_pool(&secret(&name), 4).await.expect("connect");
let next = PgEventStore::open(pool).await.expect("open");
let blobs = blob_store_with(&next, &root, TEST_KEK).await;
let next = next.with_blobs(blobs);
match next.recover().await.expect("recover").verdict {
Verdict::Verified { anchor } => assert_eq!(anchor, watermark),
other => panic!("a clean stop must restart into Verified, got {other:?}"),
}
drop_database(&maintenance, &name).await;
let _ = std::fs::remove_dir_all(&root);
let _ = std::fs::remove_dir_all(&dir);
}