Skip to main content

cargo_tangle/command/
operator.rs

1use 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/// Heartbeat status payload
19#[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    /// Encode the payload to bytes (simple big-endian serialization)
30    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
41/// Print operator status in human or JSON form.
42pub 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
77/// Submit a heartbeat to the OperatorStatusRegistry contract.
78pub 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
174/// Sign heartbeat payload: keccak256(abi.encodePacked(serviceId, blueprintId, statusCode, metrics))
175/// with Ethereum signed message prefix. Must match OperatorStatusRegistry.sol verification.
176fn 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}