Skip to content

Commit fa09650

Browse files
authored
fix(client): align with the server contract
Signed-off-by: Straw Hat Team Bot <61149376+sht-bot@users.noreply.github.com>
1 parent 41b32bc commit fa09650

10 files changed

Lines changed: 41 additions & 138 deletions

File tree

‎.github/workflows/integration.yml‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,22 +7,22 @@ on:
77
description: Server image registry
88
required: true
99
type: string
10-
default: docker.io
10+
default: ghcr.io
1111
server_repository:
1212
description: Server image repository
1313
required: true
1414
type: string
15-
default: eventstore
15+
default: trogonstack
1616
server_container:
1717
description: Server image name
1818
required: true
1919
type: string
20-
default: eventstore
20+
default: trogoneventstore
2121
server_version:
2222
description: Server image tag
2323
required: true
2424
type: string
25-
default: latest
25+
default: ci
2626

2727
permissions:
2828
contents: read

‎docker-compose.yml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ services:
2424
- volumes-provisioner
2525

2626
esdb-node1: &template
27-
image: ${ESDB_DOCKER_REGISTRY:-docker.io}/${ESDB_DOCKER_REPO:-eventstore}/${ESDB_DOCKER_CONTAINER:-eventstore}:${ESDB_DOCKER_CONTAINER_VERSION:-latest}
27+
image: ${ESDB_DOCKER_REGISTRY:-ghcr.io}/${ESDB_DOCKER_REPO:-trogonstack}/${ESDB_DOCKER_CONTAINER:-trogoneventstore}:${ESDB_DOCKER_CONTAINER_VERSION:-ci}
2828
env_file:
2929
- vars.env
3030
environment:

‎trogon-eventstore/src/commands.rs‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -841,7 +841,7 @@ impl Subscription {
841841

842842
streams::read_resp::Content::CaughtUp(args) => {
843843
let args = args.timestamp.map(|t| crate::CaughtUp {
844-
date: timestamp_to_datetime(t),
844+
timestamp: timestamp_to_datetime(t),
845845
stream_revision: args.stream_revision.map(|x| x as u64),
846846
position: args.position.map(|x| Position {
847847
commit: x.commit_position,
@@ -854,7 +854,7 @@ impl Subscription {
854854

855855
streams::read_resp::Content::FellBehind(args) => {
856856
let args = args.timestamp.map(|t| crate::FellBehind {
857-
date: timestamp_to_datetime(t),
857+
timestamp: timestamp_to_datetime(t),
858858
stream_revision: args.stream_revision.map(|x| x as u64),
859859
position: args.position.map(|x| Position {
860860
commit: x.commit_position,

‎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: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1127,14 +1127,14 @@ pub enum SubscriptionEvent {
11271127

11281128
#[derive(Debug)]
11291129
pub struct CaughtUp {
1130-
pub date: DateTime<Utc>,
1130+
pub timestamp: DateTime<Utc>,
11311131
pub stream_revision: Option<u64>,
11321132
pub position: Option<Position>,
11331133
}
11341134

11351135
#[derive(Debug)]
11361136
pub struct FellBehind {
1137-
pub date: DateTime<Utc>,
1137+
pub timestamp: DateTime<Utc>,
11381138
pub stream_revision: Option<u64>,
11391139
pub position: Option<Position>,
11401140
}

‎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
}

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

Lines changed: 16 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,5 @@
11
use crate::common::{fresh_stream_id, generate_events};
22
use chrono::{Datelike, Utc};
3-
use futures::channel::oneshot;
43
use std::collections::HashMap;
54
use std::time::Duration;
65
use tracing::{debug, warn};
@@ -265,9 +264,7 @@ async fn test_subscription(client: &Client) -> eyre::Result<()> {
265264
.subscribe_to_stream(stream_id.as_str(), &options)
266265
.await;
267266

268-
let (tx, recv) = oneshot::channel();
269-
270-
tokio::spawn(async move {
267+
let subscription = tokio::spawn(async move {
271268
let mut count = 0usize;
272269
let max = 6usize;
273270

@@ -280,20 +277,20 @@ async fn test_subscription(client: &Client) -> eyre::Result<()> {
280277
}
281278
}
282279

283-
tx.send(count).unwrap();
284-
Ok(()) as trogon_eventstore::Result<()>
280+
Ok(count) as trogon_eventstore::Result<usize>
285281
});
286282

287283
let _ = client
288284
.append_to_stream(stream_id, &Default::default(), events_after)
289285
.await?;
290286

291-
match tokio::time::timeout(Duration::from_secs(60), recv).await {
287+
match tokio::time::timeout(Duration::from_secs(60), subscription).await {
292288
Ok(test_count) => {
289+
let test_count = test_count??;
293290
assert_eq!(
294-
test_count?, 6,
291+
test_count, 6,
295292
"We are testing proper state after catchup subscription: got {} expected {}.",
296-
test_count?, 6
293+
test_count, 6
297294
);
298295
}
299296

@@ -328,25 +325,20 @@ async fn test_subscription_caughtup(client: &Client) -> trogon_eventstore::Resul
328325
.subscribe_to_stream(stream_id.clone(), &options)
329326
.await;
330327

331-
let (tx, recv) = oneshot::channel();
332-
333-
tokio::spawn(async move {
328+
let caught_up = tokio::time::timeout(Duration::from_secs(60), async move {
334329
loop {
335-
if let SubscriptionEvent::CaughtUp(_) = sub.next_subscription_event().await? {
336-
break;
330+
if let SubscriptionEvent::CaughtUp(caught_up) = sub.next_subscription_event().await? {
331+
return Ok::<_, trogon_eventstore::Error>(caught_up);
337332
}
338333
}
334+
})
335+
.await
336+
.expect("test_subscription_caughtup timed out")?
337+
.expect("server did not provide caught-up context");
339338

340-
let _ = tx.send(());
341-
Ok(()) as trogon_eventstore::Result<()>
342-
});
343-
344-
if tokio::time::timeout(Duration::from_secs(60), recv)
345-
.await
346-
.is_err()
347-
{
348-
panic!("test_subscription_caughtup timed out!");
349-
}
339+
assert!(caught_up.timestamp <= Utc::now());
340+
assert_eq!(caught_up.stream_revision, Some(9));
341+
assert_eq!(caught_up.position, None);
350342

351343
Ok(())
352344
}

‎trogon-eventstore/tests/images.rs‎

Lines changed: 5 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -6,10 +6,10 @@ use testcontainers::{
66
core::{ContainerPort, Mount, WaitFor},
77
};
88

9-
const DEFAULT_REGISTRY: &str = "docker.io";
10-
const DEFAULT_REPO: &str = "eventstore";
11-
const DEFAULT_CONTAINER: &str = "eventstore";
12-
const DEFAULT_TAG: &str = "latest";
9+
const DEFAULT_REGISTRY: &str = "ghcr.io";
10+
const DEFAULT_REPO: &str = "trogonstack";
11+
const DEFAULT_CONTAINER: &str = "trogoneventstore";
12+
const DEFAULT_TAG: &str = "ci";
1313

1414
#[derive(Debug, Clone)]
1515
pub struct EventStoreDB {
@@ -23,10 +23,6 @@ impl EventStoreDB {
2323
pub fn insecure_mode(mut self) -> Self {
2424
self.env_vars
2525
.insert("EVENTSTORE_INSECURE".to_string(), "true".to_string());
26-
self.env_vars.insert(
27-
"EVENTSTORE_ENABLE_ATOM_PUB_OVER_HTTP".to_string(),
28-
"true".to_string(),
29-
);
3026

3127
self
3228
}
@@ -163,16 +159,10 @@ impl Default for EventStoreDB {
163159
let tag = option_env!("ESDB_DOCKER_CONTAINER_VERSION").unwrap_or(DEFAULT_TAG);
164160
let repo = option_env!("ESDB_DOCKER_REPO").unwrap_or(DEFAULT_REPO);
165161
let container = option_env!("ESDB_DOCKER_CONTAINER").unwrap_or(DEFAULT_CONTAINER);
166-
let mut env_vars = HashMap::new();
167-
168-
env_vars.insert(
169-
"EVENTSTORE_GOSSIP_ON_SINGLE_NODE".to_string(),
170-
"true".to_string(),
171-
);
172162
EventStoreDB {
173163
name: format!("{}/{}/{}", registry, repo, container),
174164
tag: tag.to_string(),
175-
env_vars,
165+
env_vars: HashMap::new(),
176166
mounts: vec![],
177167
}
178168
}

‎trogon-eventstore/tests/integration.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ async fn wait_node_is_alive(
5757
match tokio::time::timeout(
5858
std::time::Duration::from_secs(1),
5959
client
60-
.get(format!("{}://localhost:{}/health/live", protocol, port))
60+
.get(format!("{}://localhost:{}/-/readiness", protocol, port))
6161
.send(),
6262
)
6363
.await

0 commit comments

Comments
 (0)