Skip to content

Commit f5373ac

Browse files
Merge pull request #125 from code0-tech/#94-add-runtime-usage-service
Add RuntimeUsageService
2 parents 9cff928 + b2731fe commit f5373ac

3 files changed

Lines changed: 68 additions & 4 deletions

File tree

crates/taurus/src/client/mod.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
11
pub mod runtime_status;
2+
pub mod runtime_usage;
Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
1+
use code0_flow::flow_service::retry::create_channel_with_retry;
2+
use tonic::transport::Channel;
3+
use tucana::{
4+
aquila::{RuntimeUsageRequest, runtime_usage_service_client::RuntimeUsageServiceClient},
5+
shared::RuntimeUsage,
6+
};
7+
8+
pub struct TaurusRuntimeUsageService {
9+
channel: Channel,
10+
}
11+
12+
impl TaurusRuntimeUsageService {
13+
pub async fn from_url(aquila_url: String) -> Self {
14+
let channel = create_channel_with_retry("Aquila", aquila_url).await;
15+
TaurusRuntimeUsageService { channel }
16+
}
17+
18+
pub async fn update_runtime_usage(&self, runtime_usage: RuntimeUsage) {
19+
log::info!("Updating the current Runtime Status!");
20+
let mut client = RuntimeUsageServiceClient::new(self.channel.clone());
21+
22+
let request = RuntimeUsageRequest {
23+
runtime_usage: vec![runtime_usage],
24+
};
25+
26+
match client.update(request).await {
27+
Ok(response) => {
28+
log::info!(
29+
"Was the update of the RuntimeStatus accepted by Sagittarius? {}",
30+
response.into_inner().success
31+
);
32+
}
33+
Err(err) => {
34+
log::error!("Failed to update RuntimeStatus: {:?}", err);
35+
}
36+
}
37+
}
38+
}
39+

crates/taurus/src/main.rs

Lines changed: 28 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ mod client;
22
mod config;
33

44
use crate::client::runtime_status::TaurusRuntimeStatusService;
5+
use crate::client::runtime_usage::TaurusRuntimeUsageService;
56
use crate::config::Config;
67
use code0_flow::flow_service::FlowUpdateService;
78

@@ -11,16 +12,20 @@ use futures_lite::StreamExt;
1112
use log::error;
1213
use prost::Message;
1314
use std::collections::HashMap;
15+
use std::time::{Instant, SystemTime, UNIX_EPOCH};
1416
use taurus_core::context::context::Context;
1517
use taurus_core::context::executor::Executor;
1618
use taurus_core::context::registry::FunctionStore;
1719
use taurus_core::context::signal::Signal;
1820
use tokio::signal;
1921
use tonic_health::pb::health_server::HealthServer;
2022
use tucana::shared::value::Kind;
21-
use tucana::shared::{ExecutionFlow, NodeFunction, RuntimeFeature, Translation, Value};
23+
use tucana::shared::{
24+
ExecutionFlow, NodeFunction, RuntimeFeature, RuntimeUsage, Translation, Value,
25+
};
2226

23-
fn handle_message(flow: ExecutionFlow, store: &FunctionStore) -> Signal {
27+
fn handle_message(flow: ExecutionFlow, store: &FunctionStore) -> (Signal, RuntimeUsage) {
28+
let start = Instant::now();
2429
let mut context = Context::default();
2530

2631
let node_functions: HashMap<i64, NodeFunction> = flow
@@ -29,7 +34,17 @@ fn handle_message(flow: ExecutionFlow, store: &FunctionStore) -> Signal {
2934
.map(|node| (node.database_id, node))
3035
.collect();
3136

32-
Executor::new(store, node_functions).execute(flow.starting_node_id, &mut context, true)
37+
let signal =
38+
Executor::new(store, node_functions).execute(flow.starting_node_id, &mut context, true);
39+
let duration_millis = start.elapsed().as_millis() as i64;
40+
41+
(
42+
signal,
43+
RuntimeUsage {
44+
flow_id: flow.flow_id,
45+
duration: duration_millis,
46+
},
47+
)
3348
}
3449

3550
#[tokio::main]
@@ -43,6 +58,7 @@ async fn main() {
4358
let config = Config::new();
4459
let store = FunctionStore::default();
4560
let mut runtime_status_service: Option<TaurusRuntimeStatusService> = None;
61+
let mut runtime_usage_service: Option<TaurusRuntimeUsageService> = None;
4662

4763
let client = match async_nats::connect(config.nats_url.clone()).await {
4864
Ok(client) => {
@@ -92,6 +108,9 @@ async fn main() {
92108
.send()
93109
.await;
94110

111+
let usage_service = TaurusRuntimeUsageService::from_url(config.aquila_url.clone()).await;
112+
runtime_usage_service = Some(usage_service);
113+
95114
let status_service = TaurusRuntimeStatusService::from_url(
96115
config.aquila_url.clone(),
97116
"taurus".into(),
@@ -143,7 +162,8 @@ async fn main() {
143162
};
144163

145164
let flow_id = flow.flow_id;
146-
let value = match handle_message(flow, &store) {
165+
let result = handle_message(flow, &store);
166+
let value = match result.0 {
147167
Signal::Failure(error) => {
148168
log::error!(
149169
"RuntimeError occurred, execution failed because: {:?}",
@@ -180,6 +200,10 @@ async fn main() {
180200
Err(err) => log::error!("Failed to send response: {:?}", err),
181201
}
182202
}
203+
204+
if let Some(usage_service) = &runtime_usage_service {
205+
usage_service.update_runtime_usage(result.1).await;
206+
}
183207
}
184208

185209
log::info!("NATS worker loop ended");

0 commit comments

Comments
 (0)