use umadb_client::UmaDbClient;
use umadb_dcb::{
DcbAppendCondition, DcbError, DcbEvent, DcbEventStoreSync, DcbQuery, DcbQueryItem, TrackingInfo,
};
use uuid::Uuid;
fn main() -> Result<(), Box<dyn std::error::Error>> {
let url = "http://localhost:50051".to_string();
let client = UmaDbClient::new(url).connect()?;
let boundary = DcbQuery::new().item(
DcbQueryItem::new()
.types(["example"])
.tags(["tag1", "tag2"]),
);
let mut read_response = client.read(Some(boundary.clone()), None, false, None, false)?;
while let Some(result) = read_response.next() {
match result {
Ok(event) => {
println!(
"Got event at position {}: {:?}",
event.position, event.event
);
}
Err(status) => panic!("gRPC stream error: {}", status),
}
}
let last_known_position = read_response.head().unwrap();
println!("Last known position is: {:?}", last_known_position);
let event = DcbEvent::default()
.event_type("example")
.tags(["tag1", "tag2"])
.data(b"Hello, world!")
.uuid(Uuid::new_v4());
let append_condition = DcbAppendCondition {
fail_if_events_match: boundary.clone(),
after: last_known_position,
};
let position1 = client.append(vec![event.clone()], Some(append_condition.clone()), None)?;
println!("Appended event at position: {}", position1);
let conflicting_event = DcbEvent::default()
.event_type("example")
.tags(["tag1", "tag2"])
.data(b"Hello, world!")
.uuid(Uuid::new_v4());
let conflicting_result = client.append(
vec![conflicting_event],
Some(append_condition.clone()),
None,
);
match conflicting_result {
Err(DcbError::IntegrityError(integrity_error)) => {
println!("Conflicting event was rejected: {:?}", integrity_error);
}
other => panic!("Expected IntegrityError, got {:?}", other),
}
println!(
"Retrying to append event at position: {:?}",
last_known_position
);
let position2 = client.append(vec![event.clone()], Some(append_condition.clone()), None)?;
if position1 == position2 {
println!("Append method returned same commit position: {}", position2);
} else {
panic!("Expected idempotent retry!")
}
let mut subscription = client.read(None, None, false, None, true)?;
while let Some(result) = subscription.next() {
match result {
Ok(ev) => {
println!("Processing event at {}: {:?}", ev.position, ev.event);
if ev.position == position2 {
println!("Projection has processed new event!");
break;
}
}
Err(status) => panic!("gRPC stream error: {}", status),
}
}
let upstream_position = client.get_tracking_info("upstream")?;
let next_upstream_position = upstream_position.unwrap_or(0) + 1;
println!("Next upstream position: {next_upstream_position}");
client.append(
vec![],
None,
Some(TrackingInfo {
source: "upstream".to_string(),
position: next_upstream_position,
}),
)?;
assert_eq!(
next_upstream_position,
client.get_tracking_info("upstream")?.unwrap()
);
println!("Upstream position tracked okay!");
let conflicting_result = client.append(
vec![],
None,
Some(TrackingInfo {
source: "upstream".to_string(),
position: next_upstream_position,
}),
);
match conflicting_result {
Err(DcbError::IntegrityError(integrity_error)) => {
println!(
"Conflicting upstream position was rejected: {:?}",
integrity_error
);
}
other => panic!("Expected IntegrityError, got {:?}", other),
}
Ok(())
}