Skip to content

Commit 6e71cfd

Browse files
Merge pull request #239 from code0-tech/deps/code0-flow
deps/code0-flow
2 parents 4b4b030 + 10c0897 commit 6e71cfd

5 files changed

Lines changed: 43 additions & 131 deletions

File tree

Cargo.lock

Lines changed: 12 additions & 88 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,8 @@ version = "0.0.0"
77
edition = "2024"
88

99
[workspace.dependencies]
10-
code0-flow = { version = "0.0.33" }
11-
tucana = { version = "0.0.73", features = ["aquila"] }
10+
code0-flow = { version = "0.0.36" }
11+
tucana = { version = "0.0.74", features = ["aquila"] }
1212
serde_json = { version = "1.0.138" }
1313
log = "0.4.27"
1414
env_logger = "0.11.8"

adapter/rest/src/main.rs

Lines changed: 13 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ use base::{
33
store::{FlowExecutionResult, FlowIdentifyResult},
44
traits::Server as ServerTrait,
55
};
6+
use code0_flow::flow_service::ModuleDefinitionAppendix;
67
use http_body_util::{BodyExt, Full};
78
use hyper::server::conn::http1;
89
use hyper::{Request, Response};
@@ -17,7 +18,7 @@ use std::sync::Arc;
1718
use tokio::net::TcpListener;
1819
use tonic::async_trait;
1920
use tucana::shared::{
20-
AdapterStatusConfiguration, Struct, ValidationFlow, Value, helper::value::ToValue, value::Kind,
21+
Endpoint, ModuleDefinition, Struct, ValidationFlow, Value, helper::value::ToValue, value::Kind,
2122
};
2223

2324
use crate::response::{error_to_http_response, value_to_http_response};
@@ -42,14 +43,18 @@ async fn main() {
4243
let addr = runner.get_server_config().port;
4344
let host = runner.get_server_config().host.clone();
4445

45-
let configs = vec![AdapterStatusConfiguration {
46-
flow_type_identifiers: vec![String::from("REST")],
47-
data: Some(
48-
tucana::shared::adapter_status_configuration::Data::Endpoint(format!(
49-
r"{}:{}/${{project_slug}}/${{flow_setting_identifier}}",
50-
host, addr
46+
let configs = vec![ModuleDefinitionAppendix {
47+
module_identifier: String::from("draco-rest"),
48+
definitions: vec![ModuleDefinition {
49+
flow_type_identifier: vec![String::from("REST")],
50+
value: Some(tucana::shared::module_definition::Value::Endpoint(
51+
Endpoint {
52+
host,
53+
port: addr as i64,
54+
endpoint: String::from(r"/${{project_slug}}${{httpURL}}"),
55+
},
5156
)),
52-
),
57+
}],
5358
}];
5459
match runner.serve(configs).await {
5560
Ok(_) => (),

crates/base/src/client/mod.rs

Lines changed: 7 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -9,15 +9,13 @@ use tonic::{
99
use tucana::{
1010
aquila::{
1111
RuntimeStatusUpdateRequest, runtime_status_service_client::RuntimeStatusServiceClient,
12-
runtime_status_update_request::Status,
1312
},
14-
shared::{AdapterRuntimeStatus, AdapterStatusConfiguration},
13+
shared::ModuleStatus,
1514
};
1615

1716
pub struct DracoRuntimeStatusService {
1817
channel: Channel,
1918
identifier: String,
20-
configs: Vec<AdapterStatusConfiguration>,
2119
aquila_token: String,
2220
}
2321

@@ -69,33 +67,22 @@ pub async fn create_channel_with_retry(channel_name: &str, url: String) -> Chann
6967
}
7068
}
7169
impl DracoRuntimeStatusService {
72-
pub async fn from_url(
73-
aquila_url: String,
74-
aquila_token: String,
75-
identifier: String,
76-
configs: Vec<AdapterStatusConfiguration>,
77-
) -> Self {
70+
pub async fn from_url(aquila_url: String, aquila_token: String, identifier: String) -> Self {
7871
let channel = create_channel_with_retry("Aquila", aquila_url).await;
79-
Self::new(channel, identifier, configs, aquila_token)
72+
Self::new(channel, identifier, aquila_token)
8073
}
8174

82-
pub fn new(
83-
channel: Channel,
84-
identifier: String,
85-
configs: Vec<AdapterStatusConfiguration>,
86-
aquila_token: String,
87-
) -> Self {
75+
pub fn new(channel: Channel, identifier: String, aquila_token: String) -> Self {
8876
DracoRuntimeStatusService {
8977
channel,
9078
identifier,
91-
configs,
9279
aquila_token,
9380
}
9481
}
9582

9683
pub async fn update_runtime_status_by_status(
9784
&self,
98-
status: tucana::shared::adapter_runtime_status::Status,
85+
status: tucana::shared::module_status::StatusVariant,
9986
) {
10087
log::info!("Updating the current runtime status!");
10188
let mut client = RuntimeStatusServiceClient::new(self.channel.clone());
@@ -113,12 +100,11 @@ impl DracoRuntimeStatusService {
113100
get_authorization_metadata(&self.aquila_token),
114101
Extensions::new(),
115102
RuntimeStatusUpdateRequest {
116-
status: Some(Status::AdapterRuntimeStatus(AdapterRuntimeStatus {
103+
status: Some(ModuleStatus {
117104
status: status.into(),
118105
timestamp: timestamp as i64,
119106
identifier: self.identifier.clone(),
120-
configurations: self.configs.clone(),
121-
})),
107+
}),
122108
},
123109
);
124110

crates/base/src/runner.rs

Lines changed: 9 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,11 @@ use crate::{
44
store::AdapterStore,
55
traits::{LoadConfig, Server as AdapterServer},
66
};
7-
use code0_flow::flow_service::FlowUpdateService;
7+
use code0_flow::flow_service::{FlowUpdateService, ModuleDefinitionAppendix};
88
use std::{sync::Arc, time::Duration};
99
use tokio::{signal, task::JoinHandle, time::sleep};
1010
use tonic::transport::Server;
1111
use tonic_health::pb::health_server::HealthServer;
12-
use tucana::shared::AdapterStatusConfiguration;
1312

1413
/// Context passed to adapter server implementations containing all shared resources
1514
pub struct ServerContext<C: LoadConfig> {
@@ -56,10 +55,7 @@ impl<C: LoadConfig> ServerRunner<C> {
5655
})
5756
}
5857

59-
pub async fn serve(
60-
self,
61-
runtime_config: Vec<AdapterStatusConfiguration>,
62-
) -> anyhow::Result<()> {
58+
pub async fn serve(self, appendix: Vec<ModuleDefinitionAppendix>) -> anyhow::Result<()> {
6359
let config = self.context.adapter_config.clone();
6460
let mut runtime_status_service: Option<Arc<DracoRuntimeStatusService>> = None;
6561
let mut runtime_status_heartbeat_task: Option<JoinHandle<()>> = None;
@@ -71,26 +67,27 @@ impl<C: LoadConfig> ServerRunner<C> {
7167
config.aquila_url.clone(),
7268
config.aquila_token.clone(),
7369
config.draco_variant.clone(),
74-
runtime_config,
7570
)
7671
.await,
7772
));
7873

7974
if let Some(ser) = &runtime_status_service {
8075
ser.update_runtime_status_by_status(
81-
tucana::shared::adapter_runtime_status::Status::NotReady,
76+
tucana::shared::module_status::StatusVariant::NotReady,
8277
)
8378
.await;
8479
};
8580

8681
let service_name = format!("draco-{}", config.draco_variant.to_lowercase());
82+
8783
let mut definition_service = FlowUpdateService::from_url(
8884
config.aquila_url.clone(),
8985
config.definition_path.as_str(),
9086
config.aquila_token.clone(),
9187
)
9288
.await
93-
.with_definition_source(service_name);
89+
.with_definition_source(service_name)
90+
.with_appendix(appendix);
9491

9592
let mut success = false;
9693
let mut count = 1;
@@ -144,7 +141,7 @@ impl<C: LoadConfig> ServerRunner<C> {
144141

145142
if let Some(ser) = &runtime_status_service {
146143
ser.update_runtime_status_by_status(
147-
tucana::shared::adapter_runtime_status::Status::Running,
144+
tucana::shared::module_status::StatusVariant::Running,
148145
)
149146
.await;
150147

@@ -163,7 +160,7 @@ impl<C: LoadConfig> ServerRunner<C> {
163160
interval.tick().await;
164161
status_service
165162
.update_runtime_status_by_status(
166-
tucana::shared::adapter_runtime_status::Status::Running,
163+
tucana::shared::module_status::StatusVariant::Running,
167164
)
168165
.await;
169166
}
@@ -251,7 +248,7 @@ impl<C: LoadConfig> ServerRunner<C> {
251248

252249
if let Some(ser) = &runtime_status_service {
253250
ser.update_runtime_status_by_status(
254-
tucana::shared::adapter_runtime_status::Status::Stopped,
251+
tucana::shared::module_status::StatusVariant::Stopped,
255252
)
256253
.await;
257254
};

0 commit comments

Comments
 (0)