From 68f2be3fdbbfc7b57766a2ce7cc3d6cceabe5f2e Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Tue, 8 Sep 2026 08:48:53 -0400 Subject: [PATCH] feat(examples): demonstrate competing inventory lifecycles Signed-off-by: Yordis Prieto --- examples/idempotent_reservation.rs | 384 +++++++++++++++++++-- trogon-eventstore/tests/api/idempotency.rs | 342 +++++++++++++++++- 2 files changed, 689 insertions(+), 37 deletions(-) diff --git a/examples/idempotent_reservation.rs b/examples/idempotent_reservation.rs index 153449d..b6a3ea8 100644 --- a/examples/idempotent_reservation.rs +++ b/examples/idempotent_reservation.rs @@ -1,12 +1,36 @@ use serde::{Deserialize, Serialize}; -use std::error::Error; -use trogon_eventstore::{AppendToStreamOptions, Client, EventData, ReadStreamOptions, StreamState}; +use std::{collections::HashMap, error::Error}; +use trogon_eventstore::{ + AppendToStreamOptions, Client, ClientSettings, CurrentRevision, Error as ClientError, + EventData, ReadStreamOptions, StreamState, WriteResult, +}; use uuid::Uuid; const DEFAULT_CONNECTION_STRING: &str = "esdb://localhost:2113?tls=false"; const CONNECTION_STRING_ENV: &str = "TROGON_EVENTSTORE_CONNECTION_STRING"; const INVENTORY_CREATED_EVENT_TYPE: &str = "inventory-created"; const INVENTORY_RESERVED_EVENT_TYPE: &str = "inventory-reserved"; +const INVENTORY_RELEASED_EVENT_TYPE: &str = "inventory-released"; +const FIRST_CLIENT: &str = "checkout-a"; +const SECOND_CLIENT: &str = "checkout-b"; + +#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)] +struct OperationId(Uuid); + +impl OperationId { + fn new() -> Self { + Self(Uuid::new_v4()) + } +} + +#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)] +struct ReservationId(Uuid); + +impl ReservationId { + fn new() -> Self { + Self(Uuid::new_v4()) + } +} #[derive(Debug, Deserialize, Serialize)] struct InventoryCreated { @@ -16,65 +40,353 @@ struct InventoryCreated { #[derive(Debug, Deserialize, Serialize)] struct InventoryReserved { - operation_id: Uuid, - order_id: String, + operation_id: OperationId, + reservation_id: ReservationId, + client: String, sku: String, quantity: u32, } +#[derive(Debug, Deserialize, Serialize)] +struct InventoryReleased { + operation_id: OperationId, + reservation_id: ReservationId, + sku: String, + quantity: u32, +} + +#[derive(Debug)] +struct InventoryState { + available: u32, + reservations: HashMap, +} + +#[derive(Clone)] +struct ReservationAttempt { + client_name: &'static str, + operation_id: OperationId, + reservation_id: ReservationId, + event: EventData, +} + +impl ReservationAttempt { + fn new(client_name: &'static str, sku: &str) -> Result> { + let operation_id = OperationId::new(); + let reservation_id = ReservationId::new(); + let reservation = 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), + }) + } +} + +#[derive(Clone)] +struct ReleaseAttempt { + operation_id: OperationId, + event: EventData, +} + +impl ReleaseAttempt { + fn new(sku: &str, reservation_id: ReservationId) -> Result> { + let operation_id = OperationId::new(); + let release = 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), + }) + } +} + +async fn read_inventory( + client: &Client, + stream: &str, +) -> Result<(InventoryState, u64), Box> { + let mut events = client + .read_stream(stream, &ReadStreamOptions::default()) + .await?; + let mut state = InventoryState { + available: 0, + reservations: HashMap::new(), + }; + let mut revision = None; + + while let Some(event) = events.next().await? { + let event = event.get_original_event(); + revision = Some(event.revision); + + match event.event_type.as_str() { + INVENTORY_CREATED_EVENT_TYPE => { + let created = event.as_json::()?; + state.available = created.available; + } + INVENTORY_RESERVED_EVENT_TYPE => { + let reserved = event.as_json::()?; + assert_eq!(event.id, reserved.operation_id.0); + assert!(state.available >= reserved.quantity); + state.available -= reserved.quantity; + assert!( + state + .reservations + .insert(reserved.reservation_id, reserved.quantity) + .is_none() + ); + } + INVENTORY_RELEASED_EVENT_TYPE => { + let released = event.as_json::()?; + assert_eq!(event.id, released.operation_id.0); + let quantity = state + .reservations + .remove(&released.reservation_id) + .expect("the released reservation to be active"); + assert_eq!(quantity, released.quantity); + state.available += released.quantity; + } + event_type => panic!("unexpected inventory event type: {event_type}"), + } + } + + Ok((state, revision.expect("the inventory stream to exist"))) +} + +fn is_revision_conflict(error: &ClientError, expected: u64, current: u64) -> bool { + matches!( + error, + ClientError::WrongExpectedVersion { + expected: StreamState::StreamRevision(actual_expected), + current: CurrentRevision::Current(actual_current), + } if *actual_expected == expected && *actual_current == current + ) +} + +async fn append_at( + client: &Client, + stream: &str, + revision: u64, + event: EventData, +) -> trogon_eventstore::Result { + let options = + AppendToStreamOptions::default().stream_state(StreamState::StreamRevision(revision)); + client.append_to_stream(stream, &options, event).await +} + #[tokio::main] async fn main() -> Result<(), Box> { let connection_string = std::env::var(CONNECTION_STRING_ENV) .unwrap_or_else(|_| DEFAULT_CONNECTION_STRING.to_owned()); - let client = Client::new(connection_string.parse()?)?; + let settings = connection_string.parse::()?; + let first_client = Client::new(settings.clone())?; + let second_client = Client::new(settings)?; - let operation_id = Uuid::new_v4(); let sku = format!("sku-{}", Uuid::new_v4()); let inventory_stream = format!("inventory-{sku}"); let inventory = InventoryCreated { sku: sku.clone(), - available: 100, + available: 1, }; - let created = client + let created = first_client .append_to_stream( inventory_stream.as_str(), &AppendToStreamOptions::default().stream_state(StreamState::NoStream), EventData::json(INVENTORY_CREATED_EVENT_TYPE, &inventory)?.id(Uuid::new_v4()), ) .await?; - let expected_revision = created.next_expected_version; - let reservation = InventoryReserved { - operation_id, - order_id: "order-123".to_owned(), - sku, - quantity: 2, - }; - let event = EventData::json(INVENTORY_RESERVED_EVENT_TYPE, &reservation)?.id(operation_id); + assert_eq!(created.next_expected_version, 0); - // Event IDs are not a stream-wide unique constraint, and `Any` only checks - // recent IDs. A durable retry must retain both this revision and event ID. - let options = AppendToStreamOptions::default() - .stream_state(StreamState::StreamRevision(expected_revision)); + let (first_state, first_revision) = + read_inventory(&first_client, inventory_stream.as_str()).await?; + let (second_state, second_revision) = + read_inventory(&second_client, inventory_stream.as_str()).await?; + assert_eq!(first_state.available, 1); + assert_eq!(second_state.available, 1); + assert_eq!(first_revision, second_revision); - let first = client - .append_to_stream(inventory_stream.as_str(), &options, event.clone()) - .await?; - let retry = client - .append_to_stream(inventory_stream.as_str(), &options, event) - .await?; + let first_attempt = ReservationAttempt::new(FIRST_CLIENT, &sku)?; + let second_attempt = ReservationAttempt::new(SECOND_CLIENT, &sku)?; + // 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(), + ), + append_at( + &second_client, + inventory_stream.as_str(), + second_revision, + second_attempt.event.clone(), + ), + ); - let mut events = client - .read_stream(inventory_stream.as_str(), &ReadStreamOptions::default()) - .await?; - let _created = events.next().await?.expect("the inventory to exist"); - let stored = events.next().await?.expect("the reservation to exist"); - assert!(events.next().await?.is_none()); - assert_eq!(stored.get_original_event().id, operation_id); - assert_eq!(first.next_expected_version, retry.next_expected_version); + let (winner_client, winner, winner_write, loser_client, loser) = + match (first_result, second_result) { + (Ok(write), Err(error)) if is_revision_conflict(&error, 0, 1) => ( + &first_client, + &first_attempt, + write, + &second_client, + &second_attempt, + ), + (Err(error), Ok(write)) if is_revision_conflict(&error, 0, 1) => ( + &second_client, + &second_attempt, + write, + &first_client, + &first_attempt, + ), + results => panic!("expected one reservation winner and one conflict: {results:?}"), + }; + assert_eq!(winner_write.next_expected_version, 1); + + // A conflict invalidates the loser's stale decision, so it must reload before deciding again. + 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()); + + 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 (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_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 (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()); + 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_release = ReleaseAttempt::new(&sku, winner_again.reservation_id)?; + let winner_again_release_write = append_at( + winner_client, + inventory_stream.as_str(), + winner_again_write.next_expected_version, + winner_again_release.event.clone(), + ) + .await?; + assert_eq!(winner_again_release_write.next_expected_version, 6); + + // Durable retries retain the original expected revision and event ID after later writes. + 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, + ) + .await?; + let winner_again_retry = append_at( + winner_client, + inventory_stream.as_str(), + available_again_revision, + winner_again.event, + ) + .await?; + let winner_again_release_retry = append_at( + winner_client, + inventory_stream.as_str(), + winner_again_write.next_expected_version, + winner_again_release.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); + assert_eq!( + winner_again_release_retry.position, + winner_again_release_write.position + ); + + 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()); + + 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 + ); println!( - "operation {operation_id} was stored once at inventory revision {}", - first.next_expected_version + "{} completed a second reserve/release cycle; all six retries remained single events after revision {final_revision}", + winner.client_name ); Ok(()) diff --git a/trogon-eventstore/tests/api/idempotency.rs b/trogon-eventstore/tests/api/idempotency.rs index 87b734a..818cd58 100644 --- a/trogon-eventstore/tests/api/idempotency.rs +++ b/trogon-eventstore/tests/api/idempotency.rs @@ -1,8 +1,88 @@ use crate::common::fresh_stream_id; +use serde::{Deserialize, Serialize}; use serde_json::{Value, json}; -use trogon_eventstore::{AppendToStreamOptions, Client, EventData, StreamState, WriteResult}; +use std::collections::HashMap; +use trogon_eventstore::{ + AppendToStreamOptions, Client, CurrentRevision, Error, EventData, StreamState, WriteResult, +}; use uuid::Uuid; +#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)] +struct OperationId(Uuid); + +impl OperationId { + fn new() -> Self { + Self(Uuid::new_v4()) + } +} + +#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)] +struct ReservationId(Uuid); + +impl ReservationId { + fn new() -> Self { + Self(Uuid::new_v4()) + } +} + +#[derive(Clone)] +struct ReservationAttempt { + operation_id: OperationId, + reservation_id: ReservationId, + event: EventData, +} + +impl ReservationAttempt { + fn new(client: &str) -> Self { + let operation_id = OperationId::new(); + let reservation_id = ReservationId::new(); + + Self { + operation_id, + reservation_id, + event: EventData::json( + "inventory-reserved", + &json!({ + "state": "reserved", + "client": client, + "operation_id": operation_id, + "reservation_id": reservation_id, + "quantity": 1, + }), + ) + .unwrap() + .id(operation_id.0), + } + } +} + +#[derive(Clone)] +struct ReleaseAttempt { + operation_id: OperationId, + event: EventData, +} + +impl ReleaseAttempt { + fn new(reservation_id: ReservationId) -> Self { + let operation_id = OperationId::new(); + + Self { + operation_id, + event: EventData::json( + "inventory-released", + &json!({ + "state": "released", + "operation_id": operation_id, + "reservation_id": reservation_id, + "quantity": 1, + }), + ) + .unwrap() + .id(operation_id.0), + } + } +} + fn event(id: Uuid, state: &str) -> EventData { EventData::json("inventory-reservation", &json!({ "state": state })) .unwrap() @@ -31,6 +111,54 @@ async fn stored_events(client: &Client, stream: &str) -> eyre::Result (u32, HashMap) { + let mut available = 0; + let mut reservations = HashMap::new(); + + for (event_id, _, 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, + "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()); + } + "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; + } + state => panic!("unexpected inventory state: {state}"), + } + } + + (available, reservations) +} + +fn is_revision_conflict(error: &Error, expected: u64, current: u64) -> bool { + matches!( + error, + Error::WrongExpectedVersion { + expected: StreamState::StreamRevision(actual_expected), + current: CurrentRevision::Current(actual_current), + } if *actual_expected == expected && *actual_current == current + ) +} + async fn no_stream_retry_is_idempotent(client: &Client) -> eyre::Result<()> { let stream = fresh_stream_id("idempotency-no-stream"); let event_id = Uuid::new_v4(); @@ -187,12 +315,224 @@ async fn event_ids_are_not_unique_across_streams(client: &Client) -> eyre::Resul Ok(()) } +async fn competing_reservations_are_serialized_and_retryable(client: &Client) -> eyre::Result<()> { + let stream = fresh_stream_id("idempotency-competing-reservations"); + let first_client = Client::new(client.settings().clone())?; + let second_client = Client::new(client.settings().clone())?; + let created_operation_id = Uuid::new_v4(); + + append( + &first_client, + &stream, + StreamState::NoStream, + EventData::json( + "inventory-created", + &json!({ "state": "created", "available": 1 }), + )? + .id(created_operation_id), + ) + .await?; + + 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())); + 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); + assert_ne!( + second_attempt.operation_id.0, + second_attempt.reservation_id.0 + ); + + let (first_result, second_result) = tokio::join!( + append( + &first_client, + &stream, + StreamState::StreamRevision(first_revision), + first_attempt.event.clone(), + ), + append( + &second_client, + &stream, + StreamState::StreamRevision(second_revision), + second_attempt.event.clone(), + ), + ); + + let (winner_client, winner, winner_write, loser_client, loser) = + match (first_result, second_result) { + (Ok(write), Err(error)) if is_revision_conflict(&error, 0, 1) => ( + &first_client, + &first_attempt, + write, + &second_client, + &second_attempt, + ), + (Err(error), Ok(write)) if is_revision_conflict(&error, 0, 1) => ( + &second_client, + &second_attempt, + write, + &first_client, + &first_attempt, + ), + results => eyre::bail!( + "expected one reservation winner and one revision conflict: {results:?}" + ), + }; + + 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_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 released_view = stored_events(loser_client, &stream).await?; + assert_eq!(fold_inventory(&released_view), (1, HashMap::new())); + 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_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 available_again_view = stored_events(winner_client, &stream).await?; + assert_eq!(fold_inventory(&available_again_view), (1, HashMap::new())); + 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_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_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, + ) + .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 + ); + + 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())); + 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, + ] + ); + + Ok(()) +} + pub async fn tests(client: Client) -> eyre::Result<()> { no_stream_retry_is_idempotent(&client).await?; explicit_revision_retry_is_idempotent(&client).await?; retry_identity_does_not_include_payload(&client).await?; same_id_at_a_different_revision_is_a_new_event(&client).await?; event_ids_are_not_unique_across_streams(&client).await?; + competing_reservations_are_serialized_and_retryable(&client).await?; Ok(()) }