cargo_tangle/command/
operator.rs1use std::time::{SystemTime, UNIX_EPOCH};
2
3use alloy_network::EthereumWallet;
4use alloy_primitives::{Address, keccak256};
5use alloy_provider::{Provider, ProviderBuilder};
6use alloy_rpc_types_eth::TransactionRequest;
7use alloy_sol_types::SolCall;
8use blueprint_client_tangle::{IOperatorStatusRegistry, OperatorStatusSnapshot};
9use blueprint_crypto::k256::K256SigningKey;
10use color_eyre::eyre::{Result, eyre};
11use dialoguer::console::style;
12use serde_json::json;
13
14use IOperatorStatusRegistry::submitHeartbeatCall;
15
16const ETH_MESSAGE_PREFIX: &[u8] = b"\x19Ethereum Signed Message:\n32";
17
18#[derive(Clone, Debug)]
20pub struct HeartbeatPayload {
21 pub block_number: u64,
22 pub timestamp: u64,
23 pub service_id: u64,
24 pub blueprint_id: u64,
25 pub status_code: u32,
26}
27
28impl HeartbeatPayload {
29 pub fn encode(&self) -> Vec<u8> {
31 let mut bytes = Vec::with_capacity(32);
32 bytes.extend_from_slice(&self.block_number.to_be_bytes());
33 bytes.extend_from_slice(&self.timestamp.to_be_bytes());
34 bytes.extend_from_slice(&self.service_id.to_be_bytes());
35 bytes.extend_from_slice(&self.blueprint_id.to_be_bytes());
36 bytes.extend_from_slice(&self.status_code.to_be_bytes());
37 bytes
38 }
39}
40
41pub fn print_status(status: &OperatorStatusSnapshot, json_output: bool) {
43 if json_output {
44 let payload = json!({
45 "service_id": status.service_id,
46 "operator": format!("{:#x}", status.operator),
47 "status_code": status.status_code,
48 "last_heartbeat": status.last_heartbeat,
49 "online": status.online,
50 });
51 println!(
52 "{}",
53 serde_json::to_string_pretty(&payload).expect("serialize operator status to json")
54 );
55 return;
56 }
57
58 println!(
59 "{}: {}",
60 style("Service ID").green().bold(),
61 style(status.service_id).green()
62 );
63 println!("{}: {:#x}", style("Operator").green(), status.operator);
64 println!("{}: {}", style("Status Code").green(), status.status_code);
65 if status.last_heartbeat == 0 {
66 println!("{}: {}", style("Last Heartbeat").green(), "never");
67 } else {
68 println!(
69 "{}: {}",
70 style("Last Heartbeat").green(),
71 status.last_heartbeat
72 );
73 }
74 println!("{}: {}", style("Online").green(), status.online);
75}
76
77pub async fn submit_heartbeat(
79 http_rpc_endpoint: &str,
80 status_registry_address: Address,
81 signing_key: &mut K256SigningKey,
82 service_id: u64,
83 blueprint_id: u64,
84 status_code: u8,
85 json_output: bool,
86) -> Result<()> {
87 let timestamp = SystemTime::now()
88 .duration_since(UNIX_EPOCH)
89 .map_err(|e| eyre!("System time error: {e}"))?
90 .as_secs();
91
92 let payload = HeartbeatPayload {
93 block_number: 0,
94 timestamp,
95 service_id,
96 blueprint_id,
97 status_code: u32::from(status_code),
98 };
99
100 let metrics_bytes = payload.encode();
101 let signature = sign_heartbeat_payload(
102 signing_key,
103 service_id,
104 blueprint_id,
105 status_code,
106 &metrics_bytes,
107 )?;
108
109 let local_signer = signing_key
110 .alloy_key()
111 .map_err(|e| eyre!("Failed to prepare wallet signer: {e}"))?;
112 let wallet = EthereumWallet::from(local_signer);
113
114 let provider = ProviderBuilder::new()
115 .wallet(wallet)
116 .connect(http_rpc_endpoint)
117 .await
118 .map_err(|e| eyre!("Failed to connect to RPC endpoint: {e}"))?;
119
120 let heartbeat_call = submitHeartbeatCall {
121 serviceId: service_id,
122 blueprintId: blueprint_id,
123 statusCode: status_code,
124 metrics: metrics_bytes.into(),
125 signature: signature.into(),
126 };
127
128 let calldata = heartbeat_call.abi_encode();
129
130 let tx_request = TransactionRequest::default()
131 .to(status_registry_address)
132 .input(calldata.into());
133
134 let pending_tx = provider
135 .send_transaction(tx_request)
136 .await
137 .map_err(|e| eyre!("Failed to submit heartbeat transaction: {e}"))?;
138
139 let receipt = pending_tx
140 .get_receipt()
141 .await
142 .map_err(|e| eyre!("Failed to finalize heartbeat transaction: {e}"))?;
143
144 if json_output {
145 let output = json!({
146 "service_id": service_id,
147 "blueprint_id": blueprint_id,
148 "status_code": status_code,
149 "timestamp": timestamp,
150 "tx_hash": format!("{:#x}", receipt.transaction_hash),
151 "success": receipt.status(),
152 });
153 println!("{}", serde_json::to_string_pretty(&output)?);
154 } else {
155 if receipt.status() {
156 println!(
157 "{} Heartbeat submitted successfully",
158 style("✓").green().bold()
159 );
160 println!(" Transaction: {:#x}", receipt.transaction_hash);
161 println!(" Service ID: {}", service_id);
162 println!(" Blueprint ID: {}", blueprint_id);
163 println!(" Status Code: {}", status_code);
164 println!(" Timestamp: {}", timestamp);
165 } else {
166 println!("{} Heartbeat transaction reverted", style("✗").red().bold());
167 println!(" Transaction: {:#x}", receipt.transaction_hash);
168 }
169 }
170
171 Ok(())
172}
173
174fn sign_heartbeat_payload(
177 signing_key: &mut K256SigningKey,
178 service_id: u64,
179 blueprint_id: u64,
180 status_code: u8,
181 metrics: &[u8],
182) -> Result<Vec<u8>> {
183 let mut payload = Vec::with_capacity(17 + metrics.len());
184 payload.extend_from_slice(&service_id.to_be_bytes());
185 payload.extend_from_slice(&blueprint_id.to_be_bytes());
186 payload.push(status_code);
187 payload.extend_from_slice(metrics);
188
189 let message_hash = keccak256(&payload);
190
191 let mut prefixed = Vec::with_capacity(ETH_MESSAGE_PREFIX.len() + message_hash.len());
192 prefixed.extend_from_slice(ETH_MESSAGE_PREFIX);
193 prefixed.extend_from_slice(message_hash.as_slice());
194
195 let prefixed_hash = keccak256(&prefixed);
196 let mut digest = [0u8; 32];
197 digest.copy_from_slice(prefixed_hash.as_slice());
198
199 let (signature, recovery_id) = signing_key
200 .0
201 .sign_prehash_recoverable(&digest)
202 .map_err(|e| eyre!("Failed to sign heartbeat payload: {e}"))?;
203
204 let mut signature_bytes = Vec::with_capacity(65);
205 signature_bytes.extend_from_slice(&signature.to_bytes());
206 signature_bytes.push(recovery_id.to_byte() + 27);
207 Ok(signature_bytes)
208}