feat(msg-bus): wire jobs follow up

This commit is contained in:
Edouard Vanbelle
2026-09-11 19:41:18 +02:00
parent 84ea005b65
commit a6138aa4d9
16 changed files with 766 additions and 82 deletions
+63 -9
View File
@@ -94,8 +94,11 @@ async fn run(registry: Arc<JobRegistry>) {
// Fire and forget from the supervisor's perspective — we
// don't care about the outcome, `dispatch` records it on the
// entry and emits the log line itself. Periodic ticks never
// force — that's an admin-trigger-only affordance.
let _ = dispatch(&name, entry, &JobRunArgs::default()).await;
// force — that's an admin-trigger-only affordance. Pass the
// bus reference so periodic runs also publish job events
// (same reasoning as the manual-trigger path).
let bus = registry.message_bus_snapshot();
let _ = dispatch(&name, entry, &JobRunArgs::default(), bus).await;
}
}
@@ -117,7 +120,12 @@ async fn run(registry: Arc<JobRegistry>) {
/// `args` is passed through to `JobHandler::run`. The supervisor's
/// periodic ticks pass `JobRunArgs::default()`; the admin trigger
/// endpoint forwards parsed query params such as `?force=true`.
pub(super) async fn dispatch(name: &str, entry: Arc<JobEntry>, args: &JobRunArgs) -> JobOutcome {
pub(super) async fn dispatch(
name: &str,
entry: Arc<JobEntry>,
args: &JobRunArgs,
bus: Option<std::sync::Arc<dyn crate::application::ports::message_bus_ports::MessageBus>>,
) -> JobOutcome {
// Try to acquire the single-permit gate. `try_acquire` is
// non-blocking — if held, we know the previous run is still
// executing and skip this tick.
@@ -161,6 +169,26 @@ pub(super) async fn dispatch(name: &str, entry: Arc<JobEntry>, args: &JobRunArgs
let started_wall = Utc::now();
let start_instant = Instant::now();
// Publish `JobRunStarted` on `Topic::Job(name)` so the admin
// job dashboard's live tab receives a "started" tick without
// polling. Silent no-op when the bus isn't wired (test setup)
// or when nobody is subscribed. `actor` is `Uuid::nil()` today
// because the scheduler doesn't carry the trigger caller
// through — the periodic supervisor has no caller, and the
// admin trigger endpoints don't thread it in. When they do,
// swap to the real UUID.
if let Some(bus) = bus.as_ref() {
use crate::application::ports::message_bus_ports::{MessageBusEvent, Topic};
bus.publish(
&Topic::Job(name.to_string()),
MessageBusEvent::JobRunStarted {
name: name.to_string(),
started_at: started_wall,
actor: uuid::Uuid::nil(),
},
);
}
// Spawn so panics land as `JoinError::is_panic()` instead of
// unwinding into the supervisor loop. Args cloned into the spawn
// scope so the borrow doesn't outlive the caller.
@@ -214,6 +242,33 @@ pub(super) async fn dispatch(name: &str, entry: Arc<JobEntry>, args: &JobRunArgs
// the diagnostic `cause` field.
log_outcome(name, &outcome, cause, elapsed_ms);
// Publish `JobRunEnded` on `Topic::Job(name)`. This is the
// signal the FE watches for to terminate its subscription
// (`useJobTopic` unsubscribes on `onEnded`). `success = false`
// covers timeout, panic, handler error — the admin dashboard
// renders the row as failed and the "click for details"
// notification (Slice E) will link to `/admin/jobs/<name>`.
// Silent no-op when the bus isn't wired.
if let Some(bus) = bus.as_ref() {
use crate::application::ports::message_bus_ports::{MessageBusEvent, Topic};
let success = outcome.is_ok();
let reason = match &outcome {
crate::infrastructure::scheduler::types::JobOutcome::Err { message } => {
Some(message.clone())
}
_ => None,
};
bus.publish(
&Topic::Job(name.to_string()),
MessageBusEvent::JobRunEnded {
name: name.to_string(),
success,
reason,
ended_at: Utc::now(),
},
);
}
drop(permit);
outcome
}
@@ -421,16 +476,15 @@ mod tests {
// Kick off dispatch 1 in the background — it holds the permit
// for ~200 ms.
let entry_bg = entry.clone();
let bg =
tokio::spawn(
async move { dispatch("overrun", entry_bg, &JobRunArgs::default()).await },
);
let bg = tokio::spawn(async move {
dispatch("overrun", entry_bg, &JobRunArgs::default(), None).await
});
// Give dispatch 1 time to grab the permit.
tokio::time::sleep(Duration::from_millis(50)).await;
// Dispatch 2 should observe the permit taken and skip.
dispatch("overrun", entry.clone(), &JobRunArgs::default()).await;
dispatch("overrun", entry.clone(), &JobRunArgs::default(), None).await;
// Only dispatch 1's handler should have actually run so far.
assert_eq!(calls.load(Ordering::SeqCst), 1);
@@ -458,7 +512,7 @@ mod tests {
.await;
let entry = registry.get("slow").await.unwrap();
dispatch("slow", entry.clone(), &JobRunArgs::default()).await;
dispatch("slow", entry.clone(), &JobRunArgs::default(), None).await;
// The timeout fired; last_outcome must be Err.
let state = entry.state.lock().unwrap();
+39 -1
View File
@@ -60,12 +60,24 @@ pub(super) struct JobState {
/// native services `register()` during DI wiring.
pub struct JobRegistry {
entries: RwLock<HashMap<String, Arc<JobEntry>>>,
/// Message bus — used by `dispatch` (via `trigger`) to publish
/// `JobRunStarted` / `JobRunProgress` / `JobRunEnded` on
/// `Topic::Job(name)` so the admin dashboard can render live
/// progress without polling. `OnceLock` because it's set exactly
/// once at DI time (after both the registry and the bus are
/// constructed) and read from many concurrent triggers; `Arc`
/// keeps consumers cheap. `None` before wiring (unit tests
/// exercise the registry without a bus).
message_bus: std::sync::OnceLock<
std::sync::Arc<dyn crate::application::ports::message_bus_ports::MessageBus>,
>,
}
impl JobRegistry {
pub fn new() -> Self {
Self {
entries: RwLock::new(HashMap::new()),
message_bus: std::sync::OnceLock::new(),
}
}
@@ -294,7 +306,33 @@ impl JobRegistry {
/// that just want a plain run pass `JobRunArgs::default()`.
pub async fn trigger(self: &Arc<Self>, name: &str, args: &JobRunArgs) -> Option<JobOutcome> {
let entry = self.get(name).await?;
Some(super::engine::dispatch(name, entry, args).await)
// Pass the bus reference through to `dispatch` so start / end
// events publish on `Topic::Job(name)`. `Option::cloned()`
// returns a fresh `Arc` clone (or None) — negligible.
let bus = self.message_bus.get().cloned();
Some(super::engine::dispatch(name, entry, args, bus).await)
}
/// Wire the message bus. Called once from DI after both the
/// registry and the bus are constructed. Idempotent: a second
/// call is a silent no-op (`OnceLock::set` returns `Err`), so
/// test setups that call this more than once don't panic.
pub fn set_message_bus(
&self,
bus: std::sync::Arc<dyn crate::application::ports::message_bus_ports::MessageBus>,
) {
let _ = self.message_bus.set(bus);
}
/// Snapshot the currently-wired bus (if any). `None` when
/// `set_message_bus` hasn't been called yet — every test setup
/// that skips DI wiring, and the very early boot before the
/// bus is constructed. Called by the periodic supervisor and
/// by `trigger` so both paths publish job events identically.
pub(super) fn message_bus_snapshot(
&self,
) -> Option<std::sync::Arc<dyn crate::application::ports::message_bus_ports::MessageBus>> {
self.message_bus.get().cloned()
}
}
@@ -112,8 +112,13 @@ impl InProcessMessageBus {
/// (for the receiver) — one code path for the map insert avoids a race
/// where publish creates a sender concurrent subscribers miss.
fn sender_for(&self, topic: &Topic) -> broadcast::Sender<MessageBusEvent> {
// `topic.clone()` because `Topic::Job(String)` isn't `Copy`.
// The clone is a String alloc on the cold path (first ever
// subscriber for a topic) and free on the hot path (existing
// entry — `entry` doesn't need to move the key when the
// entry is already present).
self.topics
.entry(*topic)
.entry(topic.clone())
.or_insert_with(|| broadcast::channel(BROADCAST_RING_CAPACITY).0)
.clone()
}
@@ -194,6 +199,9 @@ fn event_kind(event: &MessageBusEvent) -> &'static str {
MessageBusEvent::FolderMoved { .. } => "folder_moved",
MessageBusEvent::FolderDeleted { .. } => "folder_deleted",
MessageBusEvent::AuthzChanged { .. } => "authz_changed",
MessageBusEvent::JobRunStarted { .. } => "job_run_started",
MessageBusEvent::JobRunProgress { .. } => "job_run_progress",
MessageBusEvent::JobRunEnded { .. } => "job_run_ended",
}
}