Skip to main content

cargo_tangle/command/
operator.rs

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
16// EIP-712 typed struct mirroring `OperatorStatusRegistry.HEARTBEAT_TYPEHASH`.
17sol! {
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/// Heartbeat status payload
30#[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    /// Encode the payload to bytes (simple big-endian serialization)
41    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
52/// Print operator status in human or JSON form.
53pub 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
88/// Submit a heartbeat to the OperatorStatusRegistry contract.
89pub 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/// Build an EIP-712 typed-data signature for `OperatorStatusRegistry.submitHeartbeat`.
198/// See `crates/qos/src/heartbeat.rs::sign_heartbeat_eip712` for the canonical impl.
199#[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}