use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::{mpsc, oneshot};
use helix_core::ports::{Clock, FileUploader, FrameSender, HttpRequester, Storage};
use helix_core::{ExecutionShell, Tick};
pub use crate::batch_sink::BatchSink;
use crate::dispatch::dispatch_effects_with_context;
use crate::lifecycle::{
LifecycleCapability, LifecycleContext, LifecycleStage, LifecycleStatus, LifecycleTracker,
};
use crate::owned_effect::{consume_effect_batch, OwnedEffect};
use crate::pools::{spawn_http_pool_observed, spawn_upload_pool_observed, PersistJobPayload};
use crate::spawner::{feedback_channel, BoundedSpawner, FeedbackSink, Overflow};
use crate::tick_ingress::{EngineTickReceiver, EngineTickSender};
use crate::timer::TimerRegistry;
mod entrypoints;
pub(crate) mod perf_metrics;
mod rss;
mod shutdown;
mod types;
use perf_metrics::EngineMetricRecorder;
use shutdown::drain_pools;
const MAX_CONSECUTIVE_REPLIES: usize = 64;
pub use entrypoints::{run_engine_loop, run_engine_loop_stamped};
pub use types::{
register_transport, EngineDeps, TransportLifecycleEvent, TransportRegistration, TransportTable,
TransportTraceEvent, TransportTraceSink, TransportTraceStats, TRANSPORT_TRACE_QUEUE_CAPACITY,
};
#[allow(clippy::too_many_arguments)]
pub(super) async fn run_engine_loop_inner<S, H, U, E, Tr, C>(
mut shell: ExecutionShell,
mut tick_rx: EngineTickReceiver,
tick_tx: EngineTickSender,
deps: EngineDeps<S, H, U, E, C>,
mut shutdown_rx: oneshot::Receiver<()>,
mut transports: TransportTable<Tr>,
mut transport_rx: mpsc::UnboundedReceiver<TransportRegistration<Tr>>,
) where
S: Storage + Send + Sync + 'static,
H: HttpRequester + Send + Sync + 'static,
U: FileUploader + Send + Sync + 'static,
E: BatchSink + Send + Sync + 'static,
Tr: FrameSender + Send + Sync + 'static,
C: Clock,
{
let transport_lifecycle_tx = deps.transport_lifecycle_tx.clone();
let transport_trace_tx = deps.transport_trace_tx.clone();
let clock = deps.clock;
let storage = deps.storage;
let http = deps.http;
let uploader = deps.uploader;
let event_sink = deps.event_sink;
let trace = deps.trace;
let mut perf_metrics = EngineMetricRecorder::new(deps.metrics);
perf_metrics.on_engine_start();
let rss_sampler = perf_metrics.spawn_rss_sampler();
let http_n = deps.max_http_inflight.max(1);
let mut timer_registry = TimerRegistry::with_metrics(perf_metrics.sink_handle());
let lifecycle_tracker = LifecycleTracker::default();
let (reply_tx, mut reply_rx) = feedback_channel();
let feedback_sink = FeedbackSink::Stamped(reply_tx.clone());
let mut persist_workers: HashMap<String, BoundedSpawner<PersistJobPayload>> = HashMap::new();
let http_pool = spawn_http_pool_observed(
http_n,
Arc::clone(&http),
reply_tx.clone(),
Overflow::Block,
perf_metrics.sink_handle(),
);
let upload_pool = spawn_upload_pool_observed(
http_n,
Arc::clone(&uploader),
reply_tx.clone(),
Overflow::Block,
perf_metrics.sink_handle(),
);
let http_fire_pool = spawn_http_pool_observed(
http_n,
Arc::clone(&http),
reply_tx.clone(),
Overflow::DropNewest,
perf_metrics.sink_handle(),
);
let initial_batch = match shell.start() {
Ok(effects) => consume_effect_batch(effects),
Err(e) => {
perf_metrics.on_engine_stopped();
if let Some(rss_sampler) = rss_sampler {
rss_sampler.abort();
}
tracing::error!("helix engine start failed: {}", e);
return;
}
};
tracing::info!(
"helix engine started, {} initial effects",
initial_batch.effects.len()
);
if !initial_batch.effects.is_empty() {
let initial_context = LifecycleContext::new(0, None, None)
.with_capability(LifecycleCapability::Http, true)
.with_capability(LifecycleCapability::Persist, true)
.with_capability(LifecycleCapability::Effect, true)
.with_capability(LifecycleCapability::Ws, !transports.is_empty());
emit_missing_lifecycle_capabilities(&trace, &initial_context, &initial_batch.effects);
dispatch_effects_with_context(
initial_batch.effects,
&tick_tx,
&storage,
&http_pool,
&upload_pool,
&http_fire_pool,
&event_sink,
&mut timer_registry,
&feedback_sink,
&mut persist_workers,
&transports,
&trace,
&mut perf_metrics,
&initial_context,
Some(&lifecycle_tracker),
)
.await;
event_sink.flush();
}
let mut shutdown_armed = true;
let mut consecutive_replies = 0usize;
loop {
let iteration_started = perf_metrics.start_timer();
perf_metrics.record_queue_snapshot(tick_rx.len(), tick_rx.max_capacity());
let forced_ingress = if consecutive_replies >= MAX_CONSECUTIVE_REPLIES {
let forced = tick_rx.try_recv().ok();
if forced.is_some() {
perf_metrics.on_reply_fairness_forced();
}
forced
} else {
None
};
let (tick, queue_wait, ingress_carrier) = if let Some(queued) = forced_ingress {
queued
} else {
tokio::select! {
biased;
Some(feedback) = reply_rx.recv() => {
let wait = feedback.enqueued_at.elapsed();
perf_metrics.on_feedback_dequeued(&feedback.tick, wait, reply_rx.len());
(feedback.tick, None, None)
},
Some(registration) = transport_rx.recv() => {
let TransportRegistration { id, sender, registered_tx } = registration;
tracing::info!(transport_id = id.raw(), "transport 连接成功,注册入路由表");
transports.insert(id, sender);
registered_tx.send(()).ok();
continue;
}
res = &mut shutdown_rx, if shutdown_armed => {
shutdown_armed = false;
if res.is_ok() { break; }
continue;
}
queued = tick_rx.recv() => {
match queued {
Some(t) => t,
None => break, }
}
}
};
if matches!(tick, Tick::PortReply { .. } | Tick::PortProgress { .. }) {
consecutive_replies = consecutive_replies.saturating_add(1);
} else {
consecutive_replies = 0;
}
let metric_context = perf_metrics.on_tick(&tick, queue_wait);
if let Tick::PortReply { corr, outcome } = &tick {
let outcome_class = if matches!(outcome, helix_core::tick::PortOutcome::Ok(_)) {
"ok"
} else {
"err"
};
tracing::debug!(
hop = "host.port_reply",
corr = corr.raw(),
outcome = outcome_class,
"correlated port reply returned to the deterministic shell"
);
}
if let Tick::PortProgress { corr, progress } = &tick {
tracing::debug!(
hop = "host.port_progress",
corr = corr.raw(),
completed_bytes = progress.completed_bytes,
total_bytes = progress.total_bytes,
"correlated upload progress returned to the deterministic shell"
);
}
if let Tick::Disconnected(transport_id) = &tick {
if let Some(tx) = &transport_lifecycle_tx {
tx.send(TransportLifecycleEvent::Disconnected {
transport_id: *transport_id,
reason: "reader_closed_or_error",
})
.ok();
}
if let Some(tx) = &transport_trace_tx {
tx.try_emit(TransportTraceEvent {
transport_id: *transport_id,
name: "helix.ws.disconnect",
action: "disconnect",
attempt: None,
delay_ms: None,
next_delay_ms: None,
reason: Some("reader_closed_or_error"),
error_class: None,
});
}
}
let tick_id = lifecycle_tracker.next_tick_id();
let (parent_tick_id, inherited_carrier, inherited_span_parent) =
lifecycle_tracker.parent_for_tick(&tick);
let context = trace
.context_for_tick(
&tick,
tick_id,
parent_tick_id,
ingress_carrier.or(inherited_carrier),
)
.with_capability(LifecycleCapability::Http, true)
.with_capability(LifecycleCapability::Persist, true)
.with_capability(LifecycleCapability::Effect, true)
.with_capability(LifecycleCapability::Ws, !transports.is_empty());
let context = context.with_span_parent(inherited_span_parent);
let (_trace_scope, context) = trace.start_tick_with_context(&tick, &context);
let now_ms = clock.now_ms();
let step_started = perf_metrics.start_timer();
let batch = match shell.step(tick, now_ms) {
Ok(effects) => consume_effect_batch(effects),
Err(e) => {
perf_metrics.on_step_error(&metric_context);
perf_metrics.on_loop_iteration(iteration_started);
tracing::error!("engine step error: {}", e);
continue;
}
};
let context = context.with_capability(
LifecycleCapability::Ws,
context.has_capability(LifecycleCapability::Ws)
|| batch
.effects
.iter()
.any(|effect| matches!(effect, OwnedEffect::Send { .. })),
);
emit_missing_lifecycle_capabilities(&trace, &context, &batch.effects);
perf_metrics.on_step_complete(&metric_context, step_started, batch.summary);
if !batch.effects.is_empty() {
dispatch_effects_with_context(
batch.effects,
&tick_tx,
&storage,
&http_pool,
&upload_pool,
&http_fire_pool,
&event_sink,
&mut timer_registry,
&feedback_sink,
&mut persist_workers,
&transports,
&trace,
&mut perf_metrics,
&context,
Some(&lifecycle_tracker),
)
.await;
event_sink.flush();
}
perf_metrics.on_dispatch_complete(&metric_context);
perf_metrics.on_loop_iteration(iteration_started);
}
let shutdown_started = perf_metrics.start_timer();
drain_pools(
persist_workers,
upload_pool,
http_pool,
http_fire_pool,
reply_rx,
)
.await;
perf_metrics.on_shutdown(shutdown_started);
perf_metrics.on_engine_stopped();
if let Some(rss_sampler) = rss_sampler {
rss_sampler.abort();
}
tracing::info!(
"helix engine loop exited (persists drained; 必达 upload/http drained or aborted at limit)"
);
}
fn emit_missing_lifecycle_capabilities(
trace: &crate::trace::TraceHooks,
context: &LifecycleContext,
effects: &[OwnedEffect],
) {
let has_http = effects.iter().any(|effect| {
matches!(
effect,
OwnedEffect::Http { .. }
| OwnedEffect::HttpFire { .. }
| OwnedEffect::UploadFile { .. }
)
});
let has_ws = effects
.iter()
.any(|effect| matches!(effect, OwnedEffect::Send { .. }));
let has_persist = effects.iter().any(|effect| {
matches!(
effect,
OwnedEffect::Persist { .. }
| OwnedEffect::PersistAtomic { .. }
| OwnedEffect::PersistFire { .. }
)
});
let has_event = effects
.iter()
.any(|effect| matches!(effect, OwnedEffect::Emit { .. }));
if !has_http {
trace.emit_lifecycle_status_with_reason(
context,
LifecycleStage::T2,
Some(LifecycleCapability::Http),
LifecycleStatus::NotApplicable,
Some("http_effect_absent"),
);
}
if !has_ws {
trace.emit_lifecycle_status_with_reason(
context,
LifecycleStage::T3,
Some(LifecycleCapability::Ws),
LifecycleStatus::NotApplicable,
Some("ws_effect_absent"),
);
}
if !has_persist {
trace.emit_lifecycle_status_with_reason(
context,
LifecycleStage::T4,
Some(LifecycleCapability::Persist),
LifecycleStatus::NotApplicable,
Some("persist_effect_absent"),
);
}
if !has_event {
trace.emit_lifecycle_status_with_reason(
context,
LifecycleStage::T5,
Some(LifecycleCapability::Effect),
LifecycleStatus::NotApplicable,
Some("event_effect_absent"),
);
}
}