File size: 9,749 Bytes
93c7565 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 | // prover/broadcaster.js
// CommonJS β reads signed attestations, broadcasts to CircuitBreaker on Amoy
// SafeKrypte signs attestations; oracle private key signs transaction envelopes
const { ethers } = require("ethers");
const fs = require("node:fs");
const path = require("node:path");
const notifier = require("./notifier");
// ---------------------------------------------------------------------------
// Files
// ---------------------------------------------------------------------------
const ATTESTATIONS_FILE = path.resolve(
__dirname,
"..",
".local",
"state",
"submitter-attestations.json"
);
const RECEIPTS_FILE = path.resolve(
__dirname,
"..",
".local",
"state",
"broadcast-receipts.json"
);
// ---------------------------------------------------------------------------
// ABI β exactly matching CircuitBreaker.sol (MVP)
const CIRCUIT_BREAKER_ABI = [
"function updateProof(bytes32 assetId, bytes32 deedHash) external",
"function tripCircuit(string calldata reason) external",
"function circuitOpen() external view returns (bool)",
"function latestProof(bytes32 assetId) external view returns (bytes32)",
];
// ---------------------------------------------------------------------------
// Idempotency: receipt store
// ---------------------------------------------------------------------------
function loadReceipts() {
if (!fs.existsSync(RECEIPTS_FILE)) return {};
try {
return JSON.parse(fs.readFileSync(RECEIPTS_FILE, "utf-8"));
} catch {
console.warn("[broadcaster] Could not parse receipts β starting fresh.");
return {};
}
}
function saveReceipts(receipts) {
const dir = path.dirname(RECEIPTS_FILE);
if (!fs.existsSync(dir)) fs.mkdirSync(dir, { recursive: true });
fs.writeFileSync(RECEIPTS_FILE, JSON.stringify(receipts, null, 2));
}
// ---------------------------------------------------------------------------
// Attestation loader β handles { attestations: [...] } or raw array
// ---------------------------------------------------------------------------
function loadAttestations() {
if (!fs.existsSync(ATTESTATIONS_FILE)) return [];
const raw = JSON.parse(fs.readFileSync(ATTESTATIONS_FILE, "utf-8"));
const attestations = Array.isArray(raw) ? raw : raw.attestations;
if (!Array.isArray(attestations)) {
throw new Error("Attestations file must be an array or { attestations: [...] }");
}
return attestations;
}
// ---------------------------------------------------------------------------
// Validation
// ---------------------------------------------------------------------------
function validateAttestation(att) {
const required = ["digest", "signature", "keyId", "action"];
for (const f of required) {
if (!att[f]) throw new Error(`Missing field: ${f}`);
}
const { kind, assetId, deedHash, reason } = att.action;
if (!["updateProof", "tripCircuit"].includes(kind)) {
throw new Error(`Unknown action kind: ${kind}`);
}
if (kind === "updateProof") {
if (!ethers.isHexString(assetId, 32)) throw new Error("Invalid assetId");
if (!ethers.isHexString(deedHash, 32)) throw new Error("Invalid deedHash");
}
if (kind === "tripCircuit") {
if (!reason || typeof reason !== "string") throw new Error("tripCircuit requires reason");
}
}
// ---------------------------------------------------------------------------
// Single attestation broadcast
// ---------------------------------------------------------------------------
async function broadcastOne(contract, att, receipts, options = {}) {
const digest = att.digest;
// 1. Receipt guard
if (receipts[digest]) {
console.log(`[broadcaster] SKIP ${digest.slice(0, 16)}... β already broadcast (tx: ${receipts[digest].txHash})`);
return { skipped: true, digest, txHash: receipts[digest].txHash };
}
if (options.dryRun) {
console.log(`[broadcaster] DRY-RUN ${digest.slice(0, 16)}... β would send ${att.action.kind}`);
return { dryRun: true, digest };
}
const { kind, assetId, deedHash, reason } = att.action;
let tx;
if (kind === "updateProof") {
// On-chain idempotency: skip if deedHash already set
try {
const onChain = await contract.latestProof(assetId);
if (onChain.toLowerCase() === deedHash.toLowerCase()) {
console.log("[broadcaster] SKIP updateProof β deedHash already on-chain");
receipts[digest] = { txHash: "already-on-chain", timestamp: new Date().toISOString() };
saveReceipts(receipts);
return { skipped: true, digest, reason: "already-on-chain" };
}
} catch (err) {
console.warn(`[broadcaster] Could not read latestProof (continuing): ${err.message}`);
}
console.log(`[broadcaster] SEND updateProof assetId=${assetId.slice(0, 16)}... digest=${digest.slice(0, 16)}...`);
tx = await contract.updateProof(assetId, deedHash); // MVP: no signature
}
if (kind === "tripCircuit") {
// Circuit state guard
try {
const open = await contract.circuitOpen();
if (!open) {
console.log("[broadcaster] SKIP tripCircuit β circuit already tripped");
receipts[digest] = { txHash: "circuit-already-tripped", timestamp: new Date().toISOString() };
saveReceipts(receipts);
return { skipped: true, digest, reason: "circuit-already-tripped" };
}
} catch (err) {
console.warn(`[broadcaster] Could not read circuitOpen (continuing): ${err.message}`);
}
console.log(`[broadcaster] SEND tripCircuit reason="${reason}" digest=${digest.slice(0, 16)}...`);
tx = await contract.tripCircuit(reason); // MVP: no signature
}
console.log(`[broadcaster] WAIT ${tx.hash}...`);
const receipt = await tx.wait();
receipts[digest] = {
txHash: receipt.hash,
blockNumber: receipt.blockNumber,
gasUsed: receipt.gasUsed.toString(),
timestamp: new Date().toISOString(),
attestationDigest: digest,
action: att.action,
keyId: att.keyId,
};
saveReceipts(receipts);
console.log(`[broadcaster] CONFIRMED block ${receipt.blockNumber} gas ${receipt.gasUsed} tx ${receipt.hash}`);
// Notify success
try {
notifier.broadcastSuccess({
txHash: receipt.hash,
block: receipt.blockNumber,
gasUsed: receipt.gasUsed.toString(),
action: att.action.kind,
assetId: att.action.assetId,
});
// Extra alert for circuit trip β high-severity event
if (att.action.kind === 'tripCircuit') {
notifier.circuitTrip({
reason: att.action.reason,
txHash: receipt.hash,
assetId: att.action.assetId,
});
}
} catch (_) {}
return { success: true, digest, txHash: receipt.hash, blockNumber: receipt.blockNumber };
}
// ---------------------------------------------------------------------------
// Main
// ---------------------------------------------------------------------------
async function broadcast(options = {}) {
// Notify pipeline start
try { notifier.broadcastStart({}); } catch (_) {}
const attestations = loadAttestations();
if (!attestations.length) {
console.log("[broadcaster] No attestations β nothing to broadcast.");
return [];
}
console.log(`[broadcaster] Loaded ${attestations.length} attestation(s)`);
// Env check
const rpcUrl = process.env.POLYGON_AMOY_RPC_URL || process.env.RPC_URL;
const oracleKey = process.env.ORACLE_PRIVATE_KEY;
const contractAddr = process.env.CIRCUIT_BREAKER_ADDRESS;
if (!rpcUrl) throw new Error("Missing RPC_URL / POLYGON_AMOY_RPC_URL");
if (!oracleKey) throw new Error("Missing ORACLE_PRIVATE_KEY");
if (!contractAddr) throw new Error("Missing CIRCUIT_BREAKER_ADDRESS");
if (!ethers.isAddress(contractAddr)) throw new Error("Invalid CIRCUIT_BREAKER_ADDRESS");
const provider = new ethers.JsonRpcProvider(rpcUrl);
const signer = new ethers.Wallet(oracleKey, provider);
const contract = new ethers.Contract(contractAddr, CIRCUIT_BREAKER_ABI, signer);
// For V2, no oracle check needed; threshold signatures handle auth
console.log(`[broadcaster] Connected to CircuitBreakerV2 at ${contractAddr}`);
const receipts = loadReceipts();
const results = [];
for (const att of attestations) {
try {
validateAttestation(att);
const r = await broadcastOne(contract, att, receipts, options);
results.push(r);
} catch (err) {
console.error(`[broadcaster] FAILED ${att.digest?.slice(0, 16) || "unknown"}: ${err.message}`);
// Notify failure
try {
notifier.broadcastFailure({
digest: att.digest,
action: att.action?.kind,
assetId: att.action?.assetId,
error: err.message,
});
} catch (_) {}
results.push({ error: true, digest: att.digest, message: err.message });
}
}
const succeeded = results.filter(r => r.success).length;
const skipped = results.filter(r => r.skipped).length;
const dryRuns = results.filter(r => r.dryRun).length;
const failed = results.filter(r => r.error).length;
console.log(`[broadcaster] Done β ${succeeded} sent, ${skipped} skipped, ${dryRuns} dry-run, ${failed} failed`);
return results;
}
module.exports = {
broadcast,
broadcastOne,
loadAttestations,
loadReceipts,
saveReceipts,
validateAttestation,
};
// ---------------------------------------------------------------------------
// CLI
// ---------------------------------------------------------------------------
if (require.main === module) {
const dryRun = process.argv.includes("--dry-run") || process.argv.includes("--dry");
broadcast({ dryRun })
.then(() => {
process.exit(0);
})
.catch(err => {
console.error(`[broadcaster] Fatal: ${err.message}`);
process.exit(1);
});
}
|