diff --git a/examples/idempotent_reservation.rs b/examples/idempotent_reservation.rs index b6a3ea8..d9e56dc 100644 --- a/examples/idempotent_reservation.rs +++ b/examples/idempotent_reservation.rs @@ -1,5 +1,5 @@ use serde::{Deserialize, Serialize}; -use std::{collections::HashMap, error::Error}; +use std::{collections::HashMap, error::Error, fmt}; use trogon_eventstore::{ AppendToStreamOptions, Client, ClientSettings, CurrentRevision, Error as ClientError, EventData, ReadStreamOptions, StreamState, WriteResult, @@ -38,7 +38,7 @@ struct InventoryCreated { available: u32, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] struct InventoryReserved { operation_id: OperationId, reservation_id: ReservationId, @@ -47,7 +47,7 @@ struct InventoryReserved { quantity: u32, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] struct InventoryReleased { operation_id: OperationId, reservation_id: ReservationId, @@ -55,62 +55,103 @@ struct InventoryReleased { quantity: u32, } +#[derive(Clone, Debug, Eq, PartialEq)] +enum ProcessedOperation { + Reserved { + event_revision: u64, + event: InventoryReserved, + }, + Released { + event_revision: u64, + event: InventoryReleased, + }, +} + +impl ProcessedOperation { + fn event_revision(&self) -> u64 { + match self { + Self::Reserved { event_revision, .. } | Self::Released { event_revision, .. } => { + *event_revision + } + } + } +} + #[derive(Debug)] struct InventoryState { available: u32, reservations: HashMap, + // Folding this index from the inventory events keeps idempotency and inventory in one write. + processed_operations: HashMap, } #[derive(Clone)] -struct ReservationAttempt { - client_name: &'static str, - operation_id: OperationId, - reservation_id: ReservationId, - event: EventData, -} +struct ReserveCommand(InventoryReserved); -impl ReservationAttempt { - fn new(client_name: &'static str, sku: &str) -> Result> { +impl ReserveCommand { + fn new(client_name: &str, sku: &str) -> Self { let operation_id = OperationId::new(); let reservation_id = ReservationId::new(); - let reservation = InventoryReserved { + Self(InventoryReserved { operation_id, reservation_id, client: client_name.to_owned(), sku: sku.to_owned(), quantity: 1, - }; - - Ok(Self { - client_name, - operation_id, - reservation_id, - event: EventData::json(INVENTORY_RESERVED_EVENT_TYPE, &reservation)?.id(operation_id.0), }) } + + fn event(&self) -> Result> { + Ok(EventData::json(INVENTORY_RESERVED_EVENT_TYPE, &self.0)?.id(self.0.operation_id.0)) + } } #[derive(Clone)] -struct ReleaseAttempt { - operation_id: OperationId, - event: EventData, -} +struct ReleaseCommand(InventoryReleased); -impl ReleaseAttempt { - fn new(sku: &str, reservation_id: ReservationId) -> Result> { +impl ReleaseCommand { + fn new(sku: &str, reservation_id: ReservationId) -> Self { let operation_id = OperationId::new(); - let release = InventoryReleased { + Self(InventoryReleased { operation_id, reservation_id, sku: sku.to_owned(), quantity: 1, - }; - - Ok(Self { - operation_id, - event: EventData::json(INVENTORY_RELEASED_EVENT_TYPE, &release)?.id(operation_id.0), }) } + + fn event(&self) -> Result> { + Ok(EventData::json(INVENTORY_RELEASED_EVENT_TYPE, &self.0)?.id(self.0.operation_id.0)) + } +} + +#[derive(Debug)] +struct OperationIdConflict(OperationId); + +impl fmt::Display for OperationIdConflict { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(formatter, "operation ID {:?} was reused with different content", self.0) + } +} + +impl Error for OperationIdConflict {} + +#[derive(Debug)] +struct InvalidInventoryCommand(&'static str); + +impl fmt::Display for InvalidInventoryCommand { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(self.0) + } +} + +impl Error for InvalidInventoryCommand {} + +#[derive(Debug)] +struct CommandResult { + outcome: ProcessedOperation, + appended: bool, + write: Option, } async fn read_inventory( @@ -123,6 +164,7 @@ async fn read_inventory( let mut state = InventoryState { available: 0, reservations: HashMap::new(), + processed_operations: HashMap::new(), }; let mut revision = None; @@ -146,6 +188,19 @@ async fn read_inventory( .insert(reserved.reservation_id, reserved.quantity) .is_none() ); + let operation_id = reserved.operation_id; + assert!( + state + .processed_operations + .insert( + operation_id, + ProcessedOperation::Reserved { + event_revision: event.revision, + event: reserved, + }, + ) + .is_none() + ); } INVENTORY_RELEASED_EVENT_TYPE => { let released = event.as_json::()?; @@ -156,6 +211,19 @@ async fn read_inventory( .expect("the released reservation to be active"); assert_eq!(quantity, released.quantity); state.available += released.quantity; + let operation_id = released.operation_id; + assert!( + state + .processed_operations + .insert( + operation_id, + ProcessedOperation::Released { + event_revision: event.revision, + event: released, + }, + ) + .is_none() + ); } event_type => panic!("unexpected inventory event type: {event_type}"), } @@ -185,6 +253,90 @@ async fn append_at( client.append_to_stream(stream, &options, event).await } +async fn reserve( + client: &Client, + stream: &str, + command: &ReserveCommand, +) -> Result> { + let (state, revision) = read_inventory(client, stream).await?; + + if let Some(previous) = state + .processed_operations + .get(&command.0.operation_id) + { + return match previous { + ProcessedOperation::Reserved { event, .. } if event == &command.0 => { + Ok(CommandResult { + outcome: previous.clone(), + appended: false, + write: None, + }) + } + _ => Err(Box::new(OperationIdConflict(command.0.operation_id))), + }; + } + + if state.available < command.0.quantity { + return Err(Box::new(InvalidInventoryCommand( + "the requested inventory is not available", + ))); + } + + let write = append_at(client, stream, revision, command.event()?).await?; + let outcome = ProcessedOperation::Reserved { + event_revision: write.next_expected_version, + event: command.0.clone(), + }; + + Ok(CommandResult { + outcome, + appended: true, + write: Some(write), + }) +} + +async fn release( + client: &Client, + stream: &str, + command: &ReleaseCommand, +) -> Result> { + let (state, revision) = read_inventory(client, stream).await?; + + if let Some(previous) = state + .processed_operations + .get(&command.0.operation_id) + { + return match previous { + ProcessedOperation::Released { event, .. } if event == &command.0 => { + Ok(CommandResult { + outcome: previous.clone(), + appended: false, + write: None, + }) + } + _ => Err(Box::new(OperationIdConflict(command.0.operation_id))), + }; + } + + if state.reservations.get(&command.0.reservation_id) != Some(&command.0.quantity) { + return Err(Box::new(InvalidInventoryCommand( + "the reservation is not active with the requested quantity", + ))); + } + + let write = append_at(client, stream, revision, command.event()?).await?; + let outcome = ProcessedOperation::Released { + event_revision: write.next_expected_version, + event: command.0.clone(), + }; + + Ok(CommandResult { + outcome, + appended: true, + write: Some(write), + }) +} + #[tokio::main] async fn main() -> Result<(), Box> { let connection_string = std::env::var(CONNECTION_STRING_ENV) @@ -216,21 +368,23 @@ async fn main() -> Result<(), Box> { assert_eq!(second_state.available, 1); assert_eq!(first_revision, second_revision); - let first_attempt = ReservationAttempt::new(FIRST_CLIENT, &sku)?; - let second_attempt = ReservationAttempt::new(SECOND_CLIENT, &sku)?; + let first_attempt = ReserveCommand::new(FIRST_CLIENT, &sku); + let second_attempt = ReserveCommand::new(SECOND_CLIENT, &sku); + assert_ne!(first_attempt.0.operation_id.0, first_attempt.0.reservation_id.0); + assert_ne!(second_attempt.0.operation_id.0, second_attempt.0.reservation_id.0); // The shared expected revision serializes decisions made from the same inventory state. let (first_result, second_result) = tokio::join!( append_at( &first_client, inventory_stream.as_str(), first_revision, - first_attempt.event.clone(), + first_attempt.event()?, ), append_at( &second_client, inventory_stream.as_str(), second_revision, - second_attempt.event.clone(), + second_attempt.event()?, ), ); @@ -258,135 +412,136 @@ async fn main() -> Result<(), Box> { let (sold_out, sold_out_revision) = read_inventory(loser_client, inventory_stream.as_str()).await?; assert_eq!(sold_out.available, 0); - assert_eq!(sold_out.reservations, [(winner.reservation_id, 1)].into()); + assert_eq!(sold_out.reservations, [(winner.0.reservation_id, 1)].into()); + assert_eq!(sold_out_revision, 1); - let winner_release = ReleaseAttempt::new(&sku, winner.reservation_id)?; - assert_ne!(winner_release.operation_id, winner.operation_id); - let winner_release_write = append_at( - winner_client, - inventory_stream.as_str(), - sold_out_revision, - winner_release.event.clone(), - ) - .await?; - assert_eq!(winner_release_write.next_expected_version, 2); + let winner_release = ReleaseCommand::new(&sku, winner.0.reservation_id); + assert_ne!(winner_release.0.operation_id, winner.0.operation_id); + let winner_release_result = + release(winner_client, inventory_stream.as_str(), &winner_release).await?; + assert!(winner_release_result.appended); + assert_eq!(winner_release_result.outcome.event_revision(), 2); + let winner_release_outcome = winner_release_result.outcome.clone(); let (released, released_revision) = read_inventory(loser_client, inventory_stream.as_str()).await?; assert_eq!(released.available, 1); assert!(released.reservations.is_empty()); - let loser_write = append_at( - loser_client, - inventory_stream.as_str(), - released_revision, - loser.event.clone(), - ) - .await?; - assert_eq!(loser_write.next_expected_version, 3); + let loser_result = reserve(loser_client, inventory_stream.as_str(), loser).await?; + assert!(loser_result.appended); + assert_eq!(loser_result.outcome.event_revision(), 3); + let loser_outcome = loser_result.outcome.clone(); - let loser_release = ReleaseAttempt::new(&sku, loser.reservation_id)?; - assert_ne!(loser_release.operation_id, loser.operation_id); - let loser_release_write = append_at( - loser_client, - inventory_stream.as_str(), - loser_write.next_expected_version, - loser_release.event.clone(), - ) - .await?; - assert_eq!(loser_release_write.next_expected_version, 4); + let loser_release = ReleaseCommand::new(&sku, loser.0.reservation_id); + assert_ne!(loser_release.0.operation_id, loser.0.operation_id); + let loser_release_result = + release(loser_client, inventory_stream.as_str(), &loser_release).await?; + assert!(loser_release_result.appended); + assert_eq!(loser_release_result.outcome.event_revision(), 4); + let loser_release_outcome = loser_release_result.outcome.clone(); let (available_again, available_again_revision) = read_inventory(winner_client, inventory_stream.as_str()).await?; assert_eq!(available_again.available, 1); assert!(available_again.reservations.is_empty()); + assert_eq!(available_again_revision, 4); - let winner_again = ReservationAttempt::new(winner.client_name, &sku)?; - let winner_again_write = append_at( - winner_client, - inventory_stream.as_str(), - available_again_revision, - winner_again.event.clone(), - ) - .await?; - assert_eq!(winner_again_write.next_expected_version, 5); + let winner_again = ReserveCommand::new(&winner.0.client, &sku); + let winner_again_result = + reserve(winner_client, inventory_stream.as_str(), &winner_again).await?; + assert!(winner_again_result.appended); + assert_eq!(winner_again_result.outcome.event_revision(), 5); + let winner_again_outcome = winner_again_result.outcome.clone(); - let winner_again_release = ReleaseAttempt::new(&sku, winner_again.reservation_id)?; - let winner_again_release_write = append_at( + let winner_again_release = ReleaseCommand::new(&sku, winner_again.0.reservation_id); + let winner_again_release_result = release( winner_client, inventory_stream.as_str(), - winner_again_write.next_expected_version, - winner_again_release.event.clone(), + &winner_again_release, ) .await?; - assert_eq!(winner_again_release_write.next_expected_version, 6); + assert!(winner_again_release_result.appended); + assert_eq!(winner_again_release_result.outcome.event_revision(), 6); + let winner_again_release_outcome = winner_again_release_result.outcome.clone(); - // Durable retries retain the original expected revision and event ID after later writes. + // The original write tuple is enough only when a transport retry retained it unchanged. let winner_retry = append_at( winner_client, inventory_stream.as_str(), first_revision, - winner.event.clone(), - ) - .await?; - let winner_release_retry = append_at( - winner_client, - inventory_stream.as_str(), - sold_out_revision, - winner_release.event, - ) - .await?; - let loser_retry = append_at( - loser_client, - inventory_stream.as_str(), - released_revision, - loser.event.clone(), - ) - .await?; - let loser_release_retry = append_at( - loser_client, - inventory_stream.as_str(), - loser_write.next_expected_version, - loser_release.event, + winner.event()?, ) .await?; - let winner_again_retry = append_at( + assert_eq!(winner_retry.position, winner_write.position); + + let winner_outcome = ProcessedOperation::Reserved { + event_revision: winner_write.next_expected_version, + event: winner.0.clone(), + }; + + // Delayed redelivery has lost the old revision, so the folded operation outcome is authoritative. + let winner_replay = reserve(winner_client, inventory_stream.as_str(), winner).await?; + let winner_release_replay = + release(winner_client, inventory_stream.as_str(), &winner_release).await?; + let loser_replay = reserve(loser_client, inventory_stream.as_str(), loser).await?; + let loser_release_replay = + release(loser_client, inventory_stream.as_str(), &loser_release).await?; + let winner_again_replay = + reserve(winner_client, inventory_stream.as_str(), &winner_again).await?; + let winner_again_release_replay = release( winner_client, inventory_stream.as_str(), - available_again_revision, - winner_again.event, + &winner_again_release, ) .await?; - let winner_again_release_retry = append_at( + + for (replay, original) in [ + (winner_replay, winner_outcome), + (winner_release_replay, winner_release_outcome), + (loser_replay, loser_outcome), + (loser_release_replay, loser_release_outcome), + (winner_again_replay, winner_again_outcome), + (winner_again_release_replay, winner_again_release_outcome), + ] { + assert!(!replay.appended); + assert!(replay.write.is_none()); + assert_eq!(replay.outcome, original); + } + + let conflicting_command = ReserveCommand(InventoryReserved { + operation_id: winner.0.operation_id, + reservation_id: ReservationId::new(), + client: winner.0.client.clone(), + sku: sku.clone(), + quantity: 1, + }); + let conflict = reserve( winner_client, inventory_stream.as_str(), - winner_again_write.next_expected_version, - winner_again_release.event, + &conflicting_command, ) - .await?; - assert_eq!(winner_retry.position, winner_write.position); - assert_eq!(winner_release_retry.position, winner_release_write.position); - assert_eq!(loser_retry.position, loser_write.position); - assert_eq!(loser_release_retry.position, loser_release_write.position); - assert_eq!(winner_again_retry.position, winner_again_write.position); - assert_eq!( - winner_again_release_retry.position, - winner_again_release_write.position - ); + .await + .expect_err("reusing an operation ID with different content must fail"); + assert!(conflict.downcast_ref::().is_some()); let (final_state, final_revision) = read_inventory(&first_client, inventory_stream.as_str()).await?; assert_eq!(final_revision, 6); assert_eq!(final_state.available, 1); assert!(final_state.reservations.is_empty()); + assert_eq!(final_state.processed_operations.len(), 6); println!( "{} won the first race and released reservation {}; {} then reserved and released after reloading revision {}", - winner.client_name, winner.reservation_id.0, loser.client_name, released_revision + winner.0.client, winner.0.reservation_id.0, loser.0.client, released_revision + ); + println!( + "{} completed a second reserve/release cycle; six delayed command replays returned their original outcomes after revision {final_revision}", + winner.0.client ); println!( - "{} completed a second reserve/release cycle; all six retries remained single events after revision {final_revision}", - winner.client_name + "The example retains every processed operation in the inventory stream; production systems must budget for that unbounded history and index cost" ); Ok(()) diff --git a/trogon-eventstore/tests/api/idempotency.rs b/trogon-eventstore/tests/api/idempotency.rs index 818cd58..c20e1ba 100644 --- a/trogon-eventstore/tests/api/idempotency.rs +++ b/trogon-eventstore/tests/api/idempotency.rs @@ -25,64 +25,105 @@ impl ReservationId { } } -#[derive(Clone)] -struct ReservationAttempt { +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +struct InventoryReserved { + state: String, + client: String, operation_id: OperationId, reservation_id: ReservationId, - event: EventData, + quantity: u32, } -impl ReservationAttempt { +#[derive(Clone)] +struct ReserveCommand(InventoryReserved); + +impl ReserveCommand { fn new(client: &str) -> Self { let operation_id = OperationId::new(); let reservation_id = ReservationId::new(); - Self { + Self(InventoryReserved { + state: "reserved".to_owned(), + client: client.to_owned(), operation_id, reservation_id, - event: EventData::json( - "inventory-reserved", - &json!({ - "state": "reserved", - "client": client, - "operation_id": operation_id, - "reservation_id": reservation_id, - "quantity": 1, - }), - ) + quantity: 1, + }) + } + + fn event(&self) -> EventData { + EventData::json("inventory-reserved", &self.0) .unwrap() - .id(operation_id.0), - } + .id(self.0.operation_id.0) } } -#[derive(Clone)] -struct ReleaseAttempt { +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +struct InventoryReleased { + state: String, operation_id: OperationId, - event: EventData, + reservation_id: ReservationId, + quantity: u32, } -impl ReleaseAttempt { +#[derive(Clone)] +struct ReleaseCommand(InventoryReleased); + +impl ReleaseCommand { fn new(reservation_id: ReservationId) -> Self { let operation_id = OperationId::new(); - Self { + Self(InventoryReleased { + state: "released".to_owned(), operation_id, - event: EventData::json( - "inventory-released", - &json!({ - "state": "released", - "operation_id": operation_id, - "reservation_id": reservation_id, - "quantity": 1, - }), - ) + reservation_id, + quantity: 1, + }) + } + + fn event(&self) -> EventData { + EventData::json("inventory-released", &self.0) .unwrap() - .id(operation_id.0), + .id(self.0.operation_id.0) + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +enum ProcessedOperation { + Reserved { + event_revision: u64, + event: InventoryReserved, + }, + Released { + event_revision: u64, + event: InventoryReleased, + }, +} + +impl ProcessedOperation { + fn event_revision(&self) -> u64 { + match self { + Self::Reserved { event_revision, .. } | Self::Released { event_revision, .. } => { + *event_revision + } } } } +#[derive(Debug)] +struct InventoryState { + available: u32, + reservations: HashMap, + processed_operations: HashMap, +} + +#[derive(Debug)] +struct CommandResult { + outcome: ProcessedOperation, + appended: bool, + write: Option, +} + fn event(id: Uuid, state: &str) -> EventData { EventData::json("inventory-reservation", &json!({ "state": state })) .unwrap() @@ -111,42 +152,161 @@ async fn stored_events(client: &Client, stream: &str) -> eyre::Result (u32, HashMap) { - let mut available = 0; - let mut reservations = HashMap::new(); +fn fold_inventory(events: &[(Uuid, u64, Value)]) -> InventoryState { + let mut state = InventoryState { + available: 0, + reservations: HashMap::new(), + processed_operations: HashMap::new(), + }; - for (event_id, _, event) in events { + for (event_id, event_revision, event) in events { let quantity = event .get("quantity") .and_then(Value::as_u64) .unwrap_or_default() as u32; match event["state"].as_str().unwrap() { - "created" => available = event["available"].as_u64().unwrap() as u32, + "created" => state.available = event["available"].as_u64().unwrap() as u32, "reserved" => { - let operation_id: OperationId = - serde_json::from_value(event["operation_id"].clone()).unwrap(); - assert_eq!(*event_id, operation_id.0); - assert!(available >= quantity); - available -= quantity; - let reservation_id = - serde_json::from_value(event["reservation_id"].clone()).unwrap(); - assert!(reservations.insert(reservation_id, quantity).is_none()); + let reserved: InventoryReserved = serde_json::from_value(event.clone()).unwrap(); + assert_eq!(*event_id, reserved.operation_id.0); + assert!(state.available >= quantity); + state.available -= quantity; + assert!( + state + .reservations + .insert(reserved.reservation_id, quantity) + .is_none() + ); + assert!( + state + .processed_operations + .insert( + reserved.operation_id, + ProcessedOperation::Reserved { + event_revision: *event_revision, + event: reserved, + }, + ) + .is_none() + ); } "released" => { - let operation_id: OperationId = - serde_json::from_value(event["operation_id"].clone()).unwrap(); - assert_eq!(*event_id, operation_id.0); - let reservation_id = - serde_json::from_value(event["reservation_id"].clone()).unwrap(); - assert_eq!(reservations.remove(&reservation_id), Some(quantity)); - available += quantity; + let released: InventoryReleased = serde_json::from_value(event.clone()).unwrap(); + assert_eq!(*event_id, released.operation_id.0); + assert_eq!( + state.reservations.remove(&released.reservation_id), + Some(released.quantity) + ); + state.available += released.quantity; + assert!( + state + .processed_operations + .insert( + released.operation_id, + ProcessedOperation::Released { + event_revision: *event_revision, + event: released, + }, + ) + .is_none() + ); } state => panic!("unexpected inventory state: {state}"), } } - (available, reservations) + state +} + +async fn reserve( + client: &Client, + stream: &str, + command: &ReserveCommand, +) -> eyre::Result { + let events = stored_events(client, stream).await?; + let state = fold_inventory(&events); + + if let Some(previous) = state.processed_operations.get(&command.0.operation_id) { + return match previous { + ProcessedOperation::Reserved { event, .. } if event == &command.0 => { + Ok(CommandResult { + outcome: previous.clone(), + appended: false, + write: None, + }) + } + _ => eyre::bail!("operation ID was reused with different content"), + }; + } + + if state.available < command.0.quantity { + eyre::bail!("the requested inventory is not available"); + } + + let revision = events.last().unwrap().1; + let write = append( + client, + stream, + StreamState::StreamRevision(revision), + command.event(), + ) + .await?; + let outcome = ProcessedOperation::Reserved { + event_revision: write.next_expected_version, + event: command.0.clone(), + }; + + Ok(CommandResult { + outcome, + appended: true, + write: Some(write), + }) +} + +async fn release( + client: &Client, + stream: &str, + command: &ReleaseCommand, +) -> eyre::Result { + let events = stored_events(client, stream).await?; + let state = fold_inventory(&events); + + if let Some(previous) = state.processed_operations.get(&command.0.operation_id) { + return match previous { + ProcessedOperation::Released { event, .. } if event == &command.0 => { + Ok(CommandResult { + outcome: previous.clone(), + appended: false, + write: None, + }) + } + _ => eyre::bail!("operation ID was reused with different content"), + }; + } + + if state.reservations.get(&command.0.reservation_id) != Some(&command.0.quantity) { + eyre::bail!("the reservation is not active with the requested quantity"); + } + + let revision = events.last().unwrap().1; + let write = append( + client, + stream, + StreamState::StreamRevision(revision), + command.event(), + ) + .await?; + let outcome = ProcessedOperation::Released { + event_revision: write.next_expected_version, + event: command.0.clone(), + }; + + Ok(CommandResult { + outcome, + appended: true, + write: Some(write), + }) } fn is_revision_conflict(error: &Error, expected: u64, current: u64) -> bool { @@ -335,18 +495,21 @@ async fn competing_reservations_are_serialized_and_retryable(client: &Client) -> let first_view = stored_events(&first_client, &stream).await?; let second_view = stored_events(&second_client, &stream).await?; - assert_eq!(fold_inventory(&first_view), (1, HashMap::new())); - assert_eq!(fold_inventory(&second_view), (1, HashMap::new())); + assert_eq!(fold_inventory(&first_view).available, 1); + assert_eq!(fold_inventory(&second_view).available, 1); let first_revision = first_view.last().unwrap().1; let second_revision = second_view.last().unwrap().1; assert_eq!(first_revision, second_revision); - let first_attempt = ReservationAttempt::new("checkout-a"); - let second_attempt = ReservationAttempt::new("checkout-b"); - assert_ne!(first_attempt.operation_id.0, first_attempt.reservation_id.0); + let first_attempt = ReserveCommand::new("checkout-a"); + let second_attempt = ReserveCommand::new("checkout-b"); + assert_ne!( + first_attempt.0.operation_id.0, + first_attempt.0.reservation_id.0 + ); assert_ne!( - second_attempt.operation_id.0, - second_attempt.reservation_id.0 + second_attempt.0.operation_id.0, + second_attempt.0.reservation_id.0 ); let (first_result, second_result) = tokio::join!( @@ -354,13 +517,13 @@ async fn competing_reservations_are_serialized_and_retryable(client: &Client) -> &first_client, &stream, StreamState::StreamRevision(first_revision), - first_attempt.event.clone(), + first_attempt.event(), ), append( &second_client, &stream, StreamState::StreamRevision(second_revision), - second_attempt.event.clone(), + second_attempt.event(), ), ); @@ -386,140 +549,128 @@ async fn competing_reservations_are_serialized_and_retryable(client: &Client) -> }; let sold_out_view = stored_events(loser_client, &stream).await?; - assert_eq!( - fold_inventory(&sold_out_view), - (0, [(winner.reservation_id, 1)].into()) - ); + let sold_out = fold_inventory(&sold_out_view); + assert_eq!(sold_out.available, 0); + assert_eq!(sold_out.reservations, [(winner.0.reservation_id, 1)].into()); let sold_out_revision = sold_out_view.last().unwrap().1; assert_eq!(sold_out_revision, 1); - let winner_release = ReleaseAttempt::new(winner.reservation_id); - assert_ne!(winner_release.operation_id, winner.operation_id); - assert_ne!(winner_release.operation_id.0, winner.reservation_id.0); - let winner_release_write = append( - winner_client, - &stream, - StreamState::StreamRevision(sold_out_revision), - winner_release.event.clone(), - ) - .await?; + let winner_release = ReleaseCommand::new(winner.0.reservation_id); + assert_ne!(winner_release.0.operation_id, winner.0.operation_id); + assert_ne!(winner_release.0.operation_id.0, winner.0.reservation_id.0); + let winner_release_result = release(winner_client, &stream, &winner_release).await?; + assert!(winner_release_result.appended); + assert_eq!(winner_release_result.outcome.event_revision(), 2); + let winner_release_outcome = winner_release_result.outcome.clone(); let released_view = stored_events(loser_client, &stream).await?; - assert_eq!(fold_inventory(&released_view), (1, HashMap::new())); + let released = fold_inventory(&released_view); + assert_eq!(released.available, 1); + assert!(released.reservations.is_empty()); let released_revision = released_view.last().unwrap().1; assert_eq!(released_revision, 2); - let loser_write = append( - loser_client, - &stream, - StreamState::StreamRevision(released_revision), - loser.event.clone(), - ) - .await?; - assert_eq!(loser_write.next_expected_version, 3); + let loser_result = reserve(loser_client, &stream, loser).await?; + assert!(loser_result.appended); + assert_eq!(loser_result.outcome.event_revision(), 3); + let loser_outcome = loser_result.outcome.clone(); - let loser_release = ReleaseAttempt::new(loser.reservation_id); - assert_ne!(loser_release.operation_id, loser.operation_id); - let loser_release_write = append( - loser_client, - &stream, - StreamState::StreamRevision(loser_write.next_expected_version), - loser_release.event.clone(), - ) - .await?; - assert_eq!(loser_release_write.next_expected_version, 4); + let loser_release = ReleaseCommand::new(loser.0.reservation_id); + assert_ne!(loser_release.0.operation_id, loser.0.operation_id); + let loser_release_result = release(loser_client, &stream, &loser_release).await?; + assert!(loser_release_result.appended); + assert_eq!(loser_release_result.outcome.event_revision(), 4); + let loser_release_outcome = loser_release_result.outcome.clone(); let available_again_view = stored_events(winner_client, &stream).await?; - assert_eq!(fold_inventory(&available_again_view), (1, HashMap::new())); + let available_again = fold_inventory(&available_again_view); + assert_eq!(available_again.available, 1); + assert!(available_again.reservations.is_empty()); let available_again_revision = available_again_view.last().unwrap().1; assert_eq!(available_again_revision, 4); - let winner_again = ReservationAttempt::new("checkout-winner-again"); - let winner_again_write = append( - winner_client, - &stream, - StreamState::StreamRevision(available_again_revision), - winner_again.event.clone(), - ) - .await?; - assert_eq!(winner_again_write.next_expected_version, 5); + let winner_again = ReserveCommand::new("checkout-winner-again"); + let winner_again_result = reserve(winner_client, &stream, &winner_again).await?; + assert!(winner_again_result.appended); + assert_eq!(winner_again_result.outcome.event_revision(), 5); + let winner_again_outcome = winner_again_result.outcome.clone(); - let winner_again_release = ReleaseAttempt::new(winner_again.reservation_id); - let winner_again_release_write = append( - winner_client, - &stream, - StreamState::StreamRevision(winner_again_write.next_expected_version), - winner_again_release.event.clone(), - ) - .await?; - assert_eq!(winner_again_release_write.next_expected_version, 6); + let winner_again_release = ReleaseCommand::new(winner_again.0.reservation_id); + let winner_again_release_result = + release(winner_client, &stream, &winner_again_release).await?; + assert!(winner_again_release_result.appended); + assert_eq!(winner_again_release_result.outcome.event_revision(), 6); + let winner_again_release_outcome = winner_again_release_result.outcome.clone(); + // The server can recognize a transport retry because the original write tuple is unchanged. let winner_retry = append( winner_client, &stream, StreamState::StreamRevision(first_revision), - winner.event.clone(), - ) - .await?; - let winner_release_retry = append( - winner_client, - &stream, - StreamState::StreamRevision(sold_out_revision), - winner_release.event, - ) - .await?; - let loser_retry = append( - loser_client, - &stream, - StreamState::StreamRevision(released_revision), - loser.event.clone(), - ) - .await?; - let loser_release_retry = append( - loser_client, - &stream, - StreamState::StreamRevision(loser_write.next_expected_version), - loser_release.event, - ) - .await?; - let winner_again_retry = append( - winner_client, - &stream, - StreamState::StreamRevision(available_again_revision), - winner_again.event, - ) - .await?; - let winner_again_release_retry = append( - winner_client, - &stream, - StreamState::StreamRevision(winner_again_write.next_expected_version), - winner_again_release.event, + winner.event(), ) .await?; assert_eq!(winner_retry.position, winner_write.position); - assert_eq!(winner_release_retry.position, winner_release_write.position); - assert_eq!(loser_retry.position, loser_write.position); - assert_eq!(loser_release_retry.position, loser_release_write.position); - assert_eq!(winner_again_retry.position, winner_again_write.position); + + let winner_outcome = ProcessedOperation::Reserved { + event_revision: winner_write.next_expected_version, + event: winner.0.clone(), + }; + + // Reconstructed commands need the outcome retained in the same authoritative stream. + let winner_replay = reserve(winner_client, &stream, winner).await?; + let winner_release_replay = release(winner_client, &stream, &winner_release).await?; + let loser_replay = reserve(loser_client, &stream, loser).await?; + let loser_release_replay = release(loser_client, &stream, &loser_release).await?; + let winner_again_replay = reserve(winner_client, &stream, &winner_again).await?; + let winner_again_release_replay = + release(winner_client, &stream, &winner_again_release).await?; + + for (replay, original) in [ + (winner_replay, winner_outcome), + (winner_release_replay, winner_release_outcome), + (loser_replay, loser_outcome), + (loser_release_replay, loser_release_outcome), + (winner_again_replay, winner_again_outcome), + (winner_again_release_replay, winner_again_release_outcome), + ] { + assert!(!replay.appended); + assert!(replay.write.is_none()); + assert_eq!(replay.outcome, original); + } + + let conflicting_command = ReserveCommand(InventoryReserved { + state: "reserved".to_owned(), + client: winner.0.client.clone(), + operation_id: winner.0.operation_id, + reservation_id: ReservationId::new(), + quantity: 1, + }); + let conflict = reserve(winner_client, &stream, &conflicting_command) + .await + .expect_err("reusing an operation ID with different content must fail"); assert_eq!( - winner_again_release_retry.position, - winner_again_release_write.position + conflict.to_string(), + "operation ID was reused with different content" ); let final_view = stored_events(&first_client, &stream).await?; assert_eq!(final_view.len(), 7); assert_eq!(final_view.last().unwrap().1, 6); - assert_eq!(fold_inventory(&final_view), (1, HashMap::new())); + let final_state = fold_inventory(&final_view); + assert_eq!(final_state.available, 1); + assert!(final_state.reservations.is_empty()); + assert_eq!(final_state.processed_operations.len(), 6); assert_eq!( final_view.iter().map(|(id, _, _)| *id).collect::>(), [ created_operation_id, - winner.operation_id.0, - winner_release.operation_id.0, - loser.operation_id.0, - loser_release.operation_id.0, - winner_again.operation_id.0, - winner_again_release.operation_id.0, + winner.0.operation_id.0, + winner_release.0.operation_id.0, + loser.0.operation_id.0, + loser_release.0.operation_id.0, + winner_again.0.operation_id.0, + winner_again_release.0.operation_id.0, ] );