use std::sync::Arc;
use helix_driver_host::engine::{self, EngineDeps, TransportLifecycleEvent, TransportTraceSink};
use helix_driver_host::{AsyncMetricSink, NoopMetricSink, SharedFileUploader, TraceHooks};
use helix_driver_host::{TickIngressReceiver, TickIngressSender};
pub use helix_driver_host::{rows_from_reply_bytes, rows_to_reply_bytes};
use helix_core::effect::TransportId;
use helix_core::{ExecutionShell, Tick};
use std::collections::HashMap;
use tokio::sync::{mpsc, oneshot};
use crate::clock::NativeClock;
use crate::event_sink::NativeEventSink;
use crate::http::NativeHttp;
use crate::storage::NativeStorage;
use crate::transport::NativeTransport;
mod reconnect;
use reconnect::{start_native_reconnect, NativeReconnectRuntime};
pub type TransportTable = HashMap<TransportId, Arc<NativeTransport>>;
pub struct EngineConfig {
pub storage: NativeStorage,
pub http: NativeHttp,
pub uploader: SharedFileUploader,
pub event_sink: NativeEventSink,
pub metrics: Arc<dyn AsyncMetricSink>,
pub max_http_inflight: usize,
}
impl EngineConfig {
pub fn new(
storage: NativeStorage,
http: NativeHttp,
uploader: SharedFileUploader,
event_sink: NativeEventSink,
max_http_inflight: usize,
) -> Self {
Self {
storage,
http,
uploader,
event_sink,
metrics: Arc::new(NoopMetricSink),
max_http_inflight,
}
}
pub fn with_metric_sink(mut self, metrics: Arc<dyn AsyncMetricSink>) -> Self {
self.metrics = metrics;
self
}
fn into_deps(
self,
trace: TraceHooks,
transport_lifecycle_tx: Option<mpsc::UnboundedSender<TransportLifecycleEvent>>,
transport_trace_tx: Option<TransportTraceSink>,
) -> EngineDeps<NativeStorage, NativeHttp, SharedFileUploader, NativeEventSink, NativeClock>
{
EngineDeps {
storage: Arc::new(self.storage),
http: Arc::new(self.http),
uploader: Arc::new(self.uploader),
event_sink: Arc::new(self.event_sink),
clock: NativeClock,
trace,
metrics: self.metrics,
max_http_inflight: self.max_http_inflight,
transport_lifecycle_tx,
transport_trace_tx,
}
}
}
pub async fn run_engine_loop(
shell: ExecutionShell,
tick_rx: mpsc::Receiver<Tick>,
tick_tx: mpsc::Sender<Tick>,
config: EngineConfig,
shutdown_rx: oneshot::Receiver<()>,
) {
run_engine_loop_with_transports(
shell,
tick_rx,
tick_tx,
config,
shutdown_rx,
TransportTable::new(),
)
.await
}
pub async fn run_engine_loop_stamped(
shell: ExecutionShell,
tick_rx: TickIngressReceiver,
tick_tx: TickIngressSender,
config: EngineConfig,
shutdown_rx: oneshot::Receiver<()>,
) {
run_engine_loop_with_transports_and_trace_stamped(
shell,
tick_rx,
tick_tx,
config,
TraceHooks::noop(),
shutdown_rx,
TransportTable::new(),
)
.await;
}
pub async fn run_engine_loop_with_transports(
shell: ExecutionShell,
tick_rx: mpsc::Receiver<Tick>,
tick_tx: mpsc::Sender<Tick>,
config: EngineConfig,
shutdown_rx: oneshot::Receiver<()>,
transports: TransportTable,
) {
run_engine_loop_with_transports_and_trace(
shell,
tick_rx,
tick_tx,
config,
TraceHooks::noop(),
shutdown_rx,
transports,
)
.await;
}
pub async fn run_engine_loop_with_transports_and_trace(
shell: ExecutionShell,
tick_rx: mpsc::Receiver<Tick>,
tick_tx: mpsc::Sender<Tick>,
config: EngineConfig,
trace: TraceHooks,
shutdown_rx: oneshot::Receiver<()>,
transports: TransportTable,
) {
let NativeReconnectRuntime {
transport_rx,
lifecycle_tx,
trace_sink,
trace_stats,
tasks,
shutdown_tx: reconnect_shutdown_tx,
} = start_native_reconnect(&transports);
engine::run_engine_loop(
shell,
tick_rx,
tick_tx,
config.into_deps(trace, lifecycle_tx, trace_sink),
shutdown_rx,
transports,
transport_rx,
)
.await;
if let Some(stats) = trace_stats {
let dropped = stats.dropped_count();
if dropped > 0 {
tracing::warn!(dropped, "native transport trace 有界队列发生丢弃");
}
}
reconnect_shutdown_tx.send_replace(true);
for task in tasks {
task.abort();
}
}
pub async fn run_engine_loop_with_transports_and_trace_stamped(
shell: ExecutionShell,
tick_rx: TickIngressReceiver,
tick_tx: TickIngressSender,
config: EngineConfig,
trace: TraceHooks,
shutdown_rx: oneshot::Receiver<()>,
transports: TransportTable,
) {
let NativeReconnectRuntime {
transport_rx,
lifecycle_tx,
trace_sink,
trace_stats,
tasks,
shutdown_tx: reconnect_shutdown_tx,
} = start_native_reconnect(&transports);
engine::run_engine_loop_stamped(
shell,
tick_rx,
tick_tx,
config.into_deps(trace, lifecycle_tx, trace_sink),
shutdown_rx,
transports,
transport_rx,
)
.await;
if let Some(stats) = trace_stats {
let dropped = stats.dropped_count();
if dropped > 0 {
tracing::warn!(dropped, "native transport trace 有界队列发生丢弃");
}
}
reconnect_shutdown_tx.send_replace(true);
for task in tasks {
task.abort();
}
}
#[cfg(test)]
#[path = "engine_loop_tests.rs"]
mod tests;