use std::convert::TryFrom;
use std::fs;
use std::io;
use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use sha2::{Digest, Sha256};
use crate::broker::host_identity;
use crate::broker::protocol::{self, CacheManifest, Endpoint};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DaemonProcess {
pub pid: u32,
pub exe_path: PathBuf,
pub exe_hash: [u8; 32],
pub legacy_exe_sha256: [u8; 32],
pub boot_id: String,
pub ipc_endpoint: Endpoint,
pub started_at_unix_ms: u64,
pub idle_timeout_secs: Option<u32>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum DaemonIdentityHashPolicy {
#[default]
LegacyCompatible,
Blake3Only,
}
impl DaemonProcess {
pub fn current_process(
ipc_endpoint: Endpoint,
idle_timeout_secs: Option<u32>,
) -> Result<Self, IdentityError> {
Self::current_process_with_hash_policy(
ipc_endpoint,
idle_timeout_secs,
DaemonIdentityHashPolicy::LegacyCompatible,
)
}
pub fn current_process_with_hash_policy(
ipc_endpoint: Endpoint,
idle_timeout_secs: Option<u32>,
hash_policy: DaemonIdentityHashPolicy,
) -> Result<Self, IdentityError> {
let exe_path = std::env::current_exe().map_err(IdentityError::CurrentExe)?;
let exe_hash = executable_hash_file(&exe_path)?;
let legacy_exe_sha256 = match hash_policy {
DaemonIdentityHashPolicy::LegacyCompatible => sha256_file(&exe_path)?,
DaemonIdentityHashPolicy::Blake3Only => [0; 32],
};
Ok(Self {
pid: std::process::id(),
exe_path,
exe_hash,
legacy_exe_sha256,
boot_id: host_identity::current().boot_id,
ipc_endpoint,
started_at_unix_ms: unix_now_ms(),
idle_timeout_secs,
})
}
pub fn to_proto(&self) -> protocol::DaemonProcess {
protocol::DaemonProcess {
pid: self.pid,
exe_path: self.exe_path.to_string_lossy().into_owned(),
exe_hash_algorithm: EXECUTABLE_HASH_ALGORITHM.to_owned(),
exe_hash: self.exe_hash.to_vec(),
ipc_endpoint: Some(self.ipc_endpoint.clone()),
started_at_unix_ms: self.started_at_unix_ms,
boot_id: self.boot_id.clone(),
idle_timeout_secs: self.idle_timeout_secs,
}
}
pub fn encode_probe_identity(&self, output: &mut Vec<u8>) -> Result<(), prost::EncodeError> {
use prost::Message;
self.to_proto().encode(output)?;
output.push(0x1a); output.push(32); output.extend_from_slice(&self.legacy_exe_sha256);
Ok(())
}
pub fn from_manifest_current_daemon(
manifest: &CacheManifest,
) -> Result<Option<Self>, IdentityError> {
manifest
.current_daemon
.clone()
.map(Self::try_from)
.transpose()
}
}
impl TryFrom<protocol::DaemonProcess> for DaemonProcess {
type Error = IdentityError;
fn try_from(value: protocol::DaemonProcess) -> Result<Self, Self::Error> {
let ipc_endpoint = value.ipc_endpoint.ok_or(IdentityError::MissingEndpoint)?;
if value.exe_hash_algorithm != EXECUTABLE_HASH_ALGORITHM {
return Err(IdentityError::UnsupportedExecutableHashAlgorithm(
value.exe_hash_algorithm,
));
}
let exe_hash =
vec_to_hash(value.exe_hash).map_err(IdentityError::InvalidExecutableHashLength)?;
Ok(Self {
pid: value.pid,
exe_path: PathBuf::from(value.exe_path),
exe_hash,
legacy_exe_sha256: [0; 32],
boot_id: value.boot_id,
ipc_endpoint,
started_at_unix_ms: value.started_at_unix_ms,
idle_timeout_secs: value.idle_timeout_secs,
})
}
}
impl From<&DaemonProcess> for protocol::DaemonProcess {
fn from(value: &DaemonProcess) -> Self {
value.to_proto()
}
}
impl Serialize for DaemonProcess {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
DaemonProcessSerde::from(self).serialize(serializer)
}
}
impl<'de> Deserialize<'de> for DaemonProcess {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let value = DaemonProcessSerde::deserialize(deserializer)?;
if value.exe_hash_algorithm != EXECUTABLE_HASH_ALGORITHM {
return Err(<D::Error as serde::de::Error>::custom(format!(
"unsupported daemon executable hash algorithm {:?}; expected blake3",
value.exe_hash_algorithm
)));
}
Ok(Self {
pid: value.pid,
exe_path: value.exe_path,
exe_hash: value.exe_hash,
legacy_exe_sha256: value.legacy_exe_sha256,
boot_id: value.boot_id,
ipc_endpoint: value.ipc_endpoint.into(),
started_at_unix_ms: value.started_at_unix_ms,
idle_timeout_secs: value.idle_timeout_secs,
})
}
}
#[derive(Debug, thiserror::Error)]
pub enum IdentityError {
#[error("daemon process is missing ipc_endpoint")]
MissingEndpoint,
#[error("unsupported daemon executable hash algorithm {0:?}; expected blake3")]
UnsupportedExecutableHashAlgorithm(String),
#[error("daemon process exe_hash must be 32 bytes, got {0}")]
InvalidExecutableHashLength(usize),
#[error("failed to resolve current executable: {0}")]
CurrentExe(io::Error),
#[error("failed to hash executable: {0}")]
Io(#[from] io::Error),
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct DaemonProcessSerde {
pid: u32,
exe_path: PathBuf,
exe_hash_algorithm: String,
exe_hash: [u8; 32],
#[serde(default)]
legacy_exe_sha256: [u8; 32],
boot_id: String,
ipc_endpoint: EndpointSerde,
started_at_unix_ms: u64,
idle_timeout_secs: Option<u32>,
}
impl From<&DaemonProcess> for DaemonProcessSerde {
fn from(value: &DaemonProcess) -> Self {
Self {
pid: value.pid,
exe_path: value.exe_path.clone(),
exe_hash_algorithm: EXECUTABLE_HASH_ALGORITHM.to_owned(),
exe_hash: value.exe_hash,
legacy_exe_sha256: value.legacy_exe_sha256,
boot_id: value.boot_id.clone(),
ipc_endpoint: EndpointSerde::from(&value.ipc_endpoint),
started_at_unix_ms: value.started_at_unix_ms,
idle_timeout_secs: value.idle_timeout_secs,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct EndpointSerde {
namespace_id: String,
path: String,
}
impl From<&Endpoint> for EndpointSerde {
fn from(value: &Endpoint) -> Self {
Self {
namespace_id: value.namespace_id.clone(),
path: value.path.clone(),
}
}
}
impl From<EndpointSerde> for Endpoint {
fn from(value: EndpointSerde) -> Self {
Endpoint {
namespace_id: value.namespace_id,
path: value.path,
}
}
}
pub fn executable_hash_file(path: &Path) -> Result<[u8; 32], io::Error> {
crate::content_hash::blake3_file(path).map(|hash| *hash.as_bytes())
}
pub fn sha256_file(path: &Path) -> Result<[u8; 32], io::Error> {
let bytes = fs::read(path)?;
let digest = Sha256::digest(&bytes);
let mut out = [0_u8; 32];
out.copy_from_slice(&digest);
Ok(out)
}
fn vec_to_hash(bytes: Vec<u8>) -> Result<[u8; 32], usize> {
let len = bytes.len();
let Ok(out) = bytes.try_into() else {
return Err(len);
};
Ok(out)
}
const EXECUTABLE_HASH_ALGORITHM: &str = "blake3";
fn unix_now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_millis() as u64)
.unwrap_or(0)
}
#[cfg(test)]
mod broker_dance_identity_tests {
use super::*;
use prost::Message;
#[derive(Clone, PartialEq, Message)]
struct LegacyDaemonProcess {
#[prost(uint32, tag = "1")]
pid: u32,
#[prost(string, tag = "2")]
exe_path: String,
#[prost(bytes = "vec", tag = "3")]
exe_sha256: Vec<u8>,
}
#[test]
fn executable_identity_hash_uses_blake3() {
let path =
std::env::temp_dir().join(format!("running-process-946-hash-{}", std::process::id()));
std::fs::write(&path, b"daemon image bytes").expect("write fixture");
let actual = executable_hash_file(&path).expect("hash fixture");
std::fs::remove_file(&path).ok();
assert_eq!(actual, *blake3::hash(b"daemon image bytes").as_bytes());
}
#[test]
fn blake3_identity_dual_writes_a_legacy_sha256_for_stable_brokers() {
let identity = DaemonProcess::current_process(endpoint("compat.sock"), None)
.expect("current daemon identity");
let current = identity.to_proto();
let mut encoded = Vec::new();
identity
.encode_probe_identity(&mut encoded)
.expect("encode compatibility identity");
let legacy = LegacyDaemonProcess::decode(encoded.as_slice())
.expect("pre-blake3 broker decodes current identity");
assert_eq!(current.exe_hash_algorithm, EXECUTABLE_HASH_ALGORITHM);
assert_eq!(legacy.exe_sha256.len(), 32);
assert_eq!(
legacy.exe_sha256,
sha256_file(&identity.exe_path)
.expect("sha256 executable")
.to_vec(),
"the legacy wire field remains verifiable by a stable broker"
);
}
#[test]
fn blake3_only_identity_skips_the_legacy_sha256_pass_and_keeps_tag_three_zeroed() {
let identity = DaemonProcess::current_process_with_hash_policy(
endpoint("blake3-only.sock"),
None,
DaemonIdentityHashPolicy::Blake3Only,
)
.expect("current daemon identity");
let mut encoded = Vec::new();
identity
.encode_probe_identity(&mut encoded)
.expect("encode compatibility identity");
let legacy = LegacyDaemonProcess::decode(encoded.as_slice())
.expect("pre-blake3 broker decodes current identity");
assert_eq!(identity.legacy_exe_sha256, [0; 32]);
assert_eq!(legacy.exe_sha256, [0; 32]);
assert_eq!(
identity.exe_hash,
executable_hash_file(&identity.exe_path).expect("blake3 executable")
);
}
#[test]
fn blake3_only_identity_sidecar_records_a_zero_legacy_digest() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("identity.json");
let identity = DaemonProcess::current_process_with_hash_policy(
endpoint("blake3-only-sidecar.sock"),
None,
DaemonIdentityHashPolicy::Blake3Only,
)
.expect("current daemon identity");
crate::broker::backend_sdk::write_daemon_identity_file(&path, &identity)
.expect("write sidecar");
let json: serde_json::Value =
serde_json::from_slice(&std::fs::read(&path).expect("read sidecar"))
.expect("parse sidecar");
let legacy = json
.get("legacy_exe_sha256")
.and_then(serde_json::Value::as_array)
.expect("legacy digest array");
assert_eq!(legacy.len(), 32);
assert!(legacy.iter().all(|byte| byte == &serde_json::json!(0)));
}
fn endpoint(path: &str) -> Endpoint {
Endpoint {
namespace_id: "ns".to_owned(),
path: path.to_owned(),
}
}
fn identity(exe_hash: [u8; 32]) -> DaemonProcess {
DaemonProcess {
pid: 1234,
exe_path: PathBuf::from("runtime/soldr-self/v0.8.44-deadbeef/soldr.exe"),
exe_hash,
legacy_exe_sha256: [0x24; 32],
boot_id: "boot-1".to_owned(),
ipc_endpoint: endpoint("rpb-v2-soldr-daemon-0123456789abcdef-0"),
started_at_unix_ms: 1,
idle_timeout_secs: Some(600),
}
}
#[test]
fn distinct_builds_get_distinct_identities() {
let a = identity([0xAA; 32]);
let mut rebuilt = [0xAA; 32];
rebuilt[0] = 0xBB;
let b = identity(rebuilt);
assert_ne!(
a, b,
"a different executable hash must produce a distinct daemon identity, \
so the broker can never conflate two builds (no stale-version war)"
);
assert_eq!(
a,
identity([0xAA; 32]),
"identical executable bytes must yield the same identity"
);
}
#[test]
fn pre_4_10_4_json_defaults_the_missing_legacy_sha256() {
let original = identity([0x42; 32]);
let mut legacy_json = serde_json::to_value(&original).expect("serialize identity");
let object = legacy_json.as_object_mut().expect("identity JSON object");
assert!(object.remove("legacy_exe_sha256").is_some());
let restored: DaemonProcess =
serde_json::from_value(legacy_json).expect("read pre-4.10.4 identity JSON");
let mut expected = original;
expected.legacy_exe_sha256 = [0; 32];
assert_eq!(restored, expected);
}
#[test]
fn exe_sha256_survives_the_manifest_wire_round_trip() {
let original = identity([0x42; 32]);
let proto = original.to_proto();
assert_eq!(proto.exe_hash_algorithm, "blake3");
assert_eq!(
proto.exe_hash.len(),
32,
"the wire form must carry the full 32-byte BLAKE3 hash"
);
let restored = DaemonProcess::try_from(proto).expect("identity round-trips");
let mut expected = original;
expected.legacy_exe_sha256 = [0; 32];
assert_eq!(
restored, expected,
"the canonical BLAKE3 identity must survive the manifest round-trip; the legacy probe field is not persisted"
);
}
#[test]
fn legacy_sha256_wire_identity_is_rejected_actionably() {
let legacy = protocol::DaemonProcess {
pid: 1234,
exe_path: "legacy-daemon".to_owned(),
exe_hash_algorithm: String::new(),
exe_hash: Vec::new(),
ipc_endpoint: Some(endpoint("legacy.sock")),
started_at_unix_ms: 1,
boot_id: "boot-1".to_owned(),
idle_timeout_secs: None,
};
let error = DaemonProcess::try_from(legacy).expect_err("legacy SHA-256 must be fenced");
assert!(matches!(
error,
IdentityError::UnsupportedExecutableHashAlgorithm(ref algorithm)
if algorithm.is_empty()
));
assert!(error.to_string().contains("expected blake3"));
}
}