use std::time::{SystemTime, UNIX_EPOCH};
use alloy_network::EthereumWallet;
use alloy_primitives::{Address, B256, U256, keccak256};
use alloy_provider::{Provider, ProviderBuilder};
use alloy_rpc_types_eth::TransactionRequest;
use alloy_sol_types::{Eip712Domain, SolCall, SolStruct, sol};
use blueprint_client_tangle::{IOperatorStatusRegistry, OperatorStatusSnapshot};
use blueprint_crypto::k256::K256SigningKey;
use color_eyre::eyre::{Result, eyre};
use dialoguer::console::style;
use serde_json::json;
use IOperatorStatusRegistry::submitHeartbeatCall;
sol! {
#[allow(missing_docs)]
struct Heartbeat {
address operator;
uint64 serviceId;
uint64 blueprintId;
uint8 statusCode;
bytes32 metricsHash;
uint64 timestamp;
}
}
#[derive(Clone, Debug)]
pub struct HeartbeatPayload {
pub block_number: u64,
pub timestamp: u64,
pub service_id: u64,
pub blueprint_id: u64,
pub status_code: u32,
}
impl HeartbeatPayload {
pub fn encode(&self) -> Vec<u8> {
let mut bytes = Vec::with_capacity(32);
bytes.extend_from_slice(&self.block_number.to_be_bytes());
bytes.extend_from_slice(&self.timestamp.to_be_bytes());
bytes.extend_from_slice(&self.service_id.to_be_bytes());
bytes.extend_from_slice(&self.blueprint_id.to_be_bytes());
bytes.extend_from_slice(&self.status_code.to_be_bytes());
bytes
}
}
pub fn print_status(status: &OperatorStatusSnapshot, json_output: bool) {
if json_output {
let payload = json!({
"service_id": status.service_id,
"operator": format!("{:#x}", status.operator),
"status_code": status.status_code,
"last_heartbeat": status.last_heartbeat,
"online": status.online,
});
println!(
"{}",
serde_json::to_string_pretty(&payload).expect("serialize operator status to json")
);
return;
}
println!(
"{}: {}",
style("Service ID").green().bold(),
style(status.service_id).green()
);
println!("{}: {:#x}", style("Operator").green(), status.operator);
println!("{}: {}", style("Status Code").green(), status.status_code);
if status.last_heartbeat == 0 {
println!("{}: {}", style("Last Heartbeat").green(), "never");
} else {
println!(
"{}: {}",
style("Last Heartbeat").green(),
status.last_heartbeat
);
}
println!("{}: {}", style("Online").green(), status.online);
}
pub async fn submit_heartbeat(
http_rpc_endpoint: &str,
status_registry_address: Address,
signing_key: &mut K256SigningKey,
service_id: u64,
blueprint_id: u64,
status_code: u8,
json_output: bool,
) -> Result<()> {
let timestamp = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_err(|e| eyre!("System time error: {e}"))?
.as_secs();
let payload = HeartbeatPayload {
block_number: 0,
timestamp,
service_id,
blueprint_id,
status_code: u32::from(status_code),
};
let metrics_bytes = payload.encode();
let local_signer = signing_key
.alloy_key()
.map_err(|e| eyre!("Failed to prepare wallet signer: {e}"))?;
let operator_address = local_signer.address();
let wallet = EthereumWallet::from(local_signer);
let provider = ProviderBuilder::new()
.wallet(wallet)
.connect(http_rpc_endpoint)
.await
.map_err(|e| eyre!("Failed to connect to RPC endpoint: {e}"))?;
let chain_id = provider
.get_chain_id()
.await
.map_err(|e| eyre!("Failed to query chain id: {e}"))?;
let signature = sign_heartbeat_eip712(
signing_key,
chain_id,
status_registry_address,
operator_address,
service_id,
blueprint_id,
status_code,
&metrics_bytes,
timestamp,
)?;
let heartbeat_call = submitHeartbeatCall {
serviceId: service_id,
blueprintId: blueprint_id,
statusCode: status_code,
metrics: metrics_bytes.into(),
timestamp,
signature: signature.into(),
};
let calldata = heartbeat_call.abi_encode();
let tx_request = TransactionRequest::default()
.to(status_registry_address)
.input(calldata.into());
let pending_tx = provider
.send_transaction(tx_request)
.await
.map_err(|e| eyre!("Failed to submit heartbeat transaction: {e}"))?;
let receipt = pending_tx
.get_receipt()
.await
.map_err(|e| eyre!("Failed to finalize heartbeat transaction: {e}"))?;
if json_output {
let output = json!({
"service_id": service_id,
"blueprint_id": blueprint_id,
"status_code": status_code,
"timestamp": timestamp,
"tx_hash": format!("{:#x}", receipt.transaction_hash),
"success": receipt.status(),
});
println!("{}", serde_json::to_string_pretty(&output)?);
} else {
if receipt.status() {
println!(
"{} Heartbeat submitted successfully",
style("✓").green().bold()
);
println!(" Transaction: {:#x}", receipt.transaction_hash);
println!(" Service ID: {}", service_id);
println!(" Blueprint ID: {}", blueprint_id);
println!(" Status Code: {}", status_code);
println!(" Timestamp: {}", timestamp);
} else {
println!("{} Heartbeat transaction reverted", style("✗").red().bold());
println!(" Transaction: {:#x}", receipt.transaction_hash);
}
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
fn sign_heartbeat_eip712(
signing_key: &mut K256SigningKey,
chain_id: u64,
verifying_contract: Address,
operator: Address,
service_id: u64,
blueprint_id: u64,
status_code: u8,
metrics: &[u8],
timestamp: u64,
) -> Result<Vec<u8>> {
let domain = Eip712Domain {
name: Some("OperatorStatusRegistry".into()),
version: Some("1".into()),
chain_id: Some(U256::from(chain_id)),
verifying_contract: Some(verifying_contract),
salt: None,
};
let heartbeat = Heartbeat {
operator,
serviceId: service_id,
blueprintId: blueprint_id,
statusCode: status_code,
metricsHash: keccak256(metrics),
timestamp,
};
let digest: B256 = heartbeat.eip712_signing_hash(&domain);
let (signature, recovery_id) = signing_key
.0
.sign_prehash_recoverable(digest.as_slice())
.map_err(|e| eyre!("Failed to sign heartbeat payload: {e}"))?;
let mut signature_bytes = Vec::with_capacity(65);
signature_bytes.extend_from_slice(&signature.to_bytes());
signature_bytes.push(recovery_id.to_byte() + 27);
Ok(signature_bytes)
}