Skip to content

Commit 87784f2

Browse files
authored
feat(examples): demonstrate idempotent appends (#11)
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent f309717 commit 87784f2

6 files changed

Lines changed: 293 additions & 0 deletions

File tree

‎.github/workflows/integration.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ jobs:
4343
fail-fast: false
4444
matrix:
4545
test:
46+
- single_node_idempotency
4647
- single_node_streams
4748
- single_node_projections
4849
- single_node_persistent_subscriptions

‎examples/Cargo.toml‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,12 +9,17 @@ trogon-eventstore = { path = "../trogon-eventstore" }
99
futures = "0.3"
1010
uuid = { version = "1.1", features = [ "v4", "serde" ] }
1111
serde = "1"
12+
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }
1213

1314
[[example]]
1415
name = "appending_events"
1516
path = "appending_events.rs"
1617
crate-type = ["staticlib"]
1718

19+
[[example]]
20+
name = "idempotent_reservation"
21+
path = "idempotent_reservation.rs"
22+
1823
[[example]]
1924
name = "quickstart"
2025
path = "quickstart.rs"

‎examples/idempotent_reservation.rs‎

Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,81 @@
1+
use serde::{Deserialize, Serialize};
2+
use std::error::Error;
3+
use trogon_eventstore::{AppendToStreamOptions, Client, EventData, ReadStreamOptions, StreamState};
4+
use uuid::Uuid;
5+
6+
const DEFAULT_CONNECTION_STRING: &str = "esdb://localhost:2113?tls=false";
7+
const CONNECTION_STRING_ENV: &str = "TROGON_EVENTSTORE_CONNECTION_STRING";
8+
const INVENTORY_CREATED_EVENT_TYPE: &str = "inventory-created";
9+
const INVENTORY_RESERVED_EVENT_TYPE: &str = "inventory-reserved";
10+
11+
#[derive(Debug, Deserialize, Serialize)]
12+
struct InventoryCreated {
13+
sku: String,
14+
available: u32,
15+
}
16+
17+
#[derive(Debug, Deserialize, Serialize)]
18+
struct InventoryReserved {
19+
operation_id: Uuid,
20+
order_id: String,
21+
sku: String,
22+
quantity: u32,
23+
}
24+
25+
#[tokio::main]
26+
async fn main() -> Result<(), Box<dyn Error>> {
27+
let connection_string = std::env::var(CONNECTION_STRING_ENV)
28+
.unwrap_or_else(|_| DEFAULT_CONNECTION_STRING.to_owned());
29+
let client = Client::new(connection_string.parse()?)?;
30+
31+
let operation_id = Uuid::new_v4();
32+
let sku = format!("sku-{}", Uuid::new_v4());
33+
let inventory_stream = format!("inventory-{sku}");
34+
let inventory = InventoryCreated {
35+
sku: sku.clone(),
36+
available: 100,
37+
};
38+
let created = client
39+
.append_to_stream(
40+
inventory_stream.as_str(),
41+
&AppendToStreamOptions::default().stream_state(StreamState::NoStream),
42+
EventData::json(INVENTORY_CREATED_EVENT_TYPE, &inventory)?.id(Uuid::new_v4()),
43+
)
44+
.await?;
45+
let expected_revision = created.next_expected_version;
46+
let reservation = InventoryReserved {
47+
operation_id,
48+
order_id: "order-123".to_owned(),
49+
sku,
50+
quantity: 2,
51+
};
52+
let event = EventData::json(INVENTORY_RESERVED_EVENT_TYPE, &reservation)?.id(operation_id);
53+
54+
// Event IDs are not a stream-wide unique constraint, and `Any` only checks
55+
// recent IDs. A durable retry must retain both this revision and event ID.
56+
let options = AppendToStreamOptions::default()
57+
.stream_state(StreamState::StreamRevision(expected_revision));
58+
59+
let first = client
60+
.append_to_stream(inventory_stream.as_str(), &options, event.clone())
61+
.await?;
62+
let retry = client
63+
.append_to_stream(inventory_stream.as_str(), &options, event)
64+
.await?;
65+
66+
let mut events = client
67+
.read_stream(inventory_stream.as_str(), &ReadStreamOptions::default())
68+
.await?;
69+
let _created = events.next().await?.expect("the inventory to exist");
70+
let stored = events.next().await?.expect("the reservation to exist");
71+
assert!(events.next().await?.is_none());
72+
assert_eq!(stored.get_original_event().id, operation_id);
73+
assert_eq!(first.next_expected_version, retry.next_expected_version);
74+
75+
println!(
76+
"operation {operation_id} was stored once at inventory revision {}",
77+
first.next_expected_version
78+
);
79+
80+
Ok(())
81+
}
Lines changed: 198 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,198 @@
1+
use crate::common::fresh_stream_id;
2+
use serde_json::{Value, json};
3+
use trogon_eventstore::{AppendToStreamOptions, Client, EventData, StreamState, WriteResult};
4+
use uuid::Uuid;
5+
6+
fn event(id: Uuid, state: &str) -> EventData {
7+
EventData::json("inventory-reservation", &json!({ "state": state }))
8+
.unwrap()
9+
.id(id)
10+
}
11+
12+
async fn append(
13+
client: &Client,
14+
stream: &str,
15+
state: StreamState,
16+
event: EventData,
17+
) -> trogon_eventstore::Result<WriteResult> {
18+
let options = AppendToStreamOptions::default().stream_state(state);
19+
client.append_to_stream(stream, &options, event).await
20+
}
21+
22+
async fn stored_events(client: &Client, stream: &str) -> eyre::Result<Vec<(Uuid, u64, Value)>> {
23+
let mut read = client.read_stream(stream, &Default::default()).await?;
24+
let mut events = Vec::new();
25+
26+
while let Some(event) = read.next().await? {
27+
let event = event.get_original_event();
28+
events.push((event.id, event.revision, event.as_json()?));
29+
}
30+
31+
Ok(events)
32+
}
33+
34+
async fn no_stream_retry_is_idempotent(client: &Client) -> eyre::Result<()> {
35+
let stream = fresh_stream_id("idempotency-no-stream");
36+
let event_id = Uuid::new_v4();
37+
38+
let first = append(
39+
client,
40+
&stream,
41+
StreamState::NoStream,
42+
event(event_id, "reserved"),
43+
)
44+
.await?;
45+
let retry = append(
46+
client,
47+
&stream,
48+
StreamState::NoStream,
49+
event(event_id, "reserved"),
50+
)
51+
.await?;
52+
53+
assert_eq!(first.next_expected_version, 0);
54+
assert_eq!(retry.next_expected_version, 0);
55+
assert_eq!(first.position, retry.position);
56+
assert_eq!(stored_events(client, &stream).await?.len(), 1);
57+
58+
Ok(())
59+
}
60+
61+
async fn explicit_revision_retry_is_idempotent(client: &Client) -> eyre::Result<()> {
62+
let stream = fresh_stream_id("idempotency-revision");
63+
let event_id = Uuid::new_v4();
64+
65+
append(
66+
client,
67+
&stream,
68+
StreamState::NoStream,
69+
event(Uuid::new_v4(), "opened"),
70+
)
71+
.await?;
72+
let first = append(
73+
client,
74+
&stream,
75+
StreamState::StreamRevision(0),
76+
event(event_id, "reserved"),
77+
)
78+
.await?;
79+
let retry = append(
80+
client,
81+
&stream,
82+
StreamState::StreamRevision(0),
83+
event(event_id, "reserved"),
84+
)
85+
.await?;
86+
87+
assert_eq!(first.next_expected_version, 1);
88+
assert_eq!(retry.next_expected_version, 1);
89+
assert_eq!(first.position, retry.position);
90+
assert_eq!(stored_events(client, &stream).await?.len(), 2);
91+
92+
Ok(())
93+
}
94+
95+
async fn retry_identity_does_not_include_payload(client: &Client) -> eyre::Result<()> {
96+
let stream = fresh_stream_id("idempotency-payload");
97+
let seed_id = Uuid::new_v4();
98+
let event_id = Uuid::new_v4();
99+
100+
append(
101+
client,
102+
&stream,
103+
StreamState::NoStream,
104+
event(seed_id, "opened"),
105+
)
106+
.await?;
107+
append(
108+
client,
109+
&stream,
110+
StreamState::StreamRevision(0),
111+
event(event_id, "reserved"),
112+
)
113+
.await?;
114+
append(
115+
client,
116+
&stream,
117+
StreamState::StreamRevision(0),
118+
event(event_id, "released"),
119+
)
120+
.await?;
121+
122+
assert_eq!(
123+
stored_events(client, &stream).await?,
124+
[
125+
(seed_id, 0, json!({ "state": "opened" })),
126+
(event_id, 1, json!({ "state": "reserved" })),
127+
]
128+
);
129+
130+
Ok(())
131+
}
132+
133+
async fn same_id_at_a_different_revision_is_a_new_event(client: &Client) -> eyre::Result<()> {
134+
let stream = fresh_stream_id("idempotency-different-revision");
135+
let event_id = Uuid::new_v4();
136+
137+
append(
138+
client,
139+
&stream,
140+
StreamState::NoStream,
141+
event(event_id, "reserved"),
142+
)
143+
.await?;
144+
let second = append(
145+
client,
146+
&stream,
147+
StreamState::StreamRevision(0),
148+
event(event_id, "released"),
149+
)
150+
.await?;
151+
152+
assert_eq!(second.next_expected_version, 1);
153+
assert_eq!(
154+
stored_events(client, &stream).await?,
155+
[
156+
(event_id, 0, json!({ "state": "reserved" })),
157+
(event_id, 1, json!({ "state": "released" })),
158+
]
159+
);
160+
161+
Ok(())
162+
}
163+
164+
async fn event_ids_are_not_unique_across_streams(client: &Client) -> eyre::Result<()> {
165+
let first_stream = fresh_stream_id("idempotency-first-stream");
166+
let second_stream = fresh_stream_id("idempotency-second-stream");
167+
let event_id = Uuid::new_v4();
168+
169+
append(
170+
client,
171+
&first_stream,
172+
StreamState::NoStream,
173+
event(event_id, "reserved"),
174+
)
175+
.await?;
176+
append(
177+
client,
178+
&second_stream,
179+
StreamState::NoStream,
180+
event(event_id, "reserved"),
181+
)
182+
.await?;
183+
184+
assert_eq!(stored_events(client, &first_stream).await?.len(), 1);
185+
assert_eq!(stored_events(client, &second_stream).await?.len(), 1);
186+
187+
Ok(())
188+
}
189+
190+
pub async fn tests(client: Client) -> eyre::Result<()> {
191+
no_stream_retry_is_idempotent(&client).await?;
192+
explicit_revision_retry_is_idempotent(&client).await?;
193+
retry_identity_does_not_include_payload(&client).await?;
194+
same_id_at_a_different_revision_is_a_new_event(&client).await?;
195+
event_ids_are_not_unique_across_streams(&client).await?;
196+
197+
Ok(())
198+
}

‎trogon-eventstore/tests/api/mod.rs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
pub mod idempotency;
12
pub mod operations;
23
pub mod persistent_subscriptions;
34
pub mod projections;

‎trogon-eventstore/tests/integration.rs‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -259,6 +259,7 @@ impl Tests {
259259
}
260260

261261
enum ApiTests {
262+
Idempotency,
262263
Streams,
263264
PersistentSubscriptions,
264265
Projections,
@@ -353,6 +354,7 @@ async fn run_test(test: impl Into<Tests>, topology: Topologies) -> eyre::Result<
353354

354355
let result = match test {
355356
Tests::Api(test) => match test {
357+
ApiTests::Idempotency => api::idempotency::tests(predifined_client).await,
356358
ApiTests::Streams => api::streams::tests(predifined_client).await,
357359
ApiTests::PersistentSubscriptions => {
358360
api::persistent_subscriptions::tests(predifined_client).await
@@ -376,6 +378,11 @@ async fn run_test(test: impl Into<Tests>, topology: Topologies) -> eyre::Result<
376378
Ok(())
377379
}
378380

381+
#[tokio::test(flavor = "multi_thread")]
382+
async fn single_node_idempotency() -> eyre::Result<()> {
383+
run_test(ApiTests::Idempotency, Topologies::SingleNode).await
384+
}
385+
379386
#[tokio::test(flavor = "multi_thread")]
380387
async fn single_node_streams() -> eyre::Result<()> {
381388
run_test(ApiTests::Streams, Topologies::SingleNode).await

0 commit comments

Comments
 (0)