pub(crate) async fn nudge_daemon_reload(capsule_ids: &[String]) {
use astrid_core::kernel_api::{KernelRequest, KernelResponse};
if capsule_ids.is_empty() {
return;
}
if !daemon_socket_reachable().await {
return;
}
let session_uuid = uuid::Uuid::new_v4();
let session = astrid_core::SessionId::from_uuid(session_uuid);
let mut client = match classify_live_client(
crate::socket_client::connect_for_workspace(session, crate::principal::current(), None)
.await,
) {
Ok(Some(client)) => client,
Ok(None) => return,
Err(error) => {
eprintln!("Note: installed capsules to disk, but skipped live reload because {error}.");
return;
},
};
for id in capsule_ids {
let Ok(val) = serde_json::to_value(KernelRequest::ReloadCapsule { id: id.clone() }) else {
continue;
};
let correlation = uuid::Uuid::new_v4().simple().to_string();
let request_topic = astrid_types::Topic::reload_capsule_request(&correlation);
let response_topic = astrid_types::Topic::reload_capsule_response(&correlation);
let msg = astrid_types::ipc::IpcMessage::new(
request_topic,
astrid_types::ipc::IpcPayload::RawJson(val),
session_uuid,
);
if client.send_message(msg).await.is_err() {
continue;
}
let Ok(raw) = client
.read_until_topic(response_topic.as_str(), std::time::Duration::from_secs(15))
.await
else {
eprintln!(
"Note: installed '{id}' to disk, but the running daemon didn't confirm a live reload in time — it will load on the next daemon start."
);
continue;
};
match crate::socket_client::SocketClient::extract_kernel_response(&raw) {
Some(KernelResponse::Success(_)) => {
eprintln!("Live: the running daemon loaded '{id}' — no restart needed.");
},
Some(KernelResponse::Error(reason)) => {
eprintln!(
"Note: installed '{id}' to disk, but the daemon declined a live reload ({reason}); it will load on the next daemon start."
);
},
_ => {
eprintln!(
"Note: installed '{id}' to disk, but couldn't confirm a live reload; it will load on the next daemon start."
);
},
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum LiveUnload {
NoDaemon,
NotLoaded,
Unloaded,
}
pub(crate) async fn try_daemon_unload(capsule_id: &str) -> anyhow::Result<LiveUnload> {
use anyhow::{Context, bail};
use astrid_core::kernel_api::{KernelRequest, KernelResponse};
if capsule_id.is_empty() || !daemon_socket_reachable().await {
return Ok(LiveUnload::NoDaemon);
}
let session_uuid = uuid::Uuid::new_v4();
let session = astrid_core::SessionId::from_uuid(session_uuid);
let Some(mut client) = classify_live_client(
crate::socket_client::connect_for_workspace(session, crate::principal::current(), None)
.await,
)
.context("refusing to unload from a daemon with a different workspace selection")?
else {
return Ok(LiveUnload::NoDaemon);
};
let val = serde_json::to_value(KernelRequest::UnloadCapsule {
id: capsule_id.to_string(),
})
.context("failed to encode live capsule unload request")?;
let correlation = uuid::Uuid::new_v4().simple().to_string();
let request_topic =
astrid_types::Topic::kernel_request(format!("unload_capsule.{correlation}"));
let response_topic =
astrid_types::Topic::kernel_response(format!("unload_capsule.{correlation}"));
let msg = astrid_types::ipc::IpcMessage::new(
request_topic,
astrid_types::ipc::IpcPayload::RawJson(val),
session_uuid,
);
client
.send_message(msg)
.await
.context("failed to send live capsule unload request")?;
let raw = client
.read_until_topic(response_topic.as_str(), std::time::Duration::from_secs(15))
.await
.context("running daemon did not confirm live capsule unload")?;
match crate::socket_client::SocketClient::extract_kernel_response(&raw) {
Some(KernelResponse::Success(data)) => {
match data.get("status").and_then(serde_json::Value::as_str) {
Some("unloaded") => Ok(LiveUnload::Unloaded),
Some("not_loaded") => Ok(LiveUnload::NotLoaded),
Some(other) => bail!("running daemon returned unknown unload status {other:?}"),
None => bail!("running daemon returned unload success without a status"),
}
},
Some(KernelResponse::Error(reason)) => {
bail!("running daemon declined live capsule unload: {reason}")
},
_ => bail!("running daemon returned a malformed live capsule unload response"),
}
}
async fn daemon_socket_reachable() -> bool {
let path = crate::socket_client::proxy_socket_path();
path.exists() && tokio::net::UnixStream::connect(path).await.is_ok()
}
fn classify_live_client<T>(
result: crate::socket_client::WorkspaceConnectionResult<T>,
) -> anyhow::Result<Option<T>> {
match result {
Ok(client) => Ok(Some(client)),
Err(crate::socket_client::WorkspaceConnectionError::Connect(_)) => Ok(None),
Err(crate::socket_client::WorkspaceConnectionError::Selection(error)) => Err(error),
}
}
#[cfg(test)]
mod tests {
use super::classify_live_client;
use crate::socket_client::WorkspaceConnectionError;
#[test]
fn live_client_classification_keeps_offline_and_mismatch_distinct() {
assert!(
classify_live_client::<()>(Err(WorkspaceConnectionError::Connect(anyhow::anyhow!(
"unreachable socket"
))))
.unwrap()
.is_none()
);
assert!(
classify_live_client::<()>(Err(WorkspaceConnectionError::Selection(anyhow::anyhow!(
"workspace mismatch"
))))
.is_err()
);
}
}