pub struct SyncUmaDbClient { /* private fields */ }Implementations§
Source§impl SyncUmaDbClient
impl SyncUmaDbClient
pub fn connect( url: String, ca_path: Option<String>, batch_size: Option<u32>, api_key: Option<String>, ) -> DcbResult<Self>
Sourcepub fn subscribe(
&self,
query: Option<DcbQuery>,
after: Option<u64>,
) -> DcbResult<Box<dyn DcbSubscriptionSync + Send + 'static>>
pub fn subscribe( &self, query: Option<DcbQuery>, after: Option<u64>, ) -> DcbResult<Box<dyn DcbSubscriptionSync + Send + 'static>>
Subscribe to events starting from an optional position. This is a convenience wrapper around the async client’s Subscribe RPC. The returned iterator yields events indefinitely until cancelled or the stream ends.
Examples found in repository?
examples/client_sync_insecure.rs (line 93)
7fn main() -> Result<(), Box<dyn std::error::Error>> {
8 // Connect to the gRPC server
9 let url = "http://localhost:50051".to_string();
10 let client = UmaDbClient::new(url).connect()?;
11
12 // Define a consistency boundary
13 let boundary = DcbQuery::new().item(
14 DcbQueryItem::new()
15 .types(["example"])
16 .tags(["tag1", "tag2"]),
17 );
18
19 // Read events for a decision model
20 let mut read_response = client.read(Some(boundary.clone()), None, false, None)?;
21
22 // Build decision model
23 while let Some(result) = read_response.next() {
24 match result {
25 Ok(event) => {
26 println!(
27 "Got event at position {}: {:?}",
28 event.position, event.event
29 );
30 }
31 Err(status) => panic!("gRPC stream error: {}", status),
32 }
33 }
34
35 // Remember the last-known position
36 let last_known_position = read_response.head().unwrap();
37 println!("Last known position is: {:?}", last_known_position);
38
39 // Produce new event, attaching some metadata (e.g. provenance) that is
40 // stored alongside the event and returned when it is read back.
41 let event = DcbEvent::default()
42 .event_type("example")
43 .tags(["tag1", "tag2"])
44 .data(b"Hello, world!")
45 .uuid(Uuid::new_v4())
46 .metadata_entry("source", "client_sync_insecure")
47 .metadata_entry("correlation_id", Uuid::new_v4().to_string());
48
49 // Append event in consistency boundary
50 let append_condition = DcbAppendCondition {
51 fail_if_events_match: boundary.clone(),
52 after: last_known_position,
53 };
54 let position1 = client.append(vec![event.clone()], Some(append_condition.clone()), None)?;
55
56 println!("Appended event at position: {}", position1);
57
58 // Append conflicting event - expect an error
59 let conflicting_event = DcbEvent::default()
60 .event_type("example")
61 .tags(["tag1", "tag2"])
62 .data(b"Hello, world!")
63 .uuid(Uuid::new_v4()); // different UUID
64
65 let conflicting_result = client.append(
66 vec![conflicting_event],
67 Some(append_condition.clone()),
68 None,
69 );
70
71 // Expect an integrity error
72 match conflicting_result {
73 Err(DcbError::IntegrityError(integrity_error)) => {
74 println!("Conflicting event was rejected: {:?}", integrity_error);
75 }
76 other => panic!("Expected IntegrityError, got {:?}", other),
77 }
78
79 // Appending with identical event IDs and append condition is idempotent.
80 println!(
81 "Retrying to append event at position: {:?}",
82 last_known_position
83 );
84 let position2 = client.append(vec![event.clone()], Some(append_condition.clone()), None)?;
85
86 if position1 == position2 {
87 println!("Append method returned same commit position: {}", position2);
88 } else {
89 panic!("Expected idempotent retry!")
90 }
91
92 // Subscribe to all events for a projection
93 let mut subscription = client.subscribe(None, None)?;
94
95 // Build an up-to-date view
96 while let Some(result) = subscription.next() {
97 match result {
98 Ok(ev) => {
99 println!("Processing event at {}: {:?}", ev.position, ev.event);
100 if ev.position == position2 {
101 println!("Projection has processed new event!");
102 break;
103 }
104 }
105 Err(status) => panic!("gRPC stream error: {}", status),
106 }
107 }
108
109 // Track an upstream position
110 let upstream_position = client.get_tracking_info("upstream")?;
111 let next_upstream_position = upstream_position.unwrap_or(0) + 1;
112 println!("Next upstream position: {next_upstream_position}");
113 client.append(
114 vec![],
115 None,
116 Some(TrackingInfo {
117 source: "upstream".to_string(),
118 position: next_upstream_position,
119 }),
120 )?;
121 assert_eq!(
122 next_upstream_position,
123 client.get_tracking_info("upstream")?.unwrap()
124 );
125 println!("Upstream position tracked okay!");
126
127 // Try recording the same upstream position
128 let conflicting_result = client.append(
129 vec![],
130 None,
131 Some(TrackingInfo {
132 source: "upstream".to_string(),
133 position: next_upstream_position,
134 }),
135 );
136
137 // Expect an integrity error
138 match conflicting_result {
139 Err(DcbError::IntegrityError(integrity_error)) => {
140 println!(
141 "Conflicting upstream position was rejected: {:?}",
142 integrity_error
143 );
144 }
145 other => panic!("Expected IntegrityError, got {:?}", other),
146 }
147
148 Ok(())
149}More examples
examples/client_sync_secure.rs (line 116)
7fn main() -> Result<(), Box<dyn std::error::Error>> {
8 // Connect to the gRPC server
9 let url = "https://localhost:50051".to_string();
10 let client = UmaDbClient::new(url)
11 .ca_path("server.pem".to_string()) // For self-signed server certificates.
12 .api_key("umadb:example-api-key-4f7c2b1d9e5f4a038c1a".to_string())
13 .connect()?;
14
15 // Define a consistency boundary
16 let cb = DcbQuery {
17 items: vec![DcbQueryItem {
18 types: vec!["example".to_string()],
19 tags: vec!["tag1".to_string(), "tag2".to_string()],
20 }],
21 };
22
23 // Read events for a decision model
24 let mut read_response = client.read(Some(cb.clone()), None, false, None)?;
25
26 // Build decision model
27 while let Some(result) = read_response.next() {
28 match result {
29 Ok(event) => {
30 println!(
31 "Got event at position {}: {:?}",
32 event.position, event.event
33 );
34 }
35 Err(status) => panic!("gRPC stream error: {}", status),
36 }
37 }
38
39 // Remember the last-known position
40 let last_known_position = read_response.head().unwrap();
41 println!("Last known position is: {:?}", last_known_position);
42
43 // Produce new event, attaching some metadata (e.g. provenance) that is
44 // stored alongside the event and returned when it is read back.
45 let mut metadata = Vec::new();
46 metadata.push(("source".to_string(), "client_sync_secure".to_string()));
47 metadata.push(("correlation_id".to_string(), Uuid::new_v4().to_string()));
48 let event = DcbEvent {
49 event_type: "example".to_string(),
50 tags: vec!["tag1".to_string(), "tag2".to_string()],
51 data: b"Hello, world!".to_vec(),
52 uuid: Some(Uuid::new_v4()),
53 metadata,
54 };
55
56 // Append event in consistency boundary
57 let commit_position1 = client.append(
58 vec![event.clone()],
59 Some(DcbAppendCondition {
60 fail_if_events_match: cb.clone(),
61 after: last_known_position,
62 }),
63 None,
64 )?;
65 println!("Appended event at position: {}", commit_position1);
66
67 // Append conflicting event - expect an error
68 let conflicting_event = DcbEvent {
69 event_type: "example".to_string(),
70 tags: vec!["tag1".to_string(), "tag2".to_string()],
71 data: b"Hello, world!".to_vec(),
72 uuid: Some(Uuid::new_v4()), // different UUID
73 metadata: Vec::new(),
74 };
75 let conflicting_result = client.append(
76 vec![conflicting_event],
77 Some(DcbAppendCondition {
78 fail_if_events_match: cb.clone(),
79 after: last_known_position,
80 }),
81 None,
82 );
83
84 // Expect an integrity error
85 match conflicting_result {
86 Err(DcbError::IntegrityError(integrity_error)) => {
87 println!("Conflicting event was rejected: {:?}", integrity_error);
88 }
89 other => panic!("Expected IntegrityError, got {:?}", other),
90 }
91
92 // Conditional appends with event UUIDs are idempotent.
93 println!(
94 "Retrying to append event at position: {:?}",
95 last_known_position
96 );
97 let commit_position2 = client.append(
98 vec![event.clone()],
99 Some(DcbAppendCondition {
100 fail_if_events_match: cb.clone(),
101 after: last_known_position,
102 }),
103 None,
104 )?;
105
106 if commit_position1 == commit_position2 {
107 println!(
108 "Append method returned same commit position: {}",
109 commit_position2
110 );
111 } else {
112 panic!("Expected idempotent retry!")
113 }
114
115 // Subscribe to all events for a projection
116 let mut subscription = client.subscribe(None, None)?;
117
118 // Build an up-to-date view
119 while let Some(result) = subscription.next() {
120 match result {
121 Ok(ev) => {
122 println!("Processing event at {}: {:?}", ev.position, ev.event);
123 if ev.position == commit_position2 {
124 println!("Projection has processed new event!");
125 break;
126 }
127 }
128 Err(status) => panic!("gRPC stream error: {}", status),
129 }
130 }
131
132 // Track an upstream position
133 let upstream_position = client.get_tracking_info("upstream")?;
134 let next_upstream_position = upstream_position.unwrap_or(0) + 1;
135 println!("Next upstream position: {next_upstream_position}");
136 client.append(
137 vec![],
138 None,
139 Some(TrackingInfo {
140 source: "upstream".to_string(),
141 position: next_upstream_position,
142 }),
143 )?;
144 assert_eq!(
145 next_upstream_position,
146 client.get_tracking_info("upstream")?.unwrap()
147 );
148 println!("Upstream position tracked okay!");
149
150 // Try recording the same upstream position
151 let conflicting_result = client.append(
152 vec![],
153 None,
154 Some(TrackingInfo {
155 source: "upstream".to_string(),
156 position: next_upstream_position,
157 }),
158 );
159
160 // Expect an integrity error
161 match conflicting_result {
162 Err(DcbError::IntegrityError(integrity_error)) => {
163 println!(
164 "Conflicting upstream position was rejected: {:?}",
165 integrity_error
166 );
167 }
168 other => panic!("Expected IntegrityError, got {:?}", other),
169 }
170
171 Ok(())
172}Sourcepub fn close(&self)
pub fn close(&self)
Stops all streaming responses opened by this client.
Delegates to the underlying AsyncUmaDbClient; see its close() for
details. Affects only streams opened by this client instance.
Trait Implementations§
Source§impl DcbEventStoreSync for SyncUmaDbClient
impl DcbEventStoreSync for SyncUmaDbClient
Source§fn read(
&self,
query: Option<DcbQuery>,
start: Option<u64>,
backwards: bool,
limit: Option<u32>,
) -> DcbResult<Box<dyn DcbReadResponseSync + Send + 'static>>
fn read( &self, query: Option<DcbQuery>, start: Option<u64>, backwards: bool, limit: Option<u32>, ) -> DcbResult<Box<dyn DcbReadResponseSync + Send + 'static>>
Reads events from the store based on the provided query and constraints Read more
Source§fn head(&self) -> DcbResult<Option<u64>>
fn head(&self) -> DcbResult<Option<u64>>
Returns the current head position of the event store, or None if empty Read more
Source§fn get_tracking_info(&self, source: &str) -> DcbResult<Option<u64>>
fn get_tracking_info(&self, source: &str) -> DcbResult<Option<u64>>
Returns the greatest recorded upstream position for a tracking source, if any
Auto Trait Implementations§
impl !Freeze for SyncUmaDbClient
impl !RefUnwindSafe for SyncUmaDbClient
impl !UnwindSafe for SyncUmaDbClient
impl Send for SyncUmaDbClient
impl Sync for SyncUmaDbClient
impl Unpin for SyncUmaDbClient
impl UnsafeUnpin for SyncUmaDbClient
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
Wrap the input message
T in a tonic::Request