1use std::time::{SystemTime, UNIX_EPOCH};
2
3use alloy_network::EthereumWallet;
4use alloy_primitives::{Address, B256, U256, keccak256};
5use alloy_provider::{Provider, ProviderBuilder};
6use alloy_rpc_types_eth::TransactionRequest;
7use alloy_sol_types::{Eip712Domain, SolCall, SolStruct, sol};
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
16sol! {
18 #[allow(missing_docs)]
19 struct Heartbeat {
20 address operator;
21 uint64 serviceId;
22 uint64 blueprintId;
23 uint8 statusCode;
24 bytes32 metricsHash;
25 uint64 timestamp;
26 }
27}
28
29#[derive(Clone, Debug)]
31pub struct HeartbeatPayload {
32 pub block_number: u64,
33 pub timestamp: u64,
34 pub service_id: u64,
35 pub blueprint_id: u64,
36 pub status_code: u32,
37}
38
39impl HeartbeatPayload {
40 pub fn encode(&self) -> Vec<u8> {
42 let mut bytes = Vec::with_capacity(32);
43 bytes.extend_from_slice(&self.block_number.to_be_bytes());
44 bytes.extend_from_slice(&self.timestamp.to_be_bytes());
45 bytes.extend_from_slice(&self.service_id.to_be_bytes());
46 bytes.extend_from_slice(&self.blueprint_id.to_be_bytes());
47 bytes.extend_from_slice(&self.status_code.to_be_bytes());
48 bytes
49 }
50}
51
52pub fn print_status(status: &OperatorStatusSnapshot, json_output: bool) {
54 if json_output {
55 let payload = json!({
56 "service_id": status.service_id,
57 "operator": format!("{:#x}", status.operator),
58 "status_code": status.status_code,
59 "last_heartbeat": status.last_heartbeat,
60 "online": status.online,
61 });
62 println!(
63 "{}",
64 serde_json::to_string_pretty(&payload).expect("serialize operator status to json")
65 );
66 return;
67 }
68
69 println!(
70 "{}: {}",
71 style("Service ID").green().bold(),
72 style(status.service_id).green()
73 );
74 println!("{}: {:#x}", style("Operator").green(), status.operator);
75 println!("{}: {}", style("Status Code").green(), status.status_code);
76 if status.last_heartbeat == 0 {
77 println!("{}: {}", style("Last Heartbeat").green(), "never");
78 } else {
79 println!(
80 "{}: {}",
81 style("Last Heartbeat").green(),
82 status.last_heartbeat
83 );
84 }
85 println!("{}: {}", style("Online").green(), status.online);
86}
87
88pub async fn submit_heartbeat(
90 http_rpc_endpoint: &str,
91 status_registry_address: Address,
92 signing_key: &mut K256SigningKey,
93 service_id: u64,
94 blueprint_id: u64,
95 status_code: u8,
96 json_output: bool,
97) -> Result<()> {
98 let timestamp = SystemTime::now()
99 .duration_since(UNIX_EPOCH)
100 .map_err(|e| eyre!("System time error: {e}"))?
101 .as_secs();
102
103 let payload = HeartbeatPayload {
104 block_number: 0,
105 timestamp,
106 service_id,
107 blueprint_id,
108 status_code: u32::from(status_code),
109 };
110
111 let metrics_bytes = payload.encode();
112
113 let local_signer = signing_key
114 .alloy_key()
115 .map_err(|e| eyre!("Failed to prepare wallet signer: {e}"))?;
116 let operator_address = local_signer.address();
117 let wallet = EthereumWallet::from(local_signer);
118
119 let provider = ProviderBuilder::new()
120 .wallet(wallet)
121 .connect(http_rpc_endpoint)
122 .await
123 .map_err(|e| eyre!("Failed to connect to RPC endpoint: {e}"))?;
124
125 let chain_id = provider
126 .get_chain_id()
127 .await
128 .map_err(|e| eyre!("Failed to query chain id: {e}"))?;
129
130 let signature = sign_heartbeat_eip712(
131 signing_key,
132 chain_id,
133 status_registry_address,
134 operator_address,
135 service_id,
136 blueprint_id,
137 status_code,
138 &metrics_bytes,
139 timestamp,
140 )?;
141
142 let heartbeat_call = submitHeartbeatCall {
143 serviceId: service_id,
144 blueprintId: blueprint_id,
145 statusCode: status_code,
146 metrics: metrics_bytes.into(),
147 timestamp,
148 signature: signature.into(),
149 };
150
151 let calldata = heartbeat_call.abi_encode();
152
153 let tx_request = TransactionRequest::default()
154 .to(status_registry_address)
155 .input(calldata.into());
156
157 let pending_tx = provider
158 .send_transaction(tx_request)
159 .await
160 .map_err(|e| eyre!("Failed to submit heartbeat transaction: {e}"))?;
161
162 let receipt = pending_tx
163 .get_receipt()
164 .await
165 .map_err(|e| eyre!("Failed to finalize heartbeat transaction: {e}"))?;
166
167 if json_output {
168 let output = json!({
169 "service_id": service_id,
170 "blueprint_id": blueprint_id,
171 "status_code": status_code,
172 "timestamp": timestamp,
173 "tx_hash": format!("{:#x}", receipt.transaction_hash),
174 "success": receipt.status(),
175 });
176 println!("{}", serde_json::to_string_pretty(&output)?);
177 } else {
178 if receipt.status() {
179 println!(
180 "{} Heartbeat submitted successfully",
181 style("✓").green().bold()
182 );
183 println!(" Transaction: {:#x}", receipt.transaction_hash);
184 println!(" Service ID: {}", service_id);
185 println!(" Blueprint ID: {}", blueprint_id);
186 println!(" Status Code: {}", status_code);
187 println!(" Timestamp: {}", timestamp);
188 } else {
189 println!("{} Heartbeat transaction reverted", style("✗").red().bold());
190 println!(" Transaction: {:#x}", receipt.transaction_hash);
191 }
192 }
193
194 Ok(())
195}
196
197#[allow(clippy::too_many_arguments)]
200fn sign_heartbeat_eip712(
201 signing_key: &mut K256SigningKey,
202 chain_id: u64,
203 verifying_contract: Address,
204 operator: Address,
205 service_id: u64,
206 blueprint_id: u64,
207 status_code: u8,
208 metrics: &[u8],
209 timestamp: u64,
210) -> Result<Vec<u8>> {
211 let domain = Eip712Domain {
212 name: Some("OperatorStatusRegistry".into()),
213 version: Some("1".into()),
214 chain_id: Some(U256::from(chain_id)),
215 verifying_contract: Some(verifying_contract),
216 salt: None,
217 };
218 let heartbeat = Heartbeat {
219 operator,
220 serviceId: service_id,
221 blueprintId: blueprint_id,
222 statusCode: status_code,
223 metricsHash: keccak256(metrics),
224 timestamp,
225 };
226 let digest: B256 = heartbeat.eip712_signing_hash(&domain);
227
228 let (signature, recovery_id) = signing_key
229 .0
230 .sign_prehash_recoverable(digest.as_slice())
231 .map_err(|e| eyre!("Failed to sign heartbeat payload: {e}"))?;
232
233 let mut signature_bytes = Vec::with_capacity(65);
234 signature_bytes.extend_from_slice(&signature.to_bytes());
235 signature_bytes.push(recovery_id.to_byte() + 27);
236 Ok(signature_bytes)
237}