#![allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
clippy::print_stdout
)]
mod support;
use turnframe_core::case::{CaseKey, CaseRef};
use turnframe_core::ids::{AccountId, CaseRevision, CommandId, EventId};
use turnframe_core::interaction::InteractionStatus;
use turnframe_store::events::{EventBatch, EventJournalReader, EventJournalWriter, StoredEvent};
use turnframe_store::interaction::{
InteractionReader, InteractionWriter, InvalidationReason, ResolutionOutcome,
};
const SCHEMA: &str = "tf_test_rollback";
const WORKFLOW: &str = "trip";
fn event(
case: &CaseKey,
ordinal: u8,
event_type: &str,
payload: serde_json::Value,
) -> turnframe_core::event::CommittedEvent<serde_json::Value> {
let mut seed = [0_u8; 16];
for (slot, byte) in seed.iter_mut().zip(case.case_id.as_str().bytes().cycle()) {
*slot = byte;
}
seed[15] = ordinal;
turnframe_core::event::CommittedEvent {
event_id: EventId::from(uuid::Uuid::from_bytes(seed)),
event_type: event_type.to_owned(),
occurred_at: chrono::DateTime::from_timestamp(1_700_000_000, 0).expect("a fixed instant"),
payload,
}
}
async fn history_written_by_the_newer_version(
store: &turnframe_store_postgres::PgStores,
account: &AccountId,
case: &CaseKey,
) -> Vec<StoredEvent> {
let command = CommandId::new();
let batch = EventBatch::new(
account.clone(),
case.clone(),
command,
CaseRevision(2),
vec![
event(
case,
1,
"trip.name_set",
serde_json::json!({"value": "Lisbon", "set_by_version": "2"}),
),
event(
case,
2,
"trip.travel_date_set",
serde_json::json!({"value": "2026-10-31"}),
),
],
);
EventJournalWriter::append(store, batch)
.await
.expect("the newer version commits its events");
EventJournalReader::list_since(store, account, case, CaseRevision(0), 64)
.await
.expect("a server running the previous version reads the same ledger")
}
#[tokio::test]
async fn the_previous_version_reads_a_case_the_newer_one_wrote() {
let Some(url) = support::database_url() else {
support::skipped("the_previous_version_reads_a_case_the_newer_one_wrote");
return;
};
let store = support::co_tenant_of(&url, SCHEMA).await;
let account = support::unique_account("acct-rollback-reads");
let case = CaseKey::new(WORKFLOW, "trip-1");
let read_back = history_written_by_the_newer_version(&store, &account, &case).await;
assert_eq!(
read_back.len(),
2,
"a rollback reads the whole history, not the part its own version wrote"
);
assert_eq!(
read_back
.iter()
.map(|stored| stored.event_type.clone())
.collect::<Vec<_>>(),
vec!["trip.name_set", "trip.travel_date_set"],
"in the order they were committed"
);
assert_eq!(
read_back[0].payload["value"], "Lisbon",
"with the payload unchanged: an event's meaning is not edited retroactively"
);
assert_eq!(
read_back[0].payload["set_by_version"], "2",
"including a field the previous version does not know, which it ignores rather than fails on"
);
}
#[tokio::test]
async fn a_rollback_neither_rewrites_nor_removes_an_event() {
let Some(url) = support::database_url() else {
support::skipped("a_rollback_neither_rewrites_nor_removes_an_event");
return;
};
let store = support::co_tenant_of(&url, SCHEMA).await;
let account = support::unique_account("acct-rollback-immutable");
let case = CaseKey::new(WORKFLOW, "trip-2");
let before = history_written_by_the_newer_version(&store, &account, &case).await;
let after = EventJournalReader::list_since(&store, &account, &case, CaseRevision(0), 64)
.await
.expect("the ledger answers");
assert_eq!(
before.iter().map(|e| e.event_id).collect::<Vec<_>>(),
after.iter().map(|e| e.event_id).collect::<Vec<_>>(),
"the same events, with the same identifiers, in the same order"
);
assert_eq!(
before.iter().map(|e| e.sequence).collect::<Vec<_>>(),
after.iter().map(|e| e.sequence).collect::<Vec<_>>(),
"and at the same positions, so a receipt rendered before the rollback still cites what it cited"
);
}
#[tokio::test]
async fn a_card_the_newer_version_left_open_is_still_answerable_afterwards() {
let Some(url) = support::database_url() else {
support::skipped("a_card_the_newer_version_left_open_is_still_answerable_afterwards");
return;
};
let store = support::co_tenant_of(&url, SCHEMA).await;
let account = support::unique_account("acct-rollback-card-open");
let case = CaseRef::new(WORKFLOW, "trip-3", CaseRevision(2));
let card = support::card(&account, case.clone(), support::epoch());
let id = card.id;
InteractionWriter::insert(&store, card)
.await
.expect("the newer version wrote a card");
let found = InteractionReader::get(&store, &account, &id)
.await
.expect("the card survives the rollback");
assert_eq!(
found.interaction.status,
InteractionStatus::Active,
"and it is still open, so the user is not left holding an unanswerable card"
);
let option = found
.interaction
.payload
.options
.first()
.expect("the stored card carries its options")
.id
.clone();
InteractionWriter::begin_resolution(
&store,
&account,
&id,
InteractionStatus::Active,
option,
turnframe_core::ids::TurnId::new(),
)
.await
.expect("the previous version can answer a card the newer one wrote");
InteractionWriter::finish_resolution(
&store,
&account,
&id,
ResolutionOutcome::Resolved { event_ids: vec![] },
)
.await
.expect("and settle it");
let settled = InteractionReader::get(&store, &account, &id)
.await
.expect("the card is readable after settling");
assert_eq!(settled.interaction.status, InteractionStatus::Resolved);
}
#[tokio::test]
async fn the_other_half_of_the_rule_is_invalidating_the_cards_instead() {
let Some(url) = support::database_url() else {
support::skipped("the_other_half_of_the_rule_is_invalidating_the_cards_instead");
return;
};
let store = support::co_tenant_of(&url, SCHEMA).await;
let account = support::unique_account("acct-rollback-card-invalidated");
let case = CaseRef::new(WORKFLOW, "trip-4", CaseRevision(2));
let card = support::card(&account, case.clone(), support::epoch());
let id = card.id;
InteractionWriter::insert(&store, card)
.await
.expect("the newer version wrote a card");
let invalidated = InteractionWriter::invalidate_case_cards(
&store,
&account,
&case.key(),
InvalidationReason::Administrative {
code: String::from("workflow_rolled_back"),
},
)
.await
.expect("the rollback invalidates what it cannot honour");
assert!(
invalidated.contains(&id),
"the card the newer version wrote is named as invalidated"
);
let found = InteractionReader::get(&store, &account, &id)
.await
.expect("the card is still there to read");
assert_eq!(
found.interaction.status,
InteractionStatus::Invalidated,
"invalidated rather than deleted, so the record of it having existed survives"
);
}
#[tokio::test]
async fn rolling_one_workflow_back_leaves_another_alone() {
let Some(url) = support::database_url() else {
support::skipped("rolling_one_workflow_back_leaves_another_alone");
return;
};
let store = support::co_tenant_of(&url, SCHEMA).await;
let account = support::unique_account("acct-rollback-scoped");
let trip = CaseRef::new(WORKFLOW, "trip-5", CaseRevision(2));
let traveler = CaseRef::new("traveler", "trav-1", CaseRevision(2));
let trip_card = support::card(&account, trip.clone(), support::epoch());
let traveler_card = support::card(&account, traveler.clone(), support::epoch());
let traveler_id = traveler_card.id;
InteractionWriter::insert(&store, trip_card)
.await
.expect("a card on the workflow being rolled back");
InteractionWriter::insert(&store, traveler_card)
.await
.expect("a card on a workflow that is not");
InteractionWriter::invalidate_case_cards(
&store,
&account,
&trip.key(),
InvalidationReason::Administrative {
code: String::from("workflow_rolled_back"),
},
)
.await
.expect("the rollback runs against one workflow's case");
let untouched = InteractionReader::get(&store, &account, &traveler_id)
.await
.expect("the other workflow's card is still there");
assert_eq!(
untouched.interaction.status,
InteractionStatus::Active,
"rolling one workflow back leaves another's cards open: versions are independent"
);
}