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);
    });
}