feat(#1086): serialize perm changes through rebuild queue
add QueueKind::PermChange — dashboard tool-group and capability handlers no longer write the shared JSON files inline. instead they enqueue a PermChange entry; the FIFO worker applies the file write then calls rebuild_agent so the updated env var takes effect. concurrent batch-apply actions for different agents previously raced on tool-groups.json / capabilities.json (last write wins, earlier change silently dropped). serialising through the queue prevents this. dedup check extended with perm-type discriminant so tool-groups and capabilities changes for the same agent are kept as distinct entries and never collapse into one slot.
This commit is contained in:
parent
dce2bd0686
commit
eae0e875cf
5 changed files with 127 additions and 27 deletions
|
|
@ -39,6 +39,7 @@ somewhere."
|
||||||
| `Spawn` | First-deploy of a new agent (approval-driven). Same serialisation as `Rebuild` from the operator's POV. |
|
| `Spawn` | First-deploy of a new agent (approval-driven). Same serialisation as `Rebuild` from the operator's POV. |
|
||||||
| `Destroy` | For future use (`destroy --purge` does real I/O). Variant exists so the wire shape doesn't change later; not currently routed through the queue. |
|
| `Destroy` | For future use (`destroy --purge` does real I/O). Variant exists so the wire shape doesn't change later; not currently routed through the queue. |
|
||||||
| `Restart` | Stop + start a container without touching config (~5-10s). Routed through the queue so it serialises against in-flight rebuilds for the same agent — prevents a restart racing a rebuild mid-flight. Sources: dashboard ↺ button, manager `restart` MCP tool. |
|
| `Restart` | Stop + start a container without touching config (~5-10s). Routed through the queue so it serialises against in-flight rebuilds for the same agent — prevents a restart racing a rebuild mid-flight. Sources: dashboard ↺ button, manager `restart` MCP tool. |
|
||||||
|
| `PermChange` | Write a tool-group or capability change to the shared JSON file (`tool-groups.json` / `capabilities.json`), then rebuild the agent so the updated `HIVE_TOOL_GROUPS` / `HIVE_CAPABILITIES` env var takes effect. Serialising the file write through the queue prevents concurrent dashboard batch-apply actions from racing on the shared file. |
|
||||||
|
|
||||||
**Intentionally not queued** (sub-second ops): `start`, `stop`, `kill`.
|
**Intentionally not queued** (sub-second ops): `start`, `stop`, `kill`.
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -54,6 +54,7 @@ pub async fn approve(coord: Arc<Coordinator>, id: i64) -> Result<()> {
|
||||||
None,
|
None,
|
||||||
Vec::new(),
|
Vec::new(),
|
||||||
Some(id),
|
Some(id),
|
||||||
|
None,
|
||||||
);
|
);
|
||||||
coord.emit_rebuild_queue_snapshot();
|
coord.emit_rebuild_queue_snapshot();
|
||||||
Ok(())
|
Ok(())
|
||||||
|
|
@ -72,6 +73,7 @@ pub async fn approve(coord: Arc<Coordinator>, id: i64) -> Result<()> {
|
||||||
None,
|
None,
|
||||||
inputs.clone(),
|
inputs.clone(),
|
||||||
Some(id),
|
Some(id),
|
||||||
|
None,
|
||||||
);
|
);
|
||||||
// Pre-enqueue cascade rebuilds in topological order so
|
// Pre-enqueue cascade rebuilds in topological order so
|
||||||
// agents depending on updated inputs are rebuilt after the
|
// agents depending on updated inputs are rebuilt after the
|
||||||
|
|
@ -99,6 +101,7 @@ pub async fn approve(coord: Arc<Coordinator>, id: i64) -> Result<()> {
|
||||||
None,
|
None,
|
||||||
Vec::new(),
|
Vec::new(),
|
||||||
Some(id),
|
Some(id),
|
||||||
|
None,
|
||||||
);
|
);
|
||||||
coord.emit_rebuild_queue_snapshot();
|
coord.emit_rebuild_queue_snapshot();
|
||||||
Ok(())
|
Ok(())
|
||||||
|
|
|
||||||
|
|
@ -2550,16 +2550,21 @@ async fn post_tool_groups(
|
||||||
if let Some(reject) = guard_agent_name(&state, &logical).await {
|
if let Some(reject) = guard_agent_name(&state, &logical).await {
|
||||||
return reject;
|
return reject;
|
||||||
}
|
}
|
||||||
if let Err(e) = crate::tool_groups::set_groups(&logical, &body.groups) {
|
// Validate group names before queuing — fail fast so the operator
|
||||||
return error_response(&format!("set tool-groups for {logical}: {e}"));
|
// sees the error immediately rather than waiting for the worker.
|
||||||
|
if let Err(e) = crate::tool_groups::validate_groups(&body.groups) {
|
||||||
|
return error_response(&format!("invalid tool-groups for {logical}: {e}"));
|
||||||
}
|
}
|
||||||
// Trigger a rebuild so the new HIVE_TOOL_GROUPS env var takes effect.
|
// Enqueue a PermChange so the JSON file write is serialised through
|
||||||
state.coord.rebuild_queue.enqueue(
|
// the FIFO worker. Prevents concurrent batch-apply actions for
|
||||||
crate::rebuild_queue::QueueKind::Rebuild,
|
// different agents from racing on the shared tool-groups.json.
|
||||||
|
state.coord.rebuild_queue.enqueue_with_perm(
|
||||||
logical.clone(),
|
logical.clone(),
|
||||||
crate::rebuild_queue::QueueSource::Manual,
|
crate::rebuild_queue::QueueSource::Manual,
|
||||||
"tool-group change via permissions UI".to_owned(),
|
"tool-group change via permissions UI".to_owned(),
|
||||||
None,
|
crate::rebuild_queue::PermPayload::ToolGroups {
|
||||||
|
groups: body.groups.clone(),
|
||||||
|
},
|
||||||
);
|
);
|
||||||
state.coord.emit_rebuild_queue_snapshot();
|
state.coord.emit_rebuild_queue_snapshot();
|
||||||
tracing::info!(agent = %logical, groups = ?body.groups, "operator: set tool-groups via dashboard");
|
tracing::info!(agent = %logical, groups = ?body.groups, "operator: set tool-groups via dashboard");
|
||||||
|
|
@ -2614,16 +2619,16 @@ async fn post_capabilities(
|
||||||
return error_response(&format!("unknown capability: {cap}"));
|
return error_response(&format!("unknown capability: {cap}"));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if let Err(e) = crate::capabilities::set_caps(&logical, &body.caps) {
|
// Enqueue a PermChange so the JSON file write is serialised through
|
||||||
return error_response(&format!("set capabilities for {logical}: {e}"));
|
// the FIFO worker. Prevents concurrent batch-apply actions for
|
||||||
}
|
// different agents from racing on the shared capabilities.json.
|
||||||
// Trigger a rebuild so the new HIVE_CAPABILITIES env var takes effect.
|
state.coord.rebuild_queue.enqueue_with_perm(
|
||||||
state.coord.rebuild_queue.enqueue(
|
|
||||||
crate::rebuild_queue::QueueKind::Rebuild,
|
|
||||||
logical.clone(),
|
logical.clone(),
|
||||||
crate::rebuild_queue::QueueSource::Manual,
|
crate::rebuild_queue::QueueSource::Manual,
|
||||||
"capability change via dashboard".to_owned(),
|
"capability change via dashboard".to_owned(),
|
||||||
None,
|
crate::rebuild_queue::PermPayload::Capabilities {
|
||||||
|
caps: body.caps.clone(),
|
||||||
|
},
|
||||||
);
|
);
|
||||||
state.coord.emit_rebuild_queue_snapshot();
|
state.coord.emit_rebuild_queue_snapshot();
|
||||||
tracing::info!(agent = %logical, caps = ?body.caps, "operator: set capabilities via dashboard");
|
tracing::info!(agent = %logical, caps = ?body.caps, "operator: set capabilities via dashboard");
|
||||||
|
|
|
||||||
|
|
@ -7,7 +7,8 @@
|
||||||
use std::collections::VecDeque;
|
use std::collections::VecDeque;
|
||||||
use std::sync::Mutex;
|
use std::sync::Mutex;
|
||||||
|
|
||||||
use serde::Serialize;
|
use anyhow::Context as _;
|
||||||
|
use serde::{Deserialize, Serialize};
|
||||||
use tokio::sync::Notify;
|
use tokio::sync::Notify;
|
||||||
|
|
||||||
/// What the queue can run. Each variant maps to a specific worker
|
/// What the queue can run. Each variant maps to a specific worker
|
||||||
|
|
@ -36,6 +37,11 @@ pub enum QueueKind {
|
||||||
/// Queued so it serialises against in-flight rebuilds for the same
|
/// Queued so it serialises against in-flight rebuilds for the same
|
||||||
/// agent — prevents a restart racing a rebuild mid-flight.
|
/// agent — prevents a restart racing a rebuild mid-flight.
|
||||||
Restart,
|
Restart,
|
||||||
|
/// Write a tool-group or capability change to the shared JSON file,
|
||||||
|
/// then rebuild the agent so the new env var takes effect.
|
||||||
|
/// Serialised through the queue so concurrent dashboard batch-apply
|
||||||
|
/// actions for different agents never race on the shared JSON file.
|
||||||
|
PermChange,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl QueueKind {
|
impl QueueKind {
|
||||||
|
|
@ -47,10 +53,24 @@ impl QueueKind {
|
||||||
QueueKind::Destroy => "destroy",
|
QueueKind::Destroy => "destroy",
|
||||||
QueueKind::StartupSweep => "startup_sweep",
|
QueueKind::StartupSweep => "startup_sweep",
|
||||||
QueueKind::Restart => "restart",
|
QueueKind::Restart => "restart",
|
||||||
|
QueueKind::PermChange => "perm_change",
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Kind-specific payload for `QueueKind::PermChange` entries.
|
||||||
|
/// Carries the desired new value so the worker can apply the file
|
||||||
|
/// write (serialised, in FIFO order) without racing concurrent HTTP
|
||||||
|
/// handlers writing to the same shared JSON file.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||||
|
#[serde(tag = "type", rename_all = "snake_case")]
|
||||||
|
pub enum PermPayload {
|
||||||
|
/// Set the tool groups for one agent (`tool-groups.json`).
|
||||||
|
ToolGroups { groups: Vec<String> },
|
||||||
|
/// Set the capabilities for one agent (`capabilities.json`).
|
||||||
|
Capabilities { caps: Vec<String> },
|
||||||
|
}
|
||||||
|
|
||||||
/// Where the enqueue request originated. Drives the "why" chip on the
|
/// Where the enqueue request originated. Drives the "why" chip on the
|
||||||
/// dashboard and lets the UI group cascade entries under their parent
|
/// dashboard and lets the UI group cascade entries under their parent
|
||||||
/// without parsing the reason text.
|
/// without parsing the reason text.
|
||||||
|
|
@ -181,6 +201,11 @@ pub struct QueueEntry {
|
||||||
/// pipeline; the kind-specific worker is the source of truth.
|
/// pipeline; the kind-specific worker is the source of truth.
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub step: Option<String>,
|
pub step: Option<String>,
|
||||||
|
/// `PermChange`-only payload: the desired new permission value to
|
||||||
|
/// apply. Absent (`None`) on all other entry kinds — omitted from
|
||||||
|
/// the wire in those cases.
|
||||||
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
|
pub perm_payload: Option<PermPayload>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// How many terminal-state entries (`Done` / `Failed` / `Cancelled`)
|
/// How many terminal-state entries (`Done` / `Failed` / `Cancelled`)
|
||||||
|
|
@ -252,7 +277,7 @@ impl RebuildQueue {
|
||||||
reason: String,
|
reason: String,
|
||||||
parent_id: Option<u64>,
|
parent_id: Option<u64>,
|
||||||
) -> u64 {
|
) -> u64 {
|
||||||
self.enqueue_full(kind, agent, source, reason, parent_id, Vec::new(), None)
|
self.enqueue_full(kind, agent, source, reason, parent_id, Vec::new(), None, None)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Same as `enqueue` but carries an `inputs` payload — used by
|
/// Same as `enqueue` but carries an `inputs` payload — used by
|
||||||
|
|
@ -269,19 +294,40 @@ impl RebuildQueue {
|
||||||
parent_id: Option<u64>,
|
parent_id: Option<u64>,
|
||||||
inputs: Vec<String>,
|
inputs: Vec<String>,
|
||||||
) -> u64 {
|
) -> u64 {
|
||||||
self.enqueue_full(kind, agent, source, reason, parent_id, inputs, None)
|
self.enqueue_full(kind, agent, source, reason, parent_id, inputs, None, None)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Enqueue a `PermChange` entry for `agent`. The worker applies the
|
||||||
|
/// JSON file write (serialised through FIFO) then rebuilds the
|
||||||
|
/// container so the updated env var takes effect.
|
||||||
|
pub fn enqueue_with_perm(
|
||||||
|
&self,
|
||||||
|
agent: String,
|
||||||
|
source: QueueSource,
|
||||||
|
reason: String,
|
||||||
|
payload: PermPayload,
|
||||||
|
) -> u64 {
|
||||||
|
self.enqueue_full(
|
||||||
|
QueueKind::PermChange,
|
||||||
|
agent,
|
||||||
|
source,
|
||||||
|
reason,
|
||||||
|
None,
|
||||||
|
Vec::new(),
|
||||||
|
None,
|
||||||
|
Some(payload),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Full-shape enqueue — every `QueueEntry` field that's settable
|
/// Full-shape enqueue — every `QueueEntry` field that's settable
|
||||||
/// at submit time. Existing `enqueue` / `enqueue_with_inputs`
|
/// at submit time. Existing `enqueue` / `enqueue_with_inputs` /
|
||||||
/// delegate to this with `approval_id: None`; the approval-driven
|
/// `enqueue_with_perm` delegate to this; the approval-driven POST
|
||||||
/// POST handlers call it directly with the source row's id so the
|
/// handlers call it directly with the source row's id so the
|
||||||
/// worker can re-fetch the kind-specific payload.
|
/// worker can re-fetch the kind-specific payload.
|
||||||
// 8/7 args: the queue entry has 6 independent submit-time fields plus
|
// 9 args: the queue entry has 6 independent submit-time fields plus
|
||||||
// the inputs/approval_id pair specific to MetaUpdate and approval
|
// three kind-specific payload fields (inputs, approval_id, perm_payload).
|
||||||
// entries. A builder struct would obscure the call sites; the
|
// A builder struct would obscure the call sites; the shorter wrappers
|
||||||
// shorter `enqueue` / `enqueue_with_inputs` wrappers already cover
|
// already cover all common cases.
|
||||||
// the common cases.
|
|
||||||
#[allow(clippy::too_many_arguments)]
|
#[allow(clippy::too_many_arguments)]
|
||||||
pub fn enqueue_full(
|
pub fn enqueue_full(
|
||||||
&self,
|
&self,
|
||||||
|
|
@ -292,6 +338,7 @@ impl RebuildQueue {
|
||||||
parent_id: Option<u64>,
|
parent_id: Option<u64>,
|
||||||
inputs: Vec<String>,
|
inputs: Vec<String>,
|
||||||
approval_id: Option<i64>,
|
approval_id: Option<i64>,
|
||||||
|
perm_payload: Option<PermPayload>,
|
||||||
) -> u64 {
|
) -> u64 {
|
||||||
let mut inner = self.inner.lock().expect("rebuild_queue mutex poisoned");
|
let mut inner = self.inner.lock().expect("rebuild_queue mutex poisoned");
|
||||||
// Dedup against a pending entry with the same (kind, agent) —
|
// Dedup against a pending entry with the same (kind, agent) —
|
||||||
|
|
@ -301,14 +348,29 @@ impl RebuildQueue {
|
||||||
// agent never collapse into one queue slot. Rebuild (and Spawn /
|
// agent never collapse into one queue slot. Rebuild (and Spawn /
|
||||||
// Destroy) entries also require parent_id to match so a
|
// Destroy) entries also require parent_id to match so a
|
||||||
// MetaUpdate cascade rebuild is never swallowed by an unrelated
|
// MetaUpdate cascade rebuild is never swallowed by an unrelated
|
||||||
// queued rebuild (e.g. from the startup sweep).
|
// queued rebuild (e.g. from the startup sweep). PermChange
|
||||||
|
// entries additionally check the perm type discriminant — a
|
||||||
|
// tool-groups change and a capabilities change for the same
|
||||||
|
// agent are distinct operations and must not collapse into one.
|
||||||
for entry in &mut inner.entries {
|
for entry in &mut inner.entries {
|
||||||
|
let perm_type_matches = match (&entry.perm_payload, &perm_payload) {
|
||||||
|
(Some(PermPayload::ToolGroups { .. }), Some(PermPayload::ToolGroups { .. })) => {
|
||||||
|
true
|
||||||
|
}
|
||||||
|
(
|
||||||
|
Some(PermPayload::Capabilities { .. }),
|
||||||
|
Some(PermPayload::Capabilities { .. }),
|
||||||
|
) => true,
|
||||||
|
(None, None) => true,
|
||||||
|
_ => false,
|
||||||
|
};
|
||||||
if entry.state == QueueState::Queued
|
if entry.state == QueueState::Queued
|
||||||
&& entry.kind == kind
|
&& entry.kind == kind
|
||||||
&& entry.agent == agent
|
&& entry.agent == agent
|
||||||
&& (kind != QueueKind::MetaUpdate || entry.inputs == inputs)
|
&& (kind != QueueKind::MetaUpdate || entry.inputs == inputs)
|
||||||
&& entry.approval_id == approval_id
|
&& entry.approval_id == approval_id
|
||||||
&& entry.parent_id == parent_id
|
&& entry.parent_id == parent_id
|
||||||
|
&& perm_type_matches
|
||||||
{
|
{
|
||||||
if !entry.reason.contains(&reason) {
|
if !entry.reason.contains(&reason) {
|
||||||
use std::fmt::Write as _;
|
use std::fmt::Write as _;
|
||||||
|
|
@ -334,6 +396,7 @@ impl RebuildQueue {
|
||||||
inputs,
|
inputs,
|
||||||
approval_id,
|
approval_id,
|
||||||
step: None,
|
step: None,
|
||||||
|
perm_payload,
|
||||||
};
|
};
|
||||||
inner.entries.push_back(entry);
|
inner.entries.push_back(entry);
|
||||||
// Wake the worker. `notify_one` is a no-op when there's no
|
// Wake the worker. `notify_one` is a no-op when there's no
|
||||||
|
|
@ -605,6 +668,34 @@ async fn dispatch(
|
||||||
coord.rescan_containers_and_emit().await;
|
coord.rescan_containers_and_emit().await;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
(QueueKind::PermChange, _) => {
|
||||||
|
let name = &entry.agent;
|
||||||
|
// Apply the file write first — serialised here so concurrent
|
||||||
|
// dashboard batch-apply actions never race on the shared JSON.
|
||||||
|
coord.set_queue_step(Some(entry.id), "writing perm file");
|
||||||
|
match &entry.perm_payload {
|
||||||
|
Some(PermPayload::ToolGroups { groups }) => {
|
||||||
|
crate::tool_groups::set_groups(name, groups)
|
||||||
|
.with_context(|| format!("set tool-groups for {name}"))?;
|
||||||
|
}
|
||||||
|
Some(PermPayload::Capabilities { caps }) => {
|
||||||
|
crate::capabilities::set_caps(name, caps)
|
||||||
|
.map_err(|e| anyhow::anyhow!("set capabilities for {name}: {e}"))?;
|
||||||
|
}
|
||||||
|
None => {
|
||||||
|
anyhow::bail!(
|
||||||
|
"PermChange entry id={} agent={} is missing perm_payload",
|
||||||
|
entry.id,
|
||||||
|
entry.agent,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// Now rebuild so the updated HIVE_TOOL_GROUPS / HIVE_CAPABILITIES
|
||||||
|
// env var takes effect in the container.
|
||||||
|
let current_rev =
|
||||||
|
crate::auto_update::current_flake_rev(&coord.hyperhive_flake).unwrap_or_default();
|
||||||
|
crate::auto_update::rebuild_agent(coord, name, ¤t_rev, Some(entry.id)).await
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -71,7 +71,7 @@ fn write(map: &BTreeMap<String, Vec<String>>) -> std::io::Result<()> {
|
||||||
/// Validate a slice of group name strings against `ToolGroup::ALL`.
|
/// Validate a slice of group name strings against `ToolGroup::ALL`.
|
||||||
/// Returns `Ok(())` when all names are known, or `Err` listing the
|
/// Returns `Ok(())` when all names are known, or `Err` listing the
|
||||||
/// unrecognised names so callers can surface a useful error message.
|
/// unrecognised names so callers can surface a useful error message.
|
||||||
fn validate_groups(groups: &[String]) -> anyhow::Result<()> {
|
pub fn validate_groups(groups: &[String]) -> anyhow::Result<()> {
|
||||||
let valid: std::collections::BTreeSet<&str> =
|
let valid: std::collections::BTreeSet<&str> =
|
||||||
hive_sh4re::ToolGroup::ALL.iter().map(|g| g.as_str()).collect();
|
hive_sh4re::ToolGroup::ALL.iter().map(|g| g.as_str()).collect();
|
||||||
let unknown: Vec<&str> = groups
|
let unknown: Vec<&str> = groups
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue