diff --git a/.github/workflows/integration.yml b/.github/workflows/integration.yml index 0237755..ae4fac2 100644 --- a/.github/workflows/integration.yml +++ b/.github/workflows/integration.yml @@ -43,6 +43,7 @@ jobs: fail-fast: false matrix: test: + - single_node_idempotency - single_node_streams - single_node_projections - single_node_persistent_subscriptions diff --git a/examples/Cargo.toml b/examples/Cargo.toml index 56788bd..5f41c09 100644 --- a/examples/Cargo.toml +++ b/examples/Cargo.toml @@ -9,12 +9,17 @@ trogon-eventstore = { path = "../trogon-eventstore" } futures = "0.3" uuid = { version = "1.1", features = [ "v4", "serde" ] } serde = "1" +tokio = { version = "1", features = ["macros", "rt-multi-thread"] } [[example]] name = "appending_events" path = "appending_events.rs" crate-type = ["staticlib"] +[[example]] +name = "idempotent_reservation" +path = "idempotent_reservation.rs" + [[example]] name = "quickstart" path = "quickstart.rs" diff --git a/examples/idempotent_reservation.rs b/examples/idempotent_reservation.rs new file mode 100644 index 0000000..153449d --- /dev/null +++ b/examples/idempotent_reservation.rs @@ -0,0 +1,81 @@ +use serde::{Deserialize, Serialize}; +use std::error::Error; +use trogon_eventstore::{AppendToStreamOptions, Client, EventData, ReadStreamOptions, StreamState}; +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"; + +#[derive(Debug, Deserialize, Serialize)] +struct InventoryCreated { + sku: String, + available: u32, +} + +#[derive(Debug, Deserialize, Serialize)] +struct InventoryReserved { + operation_id: Uuid, + order_id: String, + sku: String, + quantity: u32, +} + +#[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 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, + }; + let created = 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); + + // 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 = 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 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); + + println!( + "operation {operation_id} was stored once at inventory revision {}", + first.next_expected_version + ); + + Ok(()) +} diff --git a/trogon-eventstore/tests/api/idempotency.rs b/trogon-eventstore/tests/api/idempotency.rs new file mode 100644 index 0000000..87b734a --- /dev/null +++ b/trogon-eventstore/tests/api/idempotency.rs @@ -0,0 +1,198 @@ +use crate::common::fresh_stream_id; +use serde_json::{Value, json}; +use trogon_eventstore::{AppendToStreamOptions, Client, EventData, StreamState, WriteResult}; +use uuid::Uuid; + +fn event(id: Uuid, state: &str) -> EventData { + EventData::json("inventory-reservation", &json!({ "state": state })) + .unwrap() + .id(id) +} + +async fn append( + client: &Client, + stream: &str, + state: StreamState, + event: EventData, +) -> trogon_eventstore::Result { + let options = AppendToStreamOptions::default().stream_state(state); + client.append_to_stream(stream, &options, event).await +} + +async fn stored_events(client: &Client, stream: &str) -> eyre::Result> { + let mut read = client.read_stream(stream, &Default::default()).await?; + let mut events = Vec::new(); + + while let Some(event) = read.next().await? { + let event = event.get_original_event(); + events.push((event.id, event.revision, event.as_json()?)); + } + + Ok(events) +} + +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(); + + let first = append( + client, + &stream, + StreamState::NoStream, + event(event_id, "reserved"), + ) + .await?; + let retry = append( + client, + &stream, + StreamState::NoStream, + event(event_id, "reserved"), + ) + .await?; + + assert_eq!(first.next_expected_version, 0); + assert_eq!(retry.next_expected_version, 0); + assert_eq!(first.position, retry.position); + assert_eq!(stored_events(client, &stream).await?.len(), 1); + + Ok(()) +} + +async fn explicit_revision_retry_is_idempotent(client: &Client) -> eyre::Result<()> { + let stream = fresh_stream_id("idempotency-revision"); + let event_id = Uuid::new_v4(); + + append( + client, + &stream, + StreamState::NoStream, + event(Uuid::new_v4(), "opened"), + ) + .await?; + let first = append( + client, + &stream, + StreamState::StreamRevision(0), + event(event_id, "reserved"), + ) + .await?; + let retry = append( + client, + &stream, + StreamState::StreamRevision(0), + event(event_id, "reserved"), + ) + .await?; + + assert_eq!(first.next_expected_version, 1); + assert_eq!(retry.next_expected_version, 1); + assert_eq!(first.position, retry.position); + assert_eq!(stored_events(client, &stream).await?.len(), 2); + + Ok(()) +} + +async fn retry_identity_does_not_include_payload(client: &Client) -> eyre::Result<()> { + let stream = fresh_stream_id("idempotency-payload"); + let seed_id = Uuid::new_v4(); + let event_id = Uuid::new_v4(); + + append( + client, + &stream, + StreamState::NoStream, + event(seed_id, "opened"), + ) + .await?; + append( + client, + &stream, + StreamState::StreamRevision(0), + event(event_id, "reserved"), + ) + .await?; + append( + client, + &stream, + StreamState::StreamRevision(0), + event(event_id, "released"), + ) + .await?; + + assert_eq!( + stored_events(client, &stream).await?, + [ + (seed_id, 0, json!({ "state": "opened" })), + (event_id, 1, json!({ "state": "reserved" })), + ] + ); + + Ok(()) +} + +async fn same_id_at_a_different_revision_is_a_new_event(client: &Client) -> eyre::Result<()> { + let stream = fresh_stream_id("idempotency-different-revision"); + let event_id = Uuid::new_v4(); + + append( + client, + &stream, + StreamState::NoStream, + event(event_id, "reserved"), + ) + .await?; + let second = append( + client, + &stream, + StreamState::StreamRevision(0), + event(event_id, "released"), + ) + .await?; + + assert_eq!(second.next_expected_version, 1); + assert_eq!( + stored_events(client, &stream).await?, + [ + (event_id, 0, json!({ "state": "reserved" })), + (event_id, 1, json!({ "state": "released" })), + ] + ); + + Ok(()) +} + +async fn event_ids_are_not_unique_across_streams(client: &Client) -> eyre::Result<()> { + let first_stream = fresh_stream_id("idempotency-first-stream"); + let second_stream = fresh_stream_id("idempotency-second-stream"); + let event_id = Uuid::new_v4(); + + append( + client, + &first_stream, + StreamState::NoStream, + event(event_id, "reserved"), + ) + .await?; + append( + client, + &second_stream, + StreamState::NoStream, + event(event_id, "reserved"), + ) + .await?; + + assert_eq!(stored_events(client, &first_stream).await?.len(), 1); + assert_eq!(stored_events(client, &second_stream).await?.len(), 1); + + 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?; + + Ok(()) +} diff --git a/trogon-eventstore/tests/api/mod.rs b/trogon-eventstore/tests/api/mod.rs index b59d60c..6f0545e 100644 --- a/trogon-eventstore/tests/api/mod.rs +++ b/trogon-eventstore/tests/api/mod.rs @@ -1,3 +1,4 @@ +pub mod idempotency; pub mod operations; pub mod persistent_subscriptions; pub mod projections; diff --git a/trogon-eventstore/tests/integration.rs b/trogon-eventstore/tests/integration.rs index c990cf8..94603be 100644 --- a/trogon-eventstore/tests/integration.rs +++ b/trogon-eventstore/tests/integration.rs @@ -259,6 +259,7 @@ impl Tests { } enum ApiTests { + Idempotency, Streams, PersistentSubscriptions, Projections, @@ -353,6 +354,7 @@ async fn run_test(test: impl Into, topology: Topologies) -> eyre::Result< let result = match test { Tests::Api(test) => match test { + ApiTests::Idempotency => api::idempotency::tests(predifined_client).await, ApiTests::Streams => api::streams::tests(predifined_client).await, ApiTests::PersistentSubscriptions => { api::persistent_subscriptions::tests(predifined_client).await @@ -376,6 +378,11 @@ async fn run_test(test: impl Into, topology: Topologies) -> eyre::Result< Ok(()) } +#[tokio::test(flavor = "multi_thread")] +async fn single_node_idempotency() -> eyre::Result<()> { + run_test(ApiTests::Idempotency, Topologies::SingleNode).await +} + #[tokio::test(flavor = "multi_thread")] async fn single_node_streams() -> eyre::Result<()> { run_test(ApiTests::Streams, Topologies::SingleNode).await