Skip to content

Commit 4d791be

Browse files
committed
fix(client): align with the server contract
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent bab6f1d commit 4d791be

8 files changed

Lines changed: 89 additions & 181 deletions

File tree

‎trogon-eventstore/protos/streams.proto‎

Lines changed: 10 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -105,34 +105,20 @@ message ReadResp {
105105

106106
// The $all or stream subscription has caught up and become live.
107107
message CaughtUp {
108-
// Current time in the server when the subscription caught up
109-
google.protobuf.Timestamp timestamp = 1;
110-
111-
// Checkpoint for resuming a stream subscription.
112-
// For stream subscriptions it is populated unless the stream is empty.
113-
// For $all subscriptions it is not populated.
114-
optional int64 stream_revision = 2;
115-
116-
// Checkpoint for resuming a $all subscription.
117-
// For stream subscriptions it is not populated.
118-
// For $all subscriptions it is populated unless the database is empty.
119-
optional Position position = 3;
108+
oneof position {
109+
uint64 stream_position = 1;
110+
event_store.client.AllStreamPosition all_stream_position = 2;
111+
event_store.client.Empty no_position = 3;
112+
}
120113
}
121114

122115
// The $all or stream subscription has fallen back into catchup mode and is no longer live.
123116
message FellBehind {
124-
// Current time in the server when the subscription fell behind
125-
google.protobuf.Timestamp timestamp = 1;
126-
127-
// Checkpoint for resuming a stream subscription.
128-
// For stream subscriptions it is populated unless the stream is empty.
129-
// For $all subscriptions it is not populated.
130-
optional int64 stream_revision = 2;
131-
132-
// Checkpoint for resuming a $all subscription.
133-
// For stream subscriptions it is not populated.
134-
// For $all subscriptions it is populated unless the database is empty.
135-
optional Position position = 3;
117+
oneof position {
118+
uint64 stream_position = 1;
119+
event_store.client.AllStreamPosition all_stream_position = 2;
120+
event_store.client.Empty no_position = 3;
121+
}
136122
}
137123

138124
message ReadEvent {
@@ -165,11 +151,6 @@ message ReadResp {
165151
google.protobuf.Timestamp timestamp = 3;
166152
}
167153

168-
message Position {
169-
uint64 commit_position = 1;
170-
uint64 prepare_position = 2;
171-
}
172-
173154
message StreamNotFound {
174155
event_store.client.StreamIdentifier stream_identifier = 1;
175156
}

‎trogon-eventstore/src/commands.rs‎

Lines changed: 35 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,8 @@
11
#![allow(clippy::large_enum_variant)]
22
//! Commands this client supports.
33
use std::ops::Add;
4-
use std::time::{Duration, SystemTime, UNIX_EPOCH};
4+
use std::time::{Duration, SystemTime};
55

6-
use chrono::{DateTime, Utc};
76
use futures::TryStreamExt;
87
use nom::AsBytes;
98
use tokio::sync::mpsc;
@@ -840,26 +839,46 @@ impl Subscription {
840839
}
841840

842841
streams::read_resp::Content::CaughtUp(args) => {
843-
let args = args.timestamp.map(|t| crate::CaughtUp {
844-
date: timestamp_to_datetime(t),
845-
stream_revision: args.stream_revision.map(|x| x as u64),
846-
position: args.position.map(|x| Position {
847-
commit: x.commit_position,
848-
prepare: x.prepare_position,
849-
}),
842+
use streams::read_resp::caught_up::Position::*;
843+
let args = args.position.map(|position| match position {
844+
StreamPosition(revision) => crate::CaughtUp {
845+
stream_revision: Some(revision),
846+
position: None,
847+
},
848+
AllStreamPosition(position) => crate::CaughtUp {
849+
stream_revision: None,
850+
position: Some(Position {
851+
commit: position.commit_position,
852+
prepare: position.prepare_position,
853+
}),
854+
},
855+
NoPosition(()) => crate::CaughtUp {
856+
stream_revision: None,
857+
position: None,
858+
},
850859
});
851860

852861
return Ok(SubscriptionEvent::CaughtUp(args));
853862
}
854863

855864
streams::read_resp::Content::FellBehind(args) => {
856-
let args = args.timestamp.map(|t| crate::FellBehind {
857-
date: timestamp_to_datetime(t),
858-
stream_revision: args.stream_revision.map(|x| x as u64),
859-
position: args.position.map(|x| Position {
860-
commit: x.commit_position,
861-
prepare: x.prepare_position,
862-
}),
865+
use streams::read_resp::fell_behind::Position::*;
866+
let args = args.position.map(|position| match position {
867+
StreamPosition(revision) => crate::FellBehind {
868+
stream_revision: Some(revision),
869+
position: None,
870+
},
871+
AllStreamPosition(position) => crate::FellBehind {
872+
stream_revision: None,
873+
position: Some(Position {
874+
commit: position.commit_position,
875+
prepare: position.prepare_position,
876+
}),
877+
},
878+
NoPosition(()) => crate::FellBehind {
879+
stream_revision: None,
880+
position: None,
881+
},
863882
});
864883

865884
return Ok(SubscriptionEvent::FellBehind(args));
@@ -927,10 +946,6 @@ impl Subscription {
927946
}
928947
}
929948

930-
fn timestamp_to_datetime(t: prost_types::Timestamp) -> DateTime<Utc> {
931-
(UNIX_EPOCH + Duration::new(t.seconds as u64, t.nanos as u32)).into()
932-
}
933-
934949
/// Runs the subscription command.
935950
pub fn subscribe_to_stream(
936951
connection: GrpcClient,

‎trogon-eventstore/src/event_store/generated/streams.rs‎

Lines changed: 28 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -194,36 +194,38 @@ pub mod read_resp {
194194
/// The $all or stream subscription has caught up and become live.
195195
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
196196
pub struct CaughtUp {
197-
/// Current time in the server when the subscription caught up
198-
#[prost(message, optional, tag = "1")]
199-
pub timestamp: ::core::option::Option<::prost_types::Timestamp>,
200-
/// Checkpoint for resuming a stream subscription.
201-
/// For stream subscriptions it is populated unless the stream is empty.
202-
/// For $all subscriptions it is not populated.
203-
#[prost(int64, optional, tag = "2")]
204-
pub stream_revision: ::core::option::Option<i64>,
205-
/// Checkpoint for resuming a $all subscription.
206-
/// For stream subscriptions it is not populated.
207-
/// For $all subscriptions it is populated unless the database is empty.
208-
#[prost(message, optional, tag = "3")]
209-
pub position: ::core::option::Option<Position>,
197+
#[prost(oneof = "caught_up::Position", tags = "1, 2, 3")]
198+
pub position: ::core::option::Option<caught_up::Position>,
199+
}
200+
/// Nested message and enum types in `CaughtUp`.
201+
pub mod caught_up {
202+
#[derive(Clone, Copy, PartialEq, ::prost::Oneof)]
203+
pub enum Position {
204+
#[prost(uint64, tag = "1")]
205+
StreamPosition(u64),
206+
#[prost(message, tag = "2")]
207+
AllStreamPosition(crate::event_store::generated::common::AllStreamPosition),
208+
#[prost(message, tag = "3")]
209+
NoPosition(()),
210+
}
210211
}
211212
/// The $all or stream subscription has fallen back into catchup mode and is no longer live.
212213
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
213214
pub struct FellBehind {
214-
/// Current time in the server when the subscription fell behind
215-
#[prost(message, optional, tag = "1")]
216-
pub timestamp: ::core::option::Option<::prost_types::Timestamp>,
217-
/// Checkpoint for resuming a stream subscription.
218-
/// For stream subscriptions it is populated unless the stream is empty.
219-
/// For $all subscriptions it is not populated.
220-
#[prost(int64, optional, tag = "2")]
221-
pub stream_revision: ::core::option::Option<i64>,
222-
/// Checkpoint for resuming a $all subscription.
223-
/// For stream subscriptions it is not populated.
224-
/// For $all subscriptions it is populated unless the database is empty.
225-
#[prost(message, optional, tag = "3")]
226-
pub position: ::core::option::Option<Position>,
215+
#[prost(oneof = "fell_behind::Position", tags = "1, 2, 3")]
216+
pub position: ::core::option::Option<fell_behind::Position>,
217+
}
218+
/// Nested message and enum types in `FellBehind`.
219+
pub mod fell_behind {
220+
#[derive(Clone, Copy, PartialEq, ::prost::Oneof)]
221+
pub enum Position {
222+
#[prost(uint64, tag = "1")]
223+
StreamPosition(u64),
224+
#[prost(message, tag = "2")]
225+
AllStreamPosition(crate::event_store::generated::common::AllStreamPosition),
226+
#[prost(message, tag = "3")]
227+
NoPosition(()),
228+
}
227229
}
228230
#[derive(Clone, PartialEq, ::prost::Message)]
229231
pub struct ReadEvent {
@@ -283,13 +285,6 @@ pub mod read_resp {
283285
#[prost(message, optional, tag = "3")]
284286
pub timestamp: ::core::option::Option<::prost_types::Timestamp>,
285287
}
286-
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
287-
pub struct Position {
288-
#[prost(uint64, tag = "1")]
289-
pub commit_position: u64,
290-
#[prost(uint64, tag = "2")]
291-
pub prepare_position: u64,
292-
}
293288
#[derive(Clone, PartialEq, ::prost::Message)]
294289
pub struct StreamNotFound {
295290
#[prost(message, optional, tag = "1")]

‎trogon-eventstore/src/operations/gossip.rs‎

Lines changed: 1 addition & 78 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,8 @@
1+
use crate::ClientSettings;
12
use crate::event_store::client::gossip as wire;
23
use crate::grpc::HyperClient;
3-
use crate::http::http_configure_auth;
44
use crate::request::build_request_metadata;
55
use crate::types::Endpoint;
6-
use crate::{ClientSettings, grpc};
76
use serde::{Deserialize, Serialize};
87
use tonic::{Request, Status};
98
use uuid::Uuid;
@@ -71,57 +70,6 @@ pub async fn read(
7170
Ok(members)
7271
}
7372

74-
pub(crate) async fn http_read(
75-
setts: &ClientSettings,
76-
handle: grpc::Handle,
77-
) -> Result<Vec<MemberInfo>, Box<dyn std::error::Error>> {
78-
let client = reqwest::Client::builder()
79-
.danger_accept_invalid_certs(!setts.tls_verify_cert)
80-
.build()?;
81-
82-
let default_auth = setts
83-
.default_user_name
84-
.as_ref()
85-
.map(|c| crate::Authentication::Basic(c.clone()));
86-
87-
let resp = http_configure_auth(
88-
client.get(format!("{}/gossip", handle.url())),
89-
default_auth.as_ref(),
90-
)
91-
.send()
92-
.await?;
93-
94-
let gossip = resp.json::<Gossip>().await?;
95-
96-
Ok(gossip
97-
.members
98-
.into_iter()
99-
.map(|i| MemberInfo {
100-
instance_id: i.instance_id,
101-
time_stamp: i.time_stamp.timestamp(),
102-
state: i.state,
103-
is_alive: i.is_alive,
104-
http_end_point: Endpoint {
105-
host: i.external_http_ip,
106-
port: i.external_http_port as u32,
107-
},
108-
last_commit_position: i.last_commit_position,
109-
writer_checkpoint: i.writer_checkpoint,
110-
chaser_checkpoint: i.chaser_checkpoint,
111-
epoch_position: i.epoch_position,
112-
epoch_number: i.epoch_number,
113-
epoch_id: i.epoch_id,
114-
node_priority: i.node_priority,
115-
})
116-
.collect())
117-
}
118-
119-
#[derive(Serialize, Deserialize, Debug)]
120-
#[serde(rename_all = "camelCase")]
121-
struct Gossip {
122-
members: Vec<HttpMemberInfo>,
123-
}
124-
12573
#[derive(Debug, Clone)]
12674
pub struct MemberInfo {
12775
pub instance_id: Uuid,
@@ -138,31 +86,6 @@ pub struct MemberInfo {
13886
pub node_priority: i64,
13987
}
14088

141-
#[derive(Deserialize, Serialize, Debug)]
142-
#[serde(rename_all = "camelCase")]
143-
pub struct HttpMemberInfo {
144-
pub instance_id: Uuid,
145-
pub time_stamp: chrono::DateTime<chrono::Utc>,
146-
pub state: VNodeState,
147-
pub is_alive: bool,
148-
pub internal_tcp_ip: String,
149-
pub internal_tcp_port: u16,
150-
pub internal_secure_tcp_port: u16,
151-
pub external_tcp_ip: String,
152-
pub external_secure_tcp_port: u16,
153-
#[serde(rename = "httpEndPointIp")]
154-
pub external_http_ip: String,
155-
#[serde(rename = "httpEndPointPort")]
156-
pub external_http_port: u16,
157-
pub last_commit_position: i64,
158-
pub writer_checkpoint: i64,
159-
pub chaser_checkpoint: i64,
160-
pub epoch_position: i64,
161-
pub epoch_number: i64,
162-
pub epoch_id: Uuid,
163-
pub node_priority: i64,
164-
}
165-
16689
#[derive(Copy, Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
16790
#[serde(rename_all = "PascalCase")]
16891
pub enum VNodeState {

‎trogon-eventstore/src/operations/mod.rs‎

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -73,12 +73,9 @@ impl Client {
7373
pub async fn read_gossip(&self) -> crate::Result<Vec<gossip::MemberInfo>> {
7474
let handle = self.inner.current_selected_node().await?;
7575

76-
// We currently use the http endpoint instead of the gRPC one because at that time
77-
// 04-25-2022, the public gRPC endpoint doesn't return all the gossip info like current
78-
// epoch and other checkpoints.
79-
gossip::http_read(self.inner.connection_settings(), handle)
76+
gossip::read(self.inner.connection_settings(), &handle.client, handle.uri)
8077
.await
81-
.map_err(|e| crate::Error::IllegalStateError(e.to_string()))
78+
.map_err(crate::Error::from_grpc)
8279
}
8380

8481
pub async fn stats(&self, options: &StatsOptions) -> crate::Result<Stats> {

‎trogon-eventstore/src/types.rs‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1127,14 +1127,12 @@ pub enum SubscriptionEvent {
11271127

11281128
#[derive(Debug)]
11291129
pub struct CaughtUp {
1130-
pub date: DateTime<Utc>,
11311130
pub stream_revision: Option<u64>,
11321131
pub position: Option<Position>,
11331132
}
11341133

11351134
#[derive(Debug)]
11361135
pub struct FellBehind {
1137-
pub date: DateTime<Utc>,
11381136
pub stream_revision: Option<u64>,
11391137
pub position: Option<Position>,
11401138
}

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

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
use std::time::Duration;
22
use tracing::debug;
3+
use trogon_eventstore::Credentials;
34
use trogon_eventstore::operations;
45
use trogon_eventstore::operations::StatsOptions;
56

@@ -186,12 +187,15 @@ async fn test_change_user_password(
186187
)
187188
.await?;
188189

190+
let options = operations::OperationalOptions::default()
191+
.authenticated(Credentials::new(login.clone(), password.clone()));
192+
189193
client
190194
.change_user_password(
191195
login.as_str(),
192-
password,
196+
password.as_str(),
193197
names.next().unwrap(),
194-
&Default::default(),
198+
&options,
195199
)
196200
.await?;
197201

@@ -242,10 +246,7 @@ async fn test_op_restart_persistent_subscription_subsystem(
242246
}
243247

244248
async fn test_scavenge(client: &operations::Client) -> trogon_eventstore::Result<()> {
245-
let result = client.start_scavenge(1, 0, &Default::default()).await?;
246-
let result = client.stop_scavenge(result.id(), &Default::default()).await;
247-
248-
assert!(result.is_ok());
249+
client.start_scavenge(1, 0, &Default::default()).await?;
249250

250251
Ok(())
251252
}

0 commit comments

Comments
 (0)