Skip to content

Commit b5ba173

Browse files
committed
ref: split taurus main file into app module
1 parent a929f9e commit b5ba173

3 files changed

Lines changed: 381 additions & 334 deletions

File tree

crates/taurus/src/app/mod.rs

Lines changed: 218 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,218 @@
1+
mod worker;
2+
3+
use std::time::Duration;
4+
5+
use code0_flow::flow_config::load_env_file;
6+
use code0_flow::flow_config::mode::Mode::DYNAMIC;
7+
use code0_flow::flow_service::FlowUpdateService;
8+
use taurus_core::runtime::engine::ExecutionEngine;
9+
use taurus_provider::providers::remote::nats_remote_runtime::NATSRemoteRuntime;
10+
use tokio::signal;
11+
use tokio::task::JoinHandle;
12+
use tokio::time::sleep;
13+
use tonic_health::pb::health_server::HealthServer;
14+
use tucana::shared::{RuntimeFeature, Translation};
15+
16+
use crate::client::runtime_status::TaurusRuntimeStatusService;
17+
use crate::client::runtime_usage::TaurusRuntimeUsageService;
18+
use crate::config::Config;
19+
20+
pub async fn run() {
21+
init_logging();
22+
load_env_file();
23+
24+
let config = Config::new();
25+
let engine = ExecutionEngine::new();
26+
let client = connect_nats(&config).await;
27+
28+
let mut health_task = spawn_health_task(&config);
29+
let (runtime_status_service, runtime_usage_service) =
30+
setup_dynamic_services_if_needed(&config).await;
31+
32+
let nats_remote = NATSRemoteRuntime::new(client.clone());
33+
let mut worker_task = worker::spawn_worker(client, engine, nats_remote, runtime_usage_service);
34+
35+
wait_for_shutdown(&mut worker_task, &mut health_task).await;
36+
update_stopped_status(runtime_status_service.as_ref()).await;
37+
38+
log::info!("Taurus shutdown complete");
39+
}
40+
41+
fn init_logging() {
42+
env_logger::Builder::from_default_env()
43+
.filter_level(log::LevelFilter::Debug)
44+
.init();
45+
}
46+
47+
async fn connect_nats(config: &Config) -> async_nats::Client {
48+
match async_nats::connect(config.nats_url.clone()).await {
49+
Ok(client) => {
50+
log::info!("Connected to NATS server");
51+
client
52+
}
53+
Err(err) => {
54+
panic!("Failed to connect to NATS server: {}", err);
55+
}
56+
}
57+
}
58+
59+
fn spawn_health_task(config: &Config) -> Option<JoinHandle<()>> {
60+
if !config.with_health_service {
61+
return None;
62+
}
63+
64+
let health_service = code0_flow::flow_health::HealthService::new(config.nats_url.clone());
65+
let address = match format!("{}:{}", config.grpc_host, config.grpc_port).parse() {
66+
Ok(address) => address,
67+
Err(err) => {
68+
log::error!("Failed to parse gRPC address: {:?}", err);
69+
return None;
70+
}
71+
};
72+
73+
log::info!("Health server starting at {}", address);
74+
Some(tokio::spawn(async move {
75+
if let Err(err) = tonic::transport::Server::builder()
76+
.add_service(HealthServer::new(health_service))
77+
.serve(address)
78+
.await
79+
{
80+
log::error!("Health server error: {:?}", err);
81+
} else {
82+
log::info!("Health server stopped gracefully");
83+
}
84+
}))
85+
}
86+
87+
async fn setup_dynamic_services_if_needed(
88+
config: &Config,
89+
) -> (
90+
Option<TaurusRuntimeStatusService>,
91+
Option<TaurusRuntimeUsageService>,
92+
) {
93+
if config.mode != DYNAMIC {
94+
return (None, None);
95+
}
96+
97+
push_definitions_until_success(config).await;
98+
99+
let runtime_usage_service = Some(
100+
TaurusRuntimeUsageService::from_url(config.aquila_url.clone(), config.aquila_token.clone())
101+
.await,
102+
);
103+
104+
let runtime_status_service = Some(
105+
TaurusRuntimeStatusService::from_url(
106+
config.aquila_url.clone(),
107+
config.aquila_token.clone(),
108+
"taurus".into(),
109+
runtime_features(),
110+
)
111+
.await,
112+
);
113+
114+
if let Some(status_service) = runtime_status_service.as_ref() {
115+
status_service
116+
.update_runtime_status(tucana::shared::execution_runtime_status::Status::Running)
117+
.await;
118+
}
119+
120+
(runtime_status_service, runtime_usage_service)
121+
}
122+
123+
async fn push_definitions_until_success(config: &Config) {
124+
let definition_service = FlowUpdateService::from_url(
125+
config.aquila_url.clone(),
126+
config.definitions.as_str(),
127+
config.aquila_token.clone(),
128+
)
129+
.await;
130+
131+
let mut retry_count = 1;
132+
loop {
133+
if definition_service.send_with_status().await {
134+
break;
135+
}
136+
137+
log::warn!(
138+
"Updating definitions failed, trying again in 3 seconds (retry #{})",
139+
retry_count
140+
);
141+
retry_count += 1;
142+
sleep(Duration::from_secs(3)).await;
143+
}
144+
}
145+
146+
fn runtime_features() -> Vec<RuntimeFeature> {
147+
vec![RuntimeFeature {
148+
name: vec![Translation {
149+
code: "en-US".to_string(),
150+
content: "Runtime".to_string(),
151+
}],
152+
description: vec![Translation {
153+
code: "en-US".to_string(),
154+
content: "Will execute incoming flows.".to_string(),
155+
}],
156+
}]
157+
}
158+
159+
async fn update_stopped_status(runtime_status_service: Option<&TaurusRuntimeStatusService>) {
160+
if let Some(status_service) = runtime_status_service {
161+
status_service
162+
.update_runtime_status(tucana::shared::execution_runtime_status::Status::Stopped)
163+
.await;
164+
}
165+
}
166+
167+
async fn wait_for_shutdown(
168+
worker_task: &mut JoinHandle<()>,
169+
health_task: &mut Option<JoinHandle<()>>,
170+
) {
171+
#[cfg(unix)]
172+
let sigterm = async {
173+
use tokio::signal::unix::{SignalKind, signal};
174+
175+
let mut term = signal(SignalKind::terminate()).expect("failed to install SIGTERM handler");
176+
term.recv().await;
177+
};
178+
179+
#[cfg(not(unix))]
180+
let sigterm = std::future::pending::<()>();
181+
182+
if let Some(health_task) = health_task.as_mut() {
183+
tokio::select! {
184+
_ = &mut *worker_task => {
185+
log::warn!("NATS worker task finished, shutting down");
186+
health_task.abort();
187+
}
188+
_ = &mut *health_task => {
189+
log::warn!("Health server task finished, shutting down");
190+
worker_task.abort();
191+
}
192+
_ = signal::ctrl_c() => {
193+
log::info!("Ctrl+C/Exit signal received, shutting down");
194+
worker_task.abort();
195+
health_task.abort();
196+
}
197+
_ = sigterm => {
198+
log::info!("SIGTERM received, shutting down");
199+
worker_task.abort();
200+
health_task.abort();
201+
}
202+
}
203+
} else {
204+
tokio::select! {
205+
_ = &mut *worker_task => {
206+
log::warn!("NATS worker task finished, shutting down");
207+
}
208+
_ = signal::ctrl_c() => {
209+
log::info!("Ctrl+C/Exit signal received, shutting down");
210+
worker_task.abort();
211+
}
212+
_ = sigterm => {
213+
log::info!("SIGTERM received, shutting down");
214+
worker_task.abort();
215+
}
216+
}
217+
}
218+
}

crates/taurus/src/app/worker.rs

Lines changed: 161 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,161 @@
1+
use std::sync::Arc;
2+
use std::sync::atomic::{AtomicUsize, Ordering};
3+
use std::time::Instant;
4+
5+
use futures_lite::StreamExt;
6+
use prost::Message;
7+
use taurus_core::runtime::engine::{EmitType, ExecutionEngine, RespondEmitter};
8+
use taurus_core::types::signal::Signal;
9+
use taurus_provider::providers::emitter::nats_emitter::NATSRespondEmitter;
10+
use taurus_provider::providers::remote::nats_remote_runtime::NATSRemoteRuntime;
11+
use tokio::task::JoinHandle;
12+
use tucana::shared::value::Kind;
13+
use tucana::shared::{ExecutionFlow, RuntimeUsage, Value};
14+
15+
use crate::client::runtime_usage::TaurusRuntimeUsageService;
16+
17+
pub fn spawn_worker(
18+
client: async_nats::Client,
19+
engine: ExecutionEngine,
20+
nats_remote: NATSRemoteRuntime,
21+
runtime_usage_service: Option<TaurusRuntimeUsageService>,
22+
) -> JoinHandle<()> {
23+
tokio::spawn(async move {
24+
let runtime_emitter = NATSRespondEmitter::new(client.clone());
25+
26+
let mut subscription = match client
27+
.queue_subscribe(String::from("execution.*"), "taurus".into())
28+
.await
29+
{
30+
Ok(subscription) => {
31+
log::info!("Subscribed to 'execution.*'");
32+
subscription
33+
}
34+
Err(err) => {
35+
log::error!("Failed to subscribe to 'execution.*': {:?}", err);
36+
return;
37+
}
38+
};
39+
40+
while let Some(message) = subscription.next().await {
41+
process_message(
42+
message,
43+
&client,
44+
&engine,
45+
&nats_remote,
46+
&runtime_emitter,
47+
runtime_usage_service.as_ref(),
48+
)
49+
.await;
50+
}
51+
52+
log::info!("NATS worker loop ended");
53+
})
54+
}
55+
56+
async fn process_message(
57+
message: async_nats::Message,
58+
client: &async_nats::Client,
59+
engine: &ExecutionEngine,
60+
nats_remote: &NATSRemoteRuntime,
61+
runtime_emitter: &NATSRespondEmitter,
62+
runtime_usage_service: Option<&TaurusRuntimeUsageService>,
63+
) {
64+
let flow: ExecutionFlow = match ExecutionFlow::decode(&*message.payload) {
65+
Ok(flow) => flow,
66+
Err(err) => {
67+
log::error!(
68+
"Failed to deserialize flow: {:?}, payload: {:?}",
69+
err,
70+
&message.payload
71+
);
72+
return;
73+
}
74+
};
75+
76+
let flow_id = flow.flow_id;
77+
let reply_subject = message.reply.clone();
78+
let respond_count = Arc::new(AtomicUsize::new(0));
79+
let respond_count_for_emitter = respond_count.clone();
80+
let respond_emitter = |execution_id, emit_type: EmitType, value: Value| {
81+
match emit_type {
82+
EmitType::OngoingExec => {
83+
respond_count_for_emitter.fetch_add(1, Ordering::Relaxed);
84+
}
85+
EmitType::StartingExec => log::debug!("Flow execution started"),
86+
EmitType::FinishedExec => log::debug!("Flow execution finished"),
87+
EmitType::FailedExec => log::debug!("Flow execution failed"),
88+
}
89+
runtime_emitter.emit(execution_id, emit_type, value);
90+
};
91+
let (signal, runtime_usage) = execute_flow(flow, engine, nats_remote, Some(&respond_emitter));
92+
93+
let has_responded = respond_count.load(Ordering::Relaxed) > 0;
94+
let final_value = signal_to_terminal_value(signal);
95+
96+
// Stream contract: if we already emitted intermediate values, do not send an extra terminal
97+
// payload on the same reply subject.
98+
if let Some(value) = final_value
99+
&& !has_responded
100+
&& let Some(reply_subject) = reply_subject
101+
{
102+
log::info!("Returning value for flow_id {}: {:?}", flow_id, value);
103+
if let Err(err) = client
104+
.publish(reply_subject, value.encode_to_vec().into())
105+
.await
106+
{
107+
log::error!("Failed to send response: {:?}", err);
108+
}
109+
}
110+
111+
if let Some(usage_service) = runtime_usage_service {
112+
usage_service.update_runtime_usage(runtime_usage).await;
113+
}
114+
}
115+
116+
fn execute_flow(
117+
flow: ExecutionFlow,
118+
engine: &ExecutionEngine,
119+
nats_remote: &NATSRemoteRuntime,
120+
respond_emitter: Option<&dyn RespondEmitter>,
121+
) -> (Signal, RuntimeUsage) {
122+
let start = Instant::now();
123+
let flow_id = flow.flow_id;
124+
let (signal, _reason) = engine.execute_flow(flow, Some(nats_remote), respond_emitter, true);
125+
let duration_millis = start.elapsed().as_millis() as i64;
126+
127+
(
128+
signal,
129+
RuntimeUsage {
130+
flow_id,
131+
duration: duration_millis,
132+
},
133+
)
134+
}
135+
136+
fn signal_to_terminal_value(signal: Signal) -> Option<Value> {
137+
match signal {
138+
Signal::Failure(error) => {
139+
log::error!("Runtime error occurred: {:?}", error);
140+
Some(error.as_value())
141+
}
142+
Signal::Success(value) => {
143+
log::debug!("Execution ended with success signal");
144+
Some(value)
145+
}
146+
Signal::Return(value) => {
147+
log::debug!("Execution ended with return signal");
148+
Some(value)
149+
}
150+
Signal::Respond(_) => {
151+
log::debug!("Execution ended with respond signal");
152+
None
153+
}
154+
Signal::Stop => {
155+
log::debug!("Received stop signal as last signal");
156+
Some(Value {
157+
kind: Some(Kind::NullValue(0)),
158+
})
159+
}
160+
}
161+
}

0 commit comments

Comments
 (0)