feat: live SSE updates for the SYST3M reminders section
Add RemindersChanged SSE event so the pending-reminders list in the SYST3M tab updates live without polling. Backend emission sites (every path that mutates the reminders table): - agent_server: store_remind (remind MCP call) - dashboard.rs: post_cancel_reminder, post_retry_reminder - questions.rs: cancel_loose_end Reminder kind - reminder_scheduler: after each delivery batch (any_delivered) Coordinator gets emit_reminders_snapshot() mirroring the existing emit_schedules_snapshot() pattern: lists PendingReminder rows from the broker and emits DashboardEvent::RemindersChanged. Frontend: applyRemindersChanged(ev) calls renderReminders(ev.reminders) and is registered as reminders_changed in MUTATION_HANDLERS. Docs: dashboard.md reminders_changed entry; CLAUDE.md file map updated.
This commit is contained in:
parent
68108fe5f8
commit
a1e46e2b3d
9 changed files with 55 additions and 1 deletions
|
|
@ -82,7 +82,8 @@ hive-c0re/ host daemon + sibling operator CLI (lib + 2 bins)
|
||||||
(`ApprovalAdded` / `ApprovalResolved`,
|
(`ApprovalAdded` / `ApprovalResolved`,
|
||||||
`QuestionAdded` / `QuestionResolved`,
|
`QuestionAdded` / `QuestionResolved`,
|
||||||
`TransientSet` / `TransientCleared`,
|
`TransientSet` / `TransientCleared`,
|
||||||
`RebuildQueueChanged`, `SchedulesChanged`).
|
`RebuildQueueChanged`, `SchedulesChanged`,
|
||||||
|
`RemindersChanged`).
|
||||||
Each frame carries a
|
Each frame carries a
|
||||||
monotonic per-process `seq` clients use to
|
monotonic per-process `seq` clients use to
|
||||||
dedupe against snapshot reads.
|
dedupe against snapshot reads.
|
||||||
|
|
|
||||||
|
|
@ -852,6 +852,14 @@ payload):
|
||||||
re-renders `schedulesState` on receipt; tab activation still
|
re-renders `schedulesState` on receipt; tab activation still
|
||||||
re-fetches as a safety net for approval-path inserts and
|
re-fetches as a safety net for approval-path inserts and
|
||||||
disconnect windows.
|
disconnect windows.
|
||||||
|
- `reminders_changed` (seq, reminders: `Vec<PendingReminder>`) —
|
||||||
|
full snapshot of all pending reminders. Emitted after every
|
||||||
|
reminder mutation: agent `remind` calls (`agent_server`),
|
||||||
|
operator cancel / retry (`/api/system/reminders/*`), `cancel_loose_end`
|
||||||
|
with Reminder kind, and the scheduler tick after each delivery
|
||||||
|
batch (`reminder_scheduler`). The SYST3M tab's reminders section
|
||||||
|
subscribes and calls `renderReminders` on receipt, so the list
|
||||||
|
updates live without polling.
|
||||||
|
|
||||||
`/api/state` is **only fetched on cold-load and on the few
|
`/api/state` is **only fetched on cold-load and on the few
|
||||||
forms that mutate non-event-derived state** (PURG3 +
|
forms that mutate non-event-derived state** (PURG3 +
|
||||||
|
|
|
||||||
|
|
@ -227,6 +227,9 @@ window.marked = marked;
|
||||||
schedulesState = (ev.schedules || []).slice();
|
schedulesState = (ev.schedules || []).slice();
|
||||||
renderSchedulesList();
|
renderSchedulesList();
|
||||||
}
|
}
|
||||||
|
function applyRemindersChanged(ev) {
|
||||||
|
renderReminders(ev.reminders || []);
|
||||||
|
}
|
||||||
// Map from agent name → highest-priority in-flight queue entry
|
// Map from agent name → highest-priority in-flight queue entry
|
||||||
// (`running` beats `queued`). Used by the container row renderer
|
// (`running` beats `queued`). Used by the container row renderer
|
||||||
// to surface "building..." / "meta-updating..." badges on the
|
// to surface "building..." / "meta-updating..." badges on the
|
||||||
|
|
@ -3604,6 +3607,7 @@ window.marked = marked;
|
||||||
meta_update_running: applyMetaUpdateRunning,
|
meta_update_running: applyMetaUpdateRunning,
|
||||||
rebuild_queue_changed: applyRebuildQueueChanged,
|
rebuild_queue_changed: applyRebuildQueueChanged,
|
||||||
schedules_changed: applySchedulesChanged,
|
schedules_changed: applySchedulesChanged,
|
||||||
|
reminders_changed: applyRemindersChanged,
|
||||||
};
|
};
|
||||||
(function bindDashboardStream() {
|
(function bindDashboardStream() {
|
||||||
// Route through the SharedWorker so all open hyperhive tabs share
|
// Route through the SharedWorker so all open hyperhive tabs share
|
||||||
|
|
|
||||||
|
|
@ -811,6 +811,7 @@ pub(crate) fn store_remind(
|
||||||
.store_reminder(agent, &stored_message, stored_path.as_deref(), due_at)
|
.store_reminder(agent, &stored_message, stored_path.as_deref(), due_at)
|
||||||
.map_err(|e| format!("failed to store reminder: {e:#}"))?;
|
.map_err(|e| format!("failed to store reminder: {e:#}"))?;
|
||||||
tracing::info!(%id, %agent, %due_at, "reminder scheduled");
|
tracing::info!(%id, %agent, %due_at, "reminder scheduled");
|
||||||
|
coord.emit_reminders_snapshot();
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -389,6 +389,24 @@ impl Coordinator {
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Emit a `RemindersChanged` snapshot event. Called from every
|
||||||
|
/// reminder mutation site (agent `remind` calls, operator cancel /
|
||||||
|
/// retry, and the scheduler after each delivery batch) so the
|
||||||
|
/// dashboard's pending-reminders list stays live without polling.
|
||||||
|
pub fn emit_reminders_snapshot(self: &Arc<Self>) {
|
||||||
|
let reminders = match self.broker.list_pending_reminders() {
|
||||||
|
Ok(rows) => rows,
|
||||||
|
Err(e) => {
|
||||||
|
tracing::warn!(error = ?e, "emit_reminders_snapshot: list failed");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
self.emit_dashboard_event(DashboardEvent::RemindersChanged {
|
||||||
|
seq: self.next_seq(),
|
||||||
|
reminders,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
/// Update the `step` label on a running queue entry and (if it
|
/// Update the `step` label on a running queue entry and (if it
|
||||||
/// actually changed) re-emit the queue snapshot so the dashboard
|
/// actually changed) re-emit the queue snapshot so the dashboard
|
||||||
/// renders the new phase. Returns `true` when the label was new
|
/// renders the new phase. Returns `true` when the label was new
|
||||||
|
|
|
||||||
|
|
@ -2141,6 +2141,7 @@ async fn post_cancel_reminder(
|
||||||
Ok(0) => error_response(&format!("reminder {id} not pending (already delivered?)")),
|
Ok(0) => error_response(&format!("reminder {id} not pending (already delivered?)")),
|
||||||
Ok(_) => {
|
Ok(_) => {
|
||||||
tracing::info!(%id, "operator cancelled reminder");
|
tracing::info!(%id, "operator cancelled reminder");
|
||||||
|
state.coord.emit_reminders_snapshot();
|
||||||
(StatusCode::OK, "ok").into_response()
|
(StatusCode::OK, "ok").into_response()
|
||||||
}
|
}
|
||||||
Err(e) => error_response(&format!("cancel reminder {id} failed: {e:#}")),
|
Err(e) => error_response(&format!("cancel reminder {id} failed: {e:#}")),
|
||||||
|
|
@ -2160,6 +2161,7 @@ async fn post_retry_reminder(
|
||||||
Ok(0) => error_response(&format!("reminder {id} not pending (already delivered?)")),
|
Ok(0) => error_response(&format!("reminder {id} not pending (already delivered?)")),
|
||||||
Ok(_) => {
|
Ok(_) => {
|
||||||
tracing::info!(%id, "operator reset reminder failure for retry");
|
tracing::info!(%id, "operator reset reminder failure for retry");
|
||||||
|
state.coord.emit_reminders_snapshot();
|
||||||
(StatusCode::OK, "ok").into_response()
|
(StatusCode::OK, "ok").into_response()
|
||||||
}
|
}
|
||||||
Err(e) => error_response(&format!("retry reminder {id} failed: {e:#}")),
|
Err(e) => error_response(&format!("retry reminder {id} failed: {e:#}")),
|
||||||
|
|
|
||||||
|
|
@ -202,6 +202,14 @@ pub enum DashboardEvent {
|
||||||
seq: u64,
|
seq: u64,
|
||||||
schedules: Vec<hive_sh4re::WireSchedule>,
|
schedules: Vec<hive_sh4re::WireSchedule>,
|
||||||
},
|
},
|
||||||
|
/// Full snapshot of all pending reminders. Emitted after every
|
||||||
|
/// reminder mutation: agent `remind` calls, operator cancel / retry,
|
||||||
|
/// and the scheduler tick after each delivery batch. Lets the
|
||||||
|
/// dashboard's reminders section stay live without polling.
|
||||||
|
RemindersChanged {
|
||||||
|
seq: u64,
|
||||||
|
reminders: Vec<crate::broker::PendingReminder>,
|
||||||
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DashboardEvent {
|
impl DashboardEvent {
|
||||||
|
|
@ -233,6 +241,7 @@ impl DashboardEvent {
|
||||||
DashboardEvent::MetaUpdateRunning { .. } => "meta_update_running",
|
DashboardEvent::MetaUpdateRunning { .. } => "meta_update_running",
|
||||||
DashboardEvent::RebuildQueueChanged { .. } => "rebuild_queue_changed",
|
DashboardEvent::RebuildQueueChanged { .. } => "rebuild_queue_changed",
|
||||||
DashboardEvent::SchedulesChanged { .. } => "schedules_changed",
|
DashboardEvent::SchedulesChanged { .. } => "schedules_changed",
|
||||||
|
DashboardEvent::RemindersChanged { .. } => "reminders_changed",
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -355,6 +364,10 @@ mod tests {
|
||||||
seq: 1,
|
seq: 1,
|
||||||
schedules: Vec::new(),
|
schedules: Vec::new(),
|
||||||
},
|
},
|
||||||
|
DashboardEvent::RemindersChanged {
|
||||||
|
seq: 1,
|
||||||
|
reminders: Vec::new(),
|
||||||
|
},
|
||||||
];
|
];
|
||||||
for ev in samples {
|
for ev in samples {
|
||||||
let v: serde_json::Value = serde_json::to_value(&ev).expect("serialise");
|
let v: serde_json::Value = serde_json::to_value(&ev).expect("serialise");
|
||||||
|
|
|
||||||
|
|
@ -178,6 +178,7 @@ pub fn handle_cancel_loose_end(
|
||||||
.cancel_reminder_as(id, canceller)
|
.cancel_reminder_as(id, canceller)
|
||||||
.map_err(|e| format!("{e:#}"))?;
|
.map_err(|e| format!("{e:#}"))?;
|
||||||
tracing::info!(%id, %canceller, %owner, "reminder cancelled");
|
tracing::info!(%id, %canceller, %owner, "reminder cancelled");
|
||||||
|
coord.emit_reminders_snapshot();
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
hive_sh4re::CancelLooseEndKind::Approval => {
|
hive_sh4re::CancelLooseEndKind::Approval => {
|
||||||
|
|
|
||||||
|
|
@ -63,6 +63,7 @@ fn tick(coord: &Arc<Coordinator>) {
|
||||||
// Single-transaction batch: one DB lock acquisition for N reminders
|
// Single-transaction batch: one DB lock acquisition for N reminders
|
||||||
// instead of N sequential lock/unlock cycles.
|
// instead of N sequential lock/unlock cycles.
|
||||||
let results = coord.broker.deliver_reminders_batch(&items);
|
let results = coord.broker.deliver_reminders_batch(&items);
|
||||||
|
let any_delivered = results.iter().any(|r| r.is_ok());
|
||||||
for ((id, agent, _body), result) in items.iter().zip(results.iter()) {
|
for ((id, agent, _body), result) in items.iter().zip(results.iter()) {
|
||||||
if let Err(e) = result {
|
if let Err(e) = result {
|
||||||
let reason = format!("{e:#}");
|
let reason = format!("{e:#}");
|
||||||
|
|
@ -82,6 +83,11 @@ fn tick(coord: &Arc<Coordinator>) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
// Emit after the batch so the dashboard's pending-reminders list
|
||||||
|
// updates when deliveries land (removes delivered rows).
|
||||||
|
if any_delivered {
|
||||||
|
coord.emit_reminders_snapshot();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Build the inbox body for a due reminder. When `file_path` is None
|
/// Build the inbox body for a due reminder. When `file_path` is None
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue