- 1
//! The one place the server touches the platform service manager. - 2
//! - 3
//! Two rules, both learned from defects this module exists to make - 4
//! impossible: - 5
//! - 6
//! **1. Every call runs on a blocking worker.** `launchctl` and - 7
//! `systemctl` are subprocesses; `vak_ops`'s health probe is - 8
//! `reqwest::blocking`. Calling either straight from an async handler - 9
//! stalls a tokio worker for as long as the service manager takes to - 10
//! answer — which, for `launchctl bootstrap`, can be indefinitely - 11
//! (AGENTS.md invariant 26). Handlers `await` these functions instead of - 12
//! calling `vak_ops` directly, and nothing in `lib.rs` may reach past - 13
//! this module to the service manager. - 14
//! - 15
//! **2. A record is not an activation.** Creating a bot, renaming it, - 16
//! setting its token, or revoking it writes `bots.json` / the credential - 17
//! store and nothing else. Registering an OS service is a separate, explicit act — - 18
//! [`reconcile`] — for the same reason `vak self install` no longer starts - 19
//! services (`docs/design/46-stabilization-install-and-onboarding.md` D6): - 20
//! configuring something and activating it are different decisions, and - 21
//! collapsing them means a routine edit quietly mutates the machine's - 22
//! service manager. It also meant a test that created a bot wrote real - 23
//! launchd plists pointing at the test binary. - 24
//! - 25
//! Revocation is deliberately **not** an exception to this. Making a - 26
//! cleared token effective by bouncing a process would put orchestration - 27
//! back in the request path — and it did: an API handler shelled out to - 28
//! `launchctl`, which blocked. A bridge watches its own credential instead - 29
//! (`surfaces::CredentialWatch`), so revoking one is a fact about the - 30
//! credential store that takes effect within a poll cycle, with nothing - 31
//! orchestrating it. - 32
- 33
use serde::Serialize; - 34
- 35
/// Run a service-manager call on a blocking worker. - 36
/// - 37
/// A panic inside the closure is reported as a failed operation rather - 38
/// than taking down the handler: the service manager's opinion is never - 39
/// worth a 500 on an unrelated request. - 40
async fn blocking<T, F>(f: F) -> Result<T, String> - 41
where - 42
F: FnOnce() -> T + Send + 'static, - 43
T: Send + 'static, - 44
{ - 45
tokio::task::spawn_blocking(f) - 46
.await - 47
.map_err(|e| format!("service manager call did not complete: {e}")) - 48
} - 49
- 50
/// What reconciling produced, per unit. - 51
#[derive(Debug, Clone, Serialize)] - 52
pub struct UnitOutcome { - 53
pub name: String, - 54
pub action: String, - 55
#[serde(skip_serializing_if = "Option::is_none")] - 56
pub error: Option<String>, - 57
} - 58
- 59
/// Bring the platform service manager in line with configuration: - 60
/// the fixed services, plus one bridge unit per configured bot. - 61
/// - 62
/// This is the **only** thing that registers, updates, or removes units, - 63
/// and it is always called deliberately — by `vak setup`'s activation - 64
/// step, by `vak self services-sync`, or by an operator pressing a button - 65
/// in the admin console. Nothing reconciles as a side effect of an edit. - 66
pub async fn reconcile(core: &vak_core::Core, port: u16) -> Result<Vec<UnitOutcome>, String> { - 67
// Units exec the running binary. In the gateway service that *is* the - 68
// installed binary; in a test or a dev build it is not, and writing - 69
// units that point at a test harness is how a test comes to mutate the - 70
// developer's machine. Refuse rather than guess. - 71
let bin_path = - 72
std::env::current_exe().map_err(|e| format!("cannot locate the running binary: {e}"))?; - 73
// A bridge unit authenticates against the gateway with this, so it has - 74
// to exist *before* the unit is registered — otherwise the unit is - 75
// written, starts, and crash-loops on "gateway token missing". - 76
vak_core::gateway_token::ensure_gateway_token()?; - 77
let data_home = core.sessions_home(); - 78
let gateway_url = vak_ops::OpsConfig { port }.base_url(); - 79
- 80
blocking(move || { - 81
let mut outcomes: Vec<UnitOutcome> = Vec::new(); - 82
let names = vak_ops::services::default_service_names(&bin_path); - 83
for outcome in vak_ops::services::services_sync( - 84
&bin_path, - 85
&names, - 86
&vak_ops::services::Paths::default(), - 87
&vak_ops::services::SystemRunner, - 88
) { - 89
outcomes.push(outcome.into()); - 90
} - 91
for outcome in vak_ops::sync_bots( - 92
&bin_path, - 93
&data_home, - 94
&gateway_url, - 95
&vak_ops::services::Paths::default(), - 96
&vak_ops::services::SystemRunner, - 97
) { - 98
outcomes.push(outcome.into()); - 99
} - 100
outcomes - 101
}) - 102
.await - 103
} - 104
- 105
impl From<vak_ops::services::SyncOutcome> for UnitOutcome { - 106
fn from(outcome: vak_ops::services::SyncOutcome) -> Self { - 107
match outcome.action { - 108
vak_ops::services::SyncAction::Failed(error) => UnitOutcome { - 109
name: outcome.name, - 110
action: "failed".into(), - 111
error: Some(error), - 112
}, - 113
action => UnitOutcome { - 114
name: outcome.name, - 115
action: format!("{action:?}").to_lowercase(), - 116
error: None, - 117
}, - 118
} - 119
} - 120
} - 121
- 122
/// How far configuration has run ahead of what is actually registered. - 123
/// - 124
/// A bot that exists in `bots.json` with no unit is not broken — it is - 125
/// **configured but not activated**, which is a state the operator chose - 126
/// and must therefore be able to see. Before this was reported, a bot - 127
/// created in the console looked identical to one that was live. - 128
#[derive(Debug, Clone, Serialize)] - 129
pub struct ActivationDrift { - 130
/// Bots configured in `bots.json`. - 131
pub configured_bots: usize, - 132
/// Bot bridge units the service manager actually has. - 133
pub registered_units: usize, - 134
/// Bot ids with no unit registered for them. - 135
pub awaiting_activation: Vec<String>, - 136
} - 137
- 138
/// Compare configured bots against registered units. - 139
pub async fn activation_drift(core: &vak_core::Core) -> Result<ActivationDrift, String> { - 140
let data_home = core.sessions_home(); - 141
blocking(move || { - 142
let configured = vak_ops::services::configured_bot_service_names_all(&data_home); - 143
let paths = vak_ops::services::Paths::default(); - 144
let awaiting: Vec<String> = configured - 145
.iter() - 146
.filter(|name| !vak_ops::services::unit_is_registered(name, &paths)) - 147
.cloned() - 148
.collect(); - 149
ActivationDrift { - 150
configured_bots: configured.len(), - 151
registered_units: configured.len() - awaiting.len(), - 152
awaiting_activation: awaiting, - 153
} - 154
}) - 155
.await - 156
} - 157
- 158
/// Start, stop, or restart a service, on a blocking worker. - 159
pub async fn act( - 160
service: vak_ops::Service, - 161
action: Action, - 162
cfg: vak_ops::OpsConfig, - 163
) -> Result<bool, String> { - 164
blocking(move || match action { - 165
Action::Start => vak_ops::start(service, &cfg), - 166
Action::Stop => vak_ops::stop(service, &cfg), - 167
Action::Restart => vak_ops::restart(service, &cfg), - 168
}) - 169
.await - 170
} - 171
- 172
/// Probe a service's state, on a blocking worker. - 173
pub async fn status( - 174
service: vak_ops::Service, - 175
cfg: vak_ops::OpsConfig, - 176
) -> Result<vak_ops::State, String> { - 177
blocking(move || vak_ops::status(service, &cfg)).await - 178
} - 179
- 180
#[derive(Debug, Clone, Copy, PartialEq, Eq)] - 181
pub enum Action { - 182
Start, - 183
Stop, - 184
Restart, - 185
} - 186
- 187
#[cfg(test)] - 188
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 189
mod tests { - 190
use super::*; - 191
- 192
/// Configured and registered are different counts, and the gap is the - 193
/// thing an operator needs to see. - 194
#[test] - 195
fn drift_separates_what_is_configured_from_what_is_registered() { - 196
let drift = ActivationDrift { - 197
configured_bots: 2, - 198
registered_units: 1, - 199
awaiting_activation: vec!["com.vak.discord-ops".into()], - 200
}; - 201
assert_eq!(drift.configured_bots - drift.registered_units, 1); - 202
assert_eq!(drift.awaiting_activation, vec!["com.vak.discord-ops"]); - 203
} - 204
- 205
#[tokio::test] - 206
async fn a_blocking_call_returns_its_value_rather_than_stalling_the_runtime() { - 207
assert_eq!(blocking(|| 7).await.unwrap(), 7); - 208
} - 209
} - 210
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.