Skip to content

Commit 6b8e9ca

Browse files
committed
fix: made config same to Aquila and Taurus
1 parent b75ade3 commit 6b8e9ca

6 files changed

Lines changed: 40 additions & 32 deletions

File tree

Cargo.lock

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

adapter/rest/src/main.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
use base::{
22
extract_flow_setting_field,
33
runner::{ServerContext, ServerRunner},
4-
store::FlowIdenfiyResult,
4+
store::FlowIdentifyResult,
55
traits::{IdentifiableFlow, LoadConfig, Server as ServerTrait},
66
};
77
use code0_flow::flow_config::env_with_default;
@@ -64,7 +64,7 @@ impl ServerTrait<HttpServerConfig> for HttpServer {
6464
};
6565

6666
match store.get_possible_flow_match(pattern, route).await {
67-
FlowIdenfiyResult::Single(flow) => {
67+
FlowIdentifyResult::Single(flow) => {
6868
execute_flow(flow, request, store).await
6969
}
7070
_ => Some(HttpResponse::internal_server_error(

crates/base/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,3 +15,4 @@ tonic-health = { workspace = true }
1515
uuid = { workspace = true }
1616
prost = { workspace = true }
1717
futures-lite = { workspace = true }
18+
log = { workspace = true }

crates/base/src/config.rs

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,11 @@ pub struct AdapterConfig {
3232
/// Port on which the adapter's Health Service server will listen.
3333
pub grpc_port: u16,
3434

35+
/// GRPC Host
36+
///
37+
/// Host on which the adapter's Health Service server will listen.
38+
pub grpc_host: String,
39+
3540
/// Aquila URL
3641
///
3742
/// URL of the Aquila server to connect to.
@@ -45,7 +50,7 @@ pub struct AdapterConfig {
4550
/// Is Monitored
4651
///
4752
/// If true the Adapter will expose a grpc health service server.
48-
pub is_monitored: bool,
53+
pub with_health_service: bool,
4954
}
5055

5156
impl AdapterConfig {
@@ -57,6 +62,7 @@ impl AdapterConfig {
5762
let nats_bucket =
5863
code0_flow::flow_config::env_with_default("NATS_BUCKET", String::from("flow_store"));
5964
let grpc_port = code0_flow::flow_config::env_with_default("GRPC_PORT", 50051);
65+
let grpc_host = code0_flow::flow_config::env_with_default("GRPC_HOST", String::from("localhost"));
6066
let aquila_url = code0_flow::flow_config::env_with_default(
6167
"AQUILA_URL",
6268
String::from("grpc://localhost:50051"),
@@ -69,17 +75,18 @@ impl AdapterConfig {
6975
"DEFINITION_PATH",
7076
String::from("./definition.yaml"),
7177
);
72-
let is_monitored = code0_flow::flow_config::env_with_default("IS_MONITORED", false);
78+
let with_health_service = code0_flow::flow_config::env_with_default("WITH_HEALTH_SERVICE", false);
7379

7480
Self {
7581
environment,
7682
nats_bucket,
7783
mode,
7884
nats_url,
7985
grpc_port,
86+
grpc_host,
8087
aquila_url,
8188
definition_path,
82-
is_monitored,
89+
with_health_service,
8390
}
8491
}
8592

crates/base/src/runner.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -61,10 +61,10 @@ impl<C: LoadConfig> ServerRunner<C> {
6161
definition_service.send().await;
6262
}
6363

64-
if config.is_monitored {
64+
if config.with_health_service {
6565
let health_service =
6666
code0_flow::flow_health::HealthService::new(config.nats_url.clone());
67-
let address = format!("127.0.0.1:{}", config.grpc_port).parse()?;
67+
let address = format!("{}:{}", config.grpc_host, config.grpc_port).parse()?;
6868

6969
tokio::spawn(async move {
7070
let _ = Server::builder()
@@ -73,7 +73,7 @@ impl<C: LoadConfig> ServerRunner<C> {
7373
.await;
7474
});
7575

76-
println!("Health server started at 127.0.0.1:{}", config.grpc_port);
76+
log::info!("Health server started at {}:{}", config.grpc_host, config.grpc_port);
7777
}
7878

7979
self.server.init(&self.context).await?;

crates/base/src/store.rs

Lines changed: 23 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ pub struct AdapterStore {
1010
kv: async_nats::jetstream::kv::Store,
1111
}
1212

13-
pub enum FlowIdenfiyResult {
13+
pub enum FlowIdentifyResult {
1414
None,
1515
Single(ValidationFlow),
1616
Multiple(Vec<ValidationFlow>),
@@ -20,19 +20,24 @@ impl AdapterStore {
2020
pub async fn from_url(url: String, bucket: String) -> Self {
2121
let client = match async_nats::connect(url).await {
2222
Ok(client) => client,
23-
Err(err) => panic!("Failed to connect to NATS server: {}", err),
23+
Err(err) => panic!("Failed to connect to NATS server: {:?}", err),
2424
};
2525

26-
let jetstream = async_nats::jetstream::new(client.clone());
26+
let stream = async_nats::jetstream::new(client.clone());
2727

28-
let _ = jetstream
28+
match stream
2929
.create_key_value(Config {
3030
bucket: bucket.clone(),
3131
..Default::default()
3232
})
33-
.await;
33+
.await {
34+
Ok(_) => {
35+
log::info!("Successfully created NATS bucket");
36+
},
37+
Err(err) => panic!("Failed to create NATS bucket: {:?}", err),
38+
}
3439

35-
let kv = match jetstream.get_key_value(bucket).await {
40+
let kv = match stream.get_key_value(bucket).await {
3641
Ok(kv) => kv,
3742
Err(err) => panic!("Failed to get key-value store: {}", err),
3843
};
@@ -63,21 +68,19 @@ impl AdapterStore {
6368
&self,
6469
pattern: String,
6570
id: I,
66-
) -> FlowIdenfiyResult {
71+
) -> FlowIdentifyResult {
6772
let mut collector = Vec::new();
6873
let mut keys = match self.kv.keys().await {
6974
Ok(keys) => keys.boxed(),
7075
Err(err) => {
71-
eprintln!("Failed to get keys: {}", err);
72-
return FlowIdenfiyResult::None;
76+
log::error!("Failed to get keys: {}", err);
77+
return FlowIdentifyResult::None;
7378
}
7479
};
7580

7681
while let Ok(Some(key)) = keys.try_next().await {
77-
println!("comparing: key: {} pattern {:?}", key, pattern);
7882

7983
if !Self::is_matching_key(&pattern, &key) {
80-
println!("Key does not match pattern: {}", key);
8184
continue;
8285
}
8386

@@ -92,9 +95,9 @@ impl AdapterStore {
9295
}
9396

9497
match collector.len() {
95-
0 => FlowIdenfiyResult::None,
96-
1 => FlowIdenfiyResult::Single(collector[0].clone()),
97-
_ => FlowIdenfiyResult::Multiple(collector),
98+
0 => FlowIdentifyResult::None,
99+
1 => FlowIdentifyResult::Single(collector[0].clone()),
100+
_ => FlowIdentifyResult::Multiple(collector),
98101
}
99102
}
100103

@@ -130,16 +133,15 @@ impl AdapterStore {
130133
match result {
131134
Ok(message) => match Value::decode(message.payload) {
132135
Ok(value) => {
133-
println!("Response: {:?}", &value);
134136
Some(value)
135137
}
136138
Err(err) => {
137-
eprintln!("Failed to decode response from NATS server: {}", err);
138-
return None;
139+
log::error!("Failed to decode response from NATS server: {:?}", err);
140+
None
139141
}
140142
},
141143
Err(err) => {
142-
eprintln!("Failed to send request to NATS server: {}", err);
144+
log::error!("Failed to send request to NATS server: {:?}", err);
143145
None
144146
}
145147
}
@@ -154,22 +156,19 @@ impl AdapterStore {
154156
}
155157

156158
fn is_matching_key(pattern: &String, key: &String) -> bool {
157-
let splitted_pattern = pattern.split(".");
158-
let splitted_key = key.split(".").collect::<Vec<&str>>();
159-
160-
let zip = splitted_pattern.into_iter().zip(splitted_key);
159+
let split_pattern = pattern.split(".");
160+
let split_key = key.split(".").collect::<Vec<&str>>();
161+
let zip = split_pattern.into_iter().zip(split_key);
161162

162163
for (pattern_part, key_part) in zip {
163164
if pattern_part == "*" {
164165
continue;
165166
}
166167

167168
if pattern_part != key_part {
168-
println!("matching: pattern: {} key: {}", pattern_part, key_part);
169169
return false;
170170
}
171171
}
172-
println!("pattern was correct");
173172
true
174173
}
175174
}

0 commit comments

Comments
 (0)