Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/workflows/integration.yml
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ jobs:
fail-fast: false
matrix:
test:
- single_node_idempotency
- single_node_streams
- single_node_projections
- single_node_persistent_subscriptions
Expand Down
5 changes: 5 additions & 0 deletions examples/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
81 changes: 81 additions & 0 deletions examples/idempotent_reservation.rs
Original file line number Diff line number Diff line change
@@ -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<dyn Error>> {
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(())
}
198 changes: 198 additions & 0 deletions trogon-eventstore/tests/api/idempotency.rs
Original file line number Diff line number Diff line change
@@ -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<WriteResult> {
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<Vec<(Uuid, u64, Value)>> {
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(())
}
1 change: 1 addition & 0 deletions trogon-eventstore/tests/api/mod.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
pub mod idempotency;
pub mod operations;
pub mod persistent_subscriptions;
pub mod projections;
Expand Down
7 changes: 7 additions & 0 deletions trogon-eventstore/tests/integration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -259,6 +259,7 @@ impl Tests {
}

enum ApiTests {
Idempotency,
Streams,
PersistentSubscriptions,
Projections,
Expand Down Expand Up @@ -353,6 +354,7 @@ async fn run_test(test: impl Into<Tests>, 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
Expand All @@ -376,6 +378,11 @@ async fn run_test(test: impl Into<Tests>, 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
Expand Down
Loading