1use std::path::{Path, PathBuf};
13use std::time::Duration;
14
15use bitcoin::{OutPoint, Txid};
16use serde::{Deserialize, Serialize};
17use sidestr_core::block::{HeaderFamily, SidestrBlock};
18use sidestr_core::blockfile::HEADER;
19use sidestr_core::chain::ChainOf;
20use sidestr_core::document::ChainDocument;
21use sidestr_core::marker::{parse_claims, parse_peg_marker};
22use sidestr_core::parent::rpc::CoreRpc;
23use sidestr_core::parent::{
24 claimable_by_transaction, new_pegins, outpoints_to_lock, paid_pegouts_in, parent_network,
25 scan_pegins_with_wallet, FoundPegin, ParentRpc, PegWallet,
26};
27use sidestr_nostr::event::Event;
28use sidestr_nostr::kinds::{
29 KIND_BLOCK_PROPOSAL, KIND_PARTIAL_SIGNATURE, KIND_PEGOUT_PSBT, KIND_PEGOUT_SIGNED,
30 KIND_SEALED_BLOCK, KIND_TRANSACTION,
31};
32use sidestr_nostr::relay::Follower;
33use sidestr_nostr::tip::{sign_tip_with_peg, TipTemplate, TIP_HEADERS};
34use tokio::sync::{mpsc, oneshot};
35
36use crate::error::{Error, Result};
37use crate::journal::{FileJournal, VoteJournal};
38use crate::pegout::{
39 burn_key, PaidPegout, PegCoin, PegoutAction, PegoutConfig, PegoutLedger, PegoutRound,
40};
41use crate::relay::{follow, ok_count, publish_all, unix_now, unix_now_ms};
42use crate::round::{Action, ClaimChecker, Round, RoundConfig};
43use crate::signer::LocalKey;
44
45#[derive(Debug, Clone)]
47pub struct Settings {
48 pub chain: PathBuf,
50 pub dir: PathBuf,
53 pub key_file: PathBuf,
55 pub port: u16,
58 pub interval: u64,
60 pub tx_interval: u64,
62 pub round: RoundConfig,
64 pub pegout: PegoutConfig,
66 pub relays: Vec<String>,
68 pub mirrors: Vec<String>,
70 pub parent: Option<ParentSettings>,
72 pub journal: Option<PathBuf>,
76 pub chain_event: Option<PathBuf>,
81}
82
83pub fn read_chain_hash(
92 chain_file: &Path,
93 explicit: Option<&Path>,
94 alias: &str,
95) -> Result<(Option<String>, PathBuf)> {
96 let path = explicit.map(Path::to_path_buf).unwrap_or_else(|| {
97 chain_file
98 .parent()
99 .unwrap_or_else(|| Path::new(""))
100 .join(sidestr_nostr::chain::CHAIN_EVENT_FILE)
101 });
102 if !path.exists() {
103 return Ok((None, path));
104 }
105 let text = std::fs::read_to_string(&path)?;
106 let v: serde_json::Value = serde_json::from_str(&text)
107 .map_err(|e| Error::Federation(format!("{}: {e}", path.display())))?;
108 let d = sidestr_nostr::chain::parse_chain_event_value(&v)?;
109 if d.alias != alias {
110 return Err(Error::Federation(format!(
111 "{} is the event of {}, not {alias}",
112 path.display(),
113 d.alias
114 )));
115 }
116 Ok((Some(d.hash), path))
117}
118
119pub fn pegout_journal_path_for(journal: &Path) -> PathBuf {
122 let stem = journal
123 .file_stem()
124 .map(|s| s.to_string_lossy().into_owned())
125 .unwrap_or_else(|| "votes".into());
126 let ext = journal
127 .extension()
128 .map(|e| format!(".{}", e.to_string_lossy()))
129 .unwrap_or_default();
130 journal.with_file_name(format!("{stem}-pegout{ext}"))
131}
132
133fn split_combined_journal(journal: &Path, pegout_journal: &Path) -> Result<()> {
140 if pegout_journal.exists() || !journal.exists() {
141 return Ok(());
142 }
143 let combined = FileJournal::open(journal)?;
144 let burns: Vec<_> = combined
145 .entries()?
146 .into_iter()
147 .filter(|e| matches!(e.scope, crate::journal::VoteScope::Burn(_)))
148 .collect();
149 drop(combined);
150 if burns.is_empty() {
151 return Ok(());
152 }
153 let mut split = FileJournal::open(pegout_journal)?;
154 for e in &burns {
155 split.record(e)?;
156 }
157 Ok(())
158}
159
160#[derive(Debug, Clone)]
162pub struct ParentSettings {
163 pub url: String,
165 pub cookie: PathBuf,
167 pub wallet: Option<String>,
169 pub poll: u64,
171 pub from: u32,
173}
174
175#[derive(Debug, Clone, Default, Serialize, Deserialize)]
177struct PeginState {
178 scanned: i64,
179 pegins: Vec<PeginRecord>,
180}
181
182#[derive(Debug, Clone, Serialize, Deserialize)]
183#[serde(rename_all = "camelCase")]
184struct PeginRecord {
185 txid: String,
186 vout: u32,
187 amount: u64,
188 script: String,
189 height: u32,
190 #[serde(default, skip_serializing_if = "Option::is_none")]
191 parent_address: Option<String>,
192}
193
194impl From<&FoundPegin> for PeginRecord {
195 fn from(p: &FoundPegin) -> Self {
196 Self {
197 txid: p.txid.clone(),
198 vout: p.vout,
199 amount: p.amount,
200 script: p.script.to_hex_string(),
201 height: p.height,
202 parent_address: p.parent_address.clone(),
203 }
204 }
205}
206
207impl PeginRecord {
208 fn found(&self) -> FoundPegin {
209 FoundPegin {
210 txid: self.txid.clone(),
211 vout: self.vout,
212 amount: self.amount,
213 script: bitcoin::ScriptBuf::from_bytes(hex::decode(&self.script).unwrap_or_default()),
214 height: self.height,
215 parent_address: self.parent_address.clone(),
216 }
217 }
218}
219
220struct ParentClaims {
224 rpc: std::rc::Rc<CoreRpc>,
225 chain_id: String,
226 need: u32,
227 challenge: bitcoin::ScriptBuf,
228}
229
230impl<F: HeaderFamily> ClaimChecker<F> for ParentClaims {
231 fn check(&self, block: &F::Block) -> Option<String> {
232 let coinbase = block.txdata().first()?;
233 let (claims, errors) = parse_claims(coinbase);
234 if let Some(e) = errors.first() {
235 return Some(e.clone());
236 }
237 for c in claims {
238 let short = &c.txid[..c.txid.len().min(12)];
239 let txid: Txid = c.txid.parse().ok()?;
240 let o = match self.rpc.tx_out(&txid, c.vout) {
241 Ok(Some(o)) => o,
242 Ok(None) => {
243 return Some(format!("{short}…:{} is not unspent on the parent", c.vout))
244 }
245 Err(e) => return Some(e.to_string()),
246 };
247 if o.confirmations < self.need {
248 return Some(format!(
249 "{short}… has {} of {} confirmations",
250 o.confirmations, self.need
251 ));
252 }
253 if o.value != c.payout.value || o.script_pubkey != self.challenge {
254 return Some(format!(
255 "{short}… does not pay the peg {} sats",
256 c.payout.value
257 ));
258 }
259 let raw = match self
260 .rpc
261 .call("getrawtransaction", serde_json::json!([c.txid, true]))
262 {
263 Ok(v) => v,
264 Err(e) => return Some(e.to_string()),
265 };
266 let marker = raw["vout"]
267 .as_array()
268 .into_iter()
269 .flatten()
270 .filter_map(|v| v["scriptPubKey"]["hex"].as_str())
271 .filter_map(|h| hex::decode(h).ok())
272 .find_map(|b| parse_peg_marker(bitcoin::Script::from_bytes(&b), &self.chain_id));
273 if marker.as_deref() != Some(c.payout.script_pubkey.as_script()) {
274 return Some(format!(
275 "{short}…'s marker names {}…, the claim pays {}…",
276 marker
277 .map(|m| m.to_hex_string())
278 .unwrap_or_else(|| "nothing".into())
279 .chars()
280 .take(12)
281 .collect::<String>(),
282 &c.payout.script_pubkey.to_hex_string()[..12]
283 ));
284 }
285 }
286 None
287 }
288}
289
290pub const MAX_TX_BODY: usize = 262_144;
292
293struct DatReply {
296 code: u16,
297 body: Vec<u8>,
298 content_range: Option<String>,
300}
301
302enum Query {
304 Status(oneshot::Sender<serde_json::Value>),
305 Dat(Option<String>, oneshot::Sender<DatReply>),
308 Tip(oneshot::Sender<serde_json::Value>),
309 Chain(oneshot::Sender<serde_json::Value>),
310 Blocks(oneshot::Sender<serde_json::Value>),
311 Pegouts(oneshot::Sender<serde_json::Value>),
312 Coins(String, oneshot::Sender<serde_json::Value>),
313 Tx(
314 String,
315 oneshot::Sender<core::result::Result<serde_json::Value, String>>,
316 ),
317}
318
319fn log(s: impl AsRef<str>) {
320 let now = unix_now();
321 let (h, m, sec) = ((now / 3600) % 24, (now / 60) % 60, now % 60);
322 println!("{h:02}:{m:02}:{sec:02} {}", s.as_ref());
323}
324
325fn read_json<T: for<'a> Deserialize<'a> + Default>(path: &Path) -> T {
326 std::fs::read_to_string(path)
327 .ok()
328 .and_then(|t| serde_json::from_str(&t).ok())
329 .unwrap_or_default()
330}
331
332fn write_json<T: Serialize>(path: &Path, v: &T) {
333 if let Ok(text) = serde_json::to_string_pretty(v) {
334 if let Err(e) = std::fs::write(path, text) {
335 log(format!("{}: {e}", path.display()));
336 }
337 }
338}
339
340struct Node<F: HeaderFamily> {
342 doc: ChainDocument,
343 chain: ChainOf<F>,
344 round: Round<F>,
345 pegout: Option<PegoutRound>,
346 parent: Option<std::rc::Rc<CoreRpc>>,
347 pegins: PeginState,
348 settings: Settings,
349 txs: Follower,
350 key: LocalKey,
351 last_block: u64,
352 announced: Option<u32>,
353 announce_retry_at: u64,
354 started: u64,
355 chain_hash: Option<String>,
357}
358
359impl<F: HeaderFamily> Node<F> {
360 fn status(&self) -> serde_json::Value {
361 let s = self.chain.state();
362 let tip = s.tip();
363 let fed = self.round.federation();
364 serde_json::json!({
365 "chain": self.doc.id, "parent": self.doc.parent,
366 "height": tip.height, "hash": tip.hash.to_string(), "time": tip.time,
367 "coins": s.utxo().len(), "mempool": s.mempool().count(), "minFeeRate": s.min_fee_rate(),
368 "relays": self.settings.relays,
369 "announce": if self.settings.mirrors.is_empty() { serde_json::Value::Null } else { serde_json::json!({"mirrors": self.settings.mirrors, "announced": self.announced.map(i64::from).unwrap_or(-1)}) },
370 "pegouts": {"burned": s.pegouts().len(), "paid": self.pegout.as_ref().map(|p| p.ledger().paid.len()).unwrap_or(0), "min": s.pegout_min(), "payer": self.settings.parent.as_ref().and_then(|p| p.wallet.clone())},
371 "pegins": self.parent.as_ref().map(|_| serde_json::json!({"scanned": self.pegins.scanned, "known": self.pegins.pegins.len(), "claimed": self.pegins.pegins.iter().filter(|p| s.claimed_tx(&p.txid)).count()})),
372 "signer": self.key.pubkey_hex(), "genesis": s.genesis_hash().to_string(), "interval": self.settings.interval,
373 "level2": {"signers": fed.signers.len(), "threshold": fed.threshold, "slot": self.round.slot() + 1,
374 "proposeAfter": self.settings.round.propose_after, "resignAfter": self.settings.round.resign_after,
375 "pending": self.round.pending().map(|p| serde_json::json!({"height": p.height, "signatures": p.sigs.len(), "at": p.at})),
376 "journal": self.settings.journal.as_ref().map(|j| j.display().to_string())},
377 "engine": "sidestr-round", "started": self.started,
378 })
379 }
380
381 fn tip_json(&self) -> serde_json::Value {
382 let t = self.chain.state().tip();
383 serde_json::json!({"height": t.height, "hash": t.hash.to_string(), "time": t.time})
384 }
385
386 fn committed_dat(&self, range: Option<&str>) -> DatReply {
392 let end = self
393 .chain
394 .index()
395 .blocks
396 .last()
397 .map(|e| e.offset + HEADER + u64::from(e.size))
398 .unwrap_or(0);
399 let bytes = match std::fs::read(self.chain.dat_path()) {
400 Ok(b) => b,
401 Err(e) => {
402 return DatReply {
403 code: 500,
404 body: serde_json::json!({"error": e.to_string()})
405 .to_string()
406 .into_bytes(),
407 content_range: None,
408 }
409 }
410 };
411 if (bytes.len() as u64) < end {
412 return DatReply {
413 code: 500,
414 body: serde_json::json!({"error": "the block file is shorter than its index"})
415 .to_string()
416 .into_bytes(),
417 content_range: None,
418 };
419 }
420 let committed = &bytes[..end as usize];
421 match range {
422 None => DatReply {
423 code: 200,
424 body: committed.to_vec(),
425 content_range: None,
426 },
427 Some(h) => match parse_range(h, end) {
428 Some((s, e)) => DatReply {
429 code: 206,
430 body: committed[s as usize..=e as usize].to_vec(),
431 content_range: Some(format!("bytes {s}-{e}/{end}")),
432 },
433 None => DatReply {
434 code: 416,
435 body: Vec::new(),
436 content_range: Some(format!("bytes */{end}")),
437 },
438 },
439 }
440 }
441
442 fn answer(&mut self, q: Query) {
443 match q {
444 Query::Status(r) => {
445 let _ = r.send(self.status());
446 }
447 Query::Tip(r) => {
448 let _ = r.send(self.tip_json());
449 }
450 Query::Dat(range, r) => {
451 let _ = r.send(self.committed_dat(range.as_deref()));
452 }
453 Query::Chain(r) => {
454 let mut d = self.doc.clone();
455 d.genesis_hash = Some(self.chain.state().genesis_hash().to_string());
456 let _ = r.send(serde_json::to_value(&d).unwrap_or_default());
457 }
458 Query::Blocks(r) => {
459 let _ = r.send(serde_json::to_value(self.chain.index()).unwrap_or_default());
460 }
461 Query::Pegouts(r) => {
462 let _ = r.send(
463 self.pegout
464 .as_ref()
465 .map(|p| serde_json::to_value(p.ledger()).unwrap_or_default())
466 .unwrap_or_else(|| serde_json::json!({"paid": {}})),
467 );
468 }
469 Query::Coins(hex_spk, r) => {
470 let list: Vec<serde_json::Value> = hex::decode(&hex_spk)
471 .map(|b| self.chain.state().coins(bitcoin::Script::from_bytes(&b)))
472 .unwrap_or_default()
473 .into_iter()
474 .map(|c| serde_json::json!({"outpoint": format!("{}:{}", c.outpoint.txid, c.outpoint.vout), "value": c.value, "height": c.height, "coinbase": c.coinbase}))
475 .collect();
476 let _ = r.send(serde_json::Value::Array(list));
477 }
478 Query::Tx(hex_tx, r) => {
479 let res = hex::decode(hex_tx.trim())
480 .map_err(|e| e.to_string())
481 .and_then(|b| self.chain.submit(&b).map_err(|e| e.to_string()))
482 .map(|s| {
483 log(format!("tx {}… accepted, fee {}", &s.txid.to_string()[..16], s.fee));
484 serde_json::json!({"txid": s.txid.to_string(), "fee": s.fee, "vsize": s.vsize, "dup": s.dup})
485 });
486 let _ = r.send(res);
487 }
488 }
489 }
490
491 fn headers_hex(&self) -> Vec<String> {
492 let s = self.chain.state();
493 let tip = s.height();
494 let from = tip.saturating_sub(TIP_HEADERS as u32 - 1);
495 (from..=tip)
496 .filter_map(|h| s.header_at(h))
497 .map(|h| hex::encode(s.family().encode_header(h)))
498 .collect()
499 }
500
501 fn peg_coins(&self) -> Vec<PegCoin> {
502 let (Some(rpc), Some(p)) = (&self.parent, &self.settings.parent) else {
503 return vec![];
504 };
505 if p.wallet.is_none() {
506 return vec![];
507 }
508 let challenge = self.chain.state().challenge().to_hex_string();
509 match rpc.wallet_call("listunspent", serde_json::json!([1, 9_999_999, [], true])) {
510 Ok(v) => v
511 .as_array()
512 .into_iter()
513 .flatten()
514 .filter(|u| u["scriptPubKey"].as_str() == Some(challenge.as_str()))
515 .filter_map(|u| {
516 Some(PegCoin {
517 outpoint: OutPoint {
518 txid: u["txid"].as_str()?.parse().ok()?,
519 vout: u32::try_from(u["vout"].as_u64()?).ok()?,
520 },
521 value: (u["amount"].as_f64()? * 1e8).round() as u64,
522 })
523 })
524 .collect(),
525 Err(e) => {
526 log(format!("peg-out round: listunspent: {e}"));
527 vec![]
528 }
529 }
530 }
531
532 fn peg_tick(&mut self) {
536 let Some(rpc) = self.parent.clone() else {
537 return;
538 };
539 let network = self.doc.parent().ok().and_then(parent_network);
540 let tip = match rpc.block_count() {
541 Ok(t) => t,
542 Err(e) => {
543 log(format!("parent: {e}"));
544 return;
545 }
546 };
547 if i64::from(tip) > self.pegins.scanned {
548 let from = u32::try_from(self.pegins.scanned + 1).unwrap_or(0);
549 let announced = self.doc.challenge_script().ok();
550 let wallet: Option<&dyn PegWallet> = rpc
551 .wallet()
552 .is_some()
553 .then_some(rpc.as_ref() as &dyn PegWallet);
554 match scan_pegins_with_wallet(
555 rpc.as_ref(),
556 &self.doc.id,
557 from,
558 tip,
559 network,
560 announced.as_deref(),
561 wallet,
562 |_| {},
563 ) {
564 Ok(found) => {
565 let known: Vec<FoundPegin> =
566 self.pegins.pegins.iter().map(PeginRecord::found).collect();
567 let fresh =
568 new_pegins(&found, &known, |txid| self.chain.state().claimed_tx(txid));
569 for p in fresh {
570 log(format!(
571 "peg-in {}…:{}: {} sats to {}…, parent h{}",
572 &p.txid[..16],
573 p.vout,
574 p.amount,
575 &p.script.to_hex_string()[..12],
576 p.height
577 ));
578 self.pegins.pegins.push((&p).into());
579 }
580 self.pegins.scanned = i64::from(tip);
581 write_json(&self.settings.dir.join("pegins.json"), &self.pegins);
582 }
583 Err(e) => log(format!("peg-in scan: {e}")),
584 }
585 }
586 let found: Vec<FoundPegin> = self.pegins.pegins.iter().map(PeginRecord::found).collect();
587 let s = self.chain.state();
588 let claims =
589 claimable_by_transaction(&found, tip, self.doc.peg_confirmations, |t| s.claimed_tx(t));
590 if self
591 .settings
592 .parent
593 .as_ref()
594 .is_some_and(|p| p.wallet.is_some())
595 {
596 let lock = outpoints_to_lock(&found, |t, _| s.claimed_tx(t));
597 let unlock: Vec<OutPoint> = found
598 .iter()
599 .filter(|p| s.claimed_tx(&p.txid))
600 .filter_map(|p| {
601 Some(OutPoint {
602 txid: p.txid.parse().ok()?,
603 vout: p.vout,
604 })
605 })
606 .collect();
607 if !lock.is_empty() {
608 if let Err(e) = rpc.lock_outputs(&lock, true) {
609 log(format!("lockunspent: {e}"));
610 }
611 }
612 if !unlock.is_empty() {
613 if let Err(e) = rpc.lock_outputs(&unlock, false) {
614 log(format!("lockunspent: {e}"));
615 }
616 }
617 }
618 self.round.want_claims(claims);
620 }
621
622 fn reconcile(&mut self) {
625 let (Some(rpc), Some(round)) = (&self.parent, &mut self.pegout) else {
626 return;
627 };
628 if self
629 .settings
630 .parent
631 .as_ref()
632 .is_none_or(|p| p.wallet.is_none())
633 {
634 return;
635 }
636 let sent = match rpc.sent_transactions() {
637 Ok(s) => s,
638 Err(e) => {
639 log(format!("peg-out reconcile: {e}"));
640 return;
641 }
642 };
643 let paid = paid_pegouts_in(&sent, &self.doc.id);
644 let mut changed = false;
645 for b in self.chain.state().pegouts() {
646 let key = burn_key(&b);
647 if round.ledger().paid.contains_key(&key) {
648 continue;
649 }
650 if let Some(txid) = paid.get(&b.txid) {
651 round.mark_paid(
652 &key,
653 PaidPegout {
654 parent_txid: txid.to_string(),
655 address: None,
656 value: b.value,
657 script: b.script.clone(),
658 height: b.height,
659 at: unix_now(),
660 signers: vec![],
661 reconciled: Some(true),
662 },
663 );
664 log(format!(
665 "peg-out {}… was paid by the federation in {}…",
666 &key[..16],
667 &txid.to_string()[..16]
668 ));
669 changed = true;
670 }
671 }
672 if changed {
673 write_json(&self.settings.dir.join("pegouts.json"), round.ledger());
674 }
675 }
676}
677
678fn parse_range(h: &str, size: u64) -> Option<(u64, u64)> {
679 let r = h.strip_prefix("bytes=")?;
680 let (a, b) = r.split_once('-')?;
681 let start: u64 = a.parse().ok()?;
682 let end: u64 = if b.is_empty() {
683 size.saturating_sub(1)
684 } else {
685 b.parse().ok()?
686 };
687 (start <= end && end < size).then_some((start, end))
688}
689
690fn blocks_etag(index: &serde_json::Value) -> String {
691 let to = index
692 .get("to")
693 .map(serde_json::Value::to_string)
694 .unwrap_or_else(|| "0".into());
695 let hash = index
696 .get("blocks")
697 .and_then(serde_json::Value::as_array)
698 .and_then(|blocks| blocks.last())
699 .and_then(|block| block.get("hash"))
700 .and_then(serde_json::Value::as_str)
701 .unwrap_or("");
702 format!("\"{to}-{}\"", &hash[..hash.len().min(16)])
703}
704
705fn serve_http(port: u16, to_loop: mpsc::UnboundedSender<Query>) -> Result<()> {
706 let server = tiny_http::Server::http(("127.0.0.1", port))
707 .map_err(|e| Error::Io(std::io::Error::other(e.to_string())))?;
708 std::thread::spawn(move || {
709 for mut req in server.incoming_requests() {
710 let path = req.url().split('?').next().unwrap_or("/").to_string();
711 let cors = [
712 ("access-control-allow-origin", "*"),
713 (
714 "access-control-allow-headers",
715 "range, content-type, if-none-match",
716 ),
717 (
718 "access-control-expose-headers",
719 "etag, accept-ranges, content-range",
720 ),
721 ("access-control-allow-methods", "GET, POST, OPTIONS"),
722 ];
723 let with_cors = |mut r: tiny_http::Response<std::io::Cursor<Vec<u8>>>| {
724 for (k, v) in cors {
725 r = r.with_header(tiny_http::Header::from_bytes(k, v).expect("static header"));
726 }
727 r
728 };
729 let json = |code: u16, v: &serde_json::Value| {
730 with_cors(
731 tiny_http::Response::from_string(v.to_string())
732 .with_status_code(code)
733 .with_header(
734 tiny_http::Header::from_bytes("content-type", "application/json")
735 .expect("static header"),
736 ),
737 )
738 };
739 let ask = |q: Query, rx: oneshot::Receiver<serde_json::Value>| {
740 let _ = to_loop.send(q);
741 rx.blocking_recv().unwrap_or(serde_json::Value::Null)
742 };
743 if req.method() == &tiny_http::Method::Options {
744 let _ = req.respond(with_cors(
745 tiny_http::Response::from_data(Vec::new()).with_status_code(204),
746 ));
747 continue;
748 }
749 let response = match (req.method().as_str(), path.as_str()) {
750 ("GET", "/") | ("GET", "/status.json") => {
751 let (tx, rx) = oneshot::channel();
752 json(200, &ask(Query::Status(tx), rx))
753 }
754 ("GET", "/tip") => {
755 let (tx, rx) = oneshot::channel();
756 json(200, &ask(Query::Tip(tx), rx))
757 }
758 ("GET", "/chain.json") => {
759 let (tx, rx) = oneshot::channel();
760 json(200, &ask(Query::Chain(tx), rx))
761 }
762 ("GET", "/blocks.json") => {
763 let (tx, rx) = oneshot::channel();
764 let index = ask(Query::Blocks(tx), rx);
765 let tag = blocks_etag(&index);
766 let unchanged = req
767 .headers()
768 .iter()
769 .find(|h| h.field.equiv("if-none-match"))
770 .is_some_and(|h| h.value.as_str() == tag.as_str());
771 if unchanged {
772 with_cors(tiny_http::Response::from_data(Vec::new()).with_status_code(304))
773 .with_header(
774 tiny_http::Header::from_bytes("etag", tag.as_str())
775 .expect("static header name"),
776 )
777 } else {
778 json(200, &index).with_header(
779 tiny_http::Header::from_bytes("etag", tag.as_str())
780 .expect("static header name"),
781 )
782 }
783 }
784 ("GET", "/pegouts.json") => {
785 let (tx, rx) = oneshot::channel();
786 json(200, &ask(Query::Pegouts(tx), rx))
787 }
788 ("GET", "/blocks.dat") => {
789 let range = req
790 .headers()
791 .iter()
792 .find(|h| h.field.equiv("range"))
793 .map(|h| h.value.as_str().to_string());
794 let (tx, rx) = oneshot::channel();
795 let _ = to_loop.send(Query::Dat(range, tx));
796 match rx.blocking_recv() {
797 Ok(d) if d.code == 500 => {
798 let v: serde_json::Value =
799 serde_json::from_slice(&d.body).unwrap_or_default();
800 json(500, &v)
801 }
802 Ok(d) => {
803 let mut r = with_cors(
804 tiny_http::Response::from_data(d.body).with_status_code(d.code),
805 )
806 .with_header(
807 tiny_http::Header::from_bytes(
808 "content-type",
809 "application/octet-stream",
810 )
811 .expect("static header"),
812 )
813 .with_header(
814 tiny_http::Header::from_bytes("accept-ranges", "bytes")
815 .expect("static header"),
816 );
817 if let Some(cr) = d.content_range {
818 r = r.with_header(
819 tiny_http::Header::from_bytes("content-range", cr.as_str())
820 .expect("static header"),
821 );
822 }
823 r
824 }
825 Err(_) => json(500, &serde_json::json!({"error": "the signer is gone"})),
826 }
827 }
828 ("GET", p) if p.starts_with("/coins/") => {
829 let (tx, rx) = oneshot::channel();
830 json(200, &ask(Query::Coins(p[7..].to_ascii_lowercase(), tx), rx))
831 }
832 ("POST", "/tx") => {
833 let too_large = json(
834 413,
835 &serde_json::json!({"error": format!("the body is over {MAX_TX_BODY} bytes")}),
836 );
837 if req.body_length().is_some_and(|n| n > MAX_TX_BODY) {
838 let _ = req.respond(too_large);
839 continue;
840 }
841 let mut body = String::new();
842 let read = std::io::Read::read_to_string(
843 &mut std::io::Read::take(req.as_reader(), MAX_TX_BODY as u64 + 1),
844 &mut body,
845 );
846 if read.is_err() || body.len() > MAX_TX_BODY {
847 let _ = req.respond(too_large);
848 continue;
849 }
850 let (tx, rx) = oneshot::channel();
851 let _ = to_loop.send(Query::Tx(body, tx));
852 match rx.blocking_recv() {
853 Ok(Ok(v)) => json(200, &v),
854 Ok(Err(e)) => json(400, &serde_json::json!({"error": e})),
855 Err(_) => json(500, &serde_json::json!({"error": "the signer is gone"})),
856 }
857 }
858 _ => json(404, &serde_json::json!({"error": "not found"})),
859 };
860 let _ = req.respond(response);
861 }
862 });
863 Ok(())
864}
865
866pub async fn run_as<F: HeaderFamily>(doc: ChainDocument, settings: Settings) -> Result<()> {
869 let key_text = std::fs::read_to_string(&settings.key_file)
870 .map_err(|e| Error::Key(format!("{}: {e}", settings.key_file.display())))?;
871 let key = LocalKey::from_hex(&key_text)?;
872 let dir = settings.dir.clone();
873 std::fs::create_dir_all(&dir)?;
874 let chain = ChainOf::<F>::open_sealed(doc.clone(), &dir, |_| {
875 Err(sidestr_core::Error::Chain(
876 "no block file: copy blocks.dat and blocks.json from a mirror of this chain first"
877 .into(),
878 ))
879 })?;
880 let journal_path = settings
881 .journal
882 .clone()
883 .unwrap_or_else(|| dir.join("votes.jsonl"));
884 let pegout_journal_path = pegout_journal_path_for(&journal_path);
887 split_combined_journal(&journal_path, &pegout_journal_path)?;
888 let journal = FileJournal::open(&journal_path)?;
889 let loaded = journal.entries()?.len();
890 let mut round = Round::new(
891 chain.state(),
892 Box::new(LocalKey::from_hex(&key_text)?),
893 Box::new(journal),
894 settings.round.clone(),
895 )?;
896 let fed = round.federation().clone();
897 let network = doc.parent().ok().and_then(parent_network);
898 let parent = settings
899 .parent
900 .as_ref()
901 .map(|p| std::rc::Rc::new(CoreRpc::new(&p.url, &p.cookie, p.wallet.as_deref())));
902 if let Some(rpc) = &parent {
903 round = round.with_claim_checker(Box::new(ParentClaims {
904 rpc: rpc.clone(),
905 chain_id: doc.id.clone(),
906 need: doc.peg_confirmations,
907 challenge: chain.state().challenge().to_owned(),
908 }));
909 }
910 let pegout = match (&parent, &settings.parent) {
911 (Some(_), Some(p)) if p.wallet.is_some() => Some(PegoutRound::new(
912 fed.clone(),
913 &doc.id,
914 Box::new(LocalKey::from_hex(&key_text)?),
915 Box::new(FileJournal::open(&pegout_journal_path)?),
916 PegoutConfig {
917 network,
918 ..settings.pegout.clone()
919 },
920 read_json::<PegoutLedger>(&dir.join("pegouts.json")),
921 )?),
922 _ => None,
923 };
924 let (chain_hash, chain_event) =
927 read_chain_hash(&settings.chain, settings.chain_event.as_deref(), &doc.id)?;
928 match &chain_hash {
929 Some(h) => log(format!("chain hash {h} ({})", chain_event.display())),
930 None => log(
931 "no chain-event.json beside the document: tips carry no chain hash (siding chain-event or sidestr-agent chain-event makes one)",
932 ),
933 }
934 let mut pegins: PeginState = read_json(&dir.join("pegins.json"));
935 if pegins.pegins.is_empty() && pegins.scanned == 0 {
936 pegins.scanned = i64::from(settings.parent.as_ref().map(|p| p.from).unwrap_or(0)) - 1;
937 }
938 let now = unix_now();
939 let mut node = Node {
940 txs: Follower::new(KIND_TRANSACTION, &doc.id),
941 doc,
942 chain,
943 round,
944 pegout,
945 parent,
946 pegins,
947 settings: settings.clone(),
948 key,
949 last_block: now,
950 announced: None,
951 announce_retry_at: 0,
952 started: now,
953 chain_hash,
954 };
955 log(format!(
956 "level 2: signer {} of {}, threshold {}, proposing after {} s when it is another signer's turn; journal {} ({loaded} entries)",
957 node.round.slot() + 1,
958 fed.signers.len(),
959 fed.threshold,
960 settings.round.propose_after,
961 journal_path.display()
962 ));
963 log(format!(
964 "chain {} at {} ({} coins), {} relay(s), {} mirror(s), parent {}",
965 node.doc.id,
966 node.chain.state().height(),
967 node.chain.state().utxo().len(),
968 settings.relays.len(),
969 settings.mirrors.len(),
970 settings
971 .parent
972 .as_ref()
973 .map(|p| p.url.as_str())
974 .unwrap_or("none")
975 ));
976 if let Some(p) = node.pegout.as_ref() {
977 log(format!(
978 "parent wallet {}: peg-outs for {} are paid by the federation's PSBT round ({} of {}), {} paid so far",
979 settings.parent.as_ref().and_then(|p| p.wallet.clone()).unwrap_or_default(),
980 node.doc.id,
981 fed.threshold,
982 fed.signers.len(),
983 p.ledger().paid.len()
984 ));
985 }
986
987 let (to_loop, mut queries) = mpsc::unbounded_channel::<Query>();
988 serve_http(settings.port, to_loop)?;
989 log(format!(
990 "producer on http://127.0.0.1:{}/ every {} s ({} s with transactions)",
991 settings.port, settings.interval, settings.tx_interval
992 ));
993 let mut events = follow(
994 settings.relays.clone(),
995 vec![
996 KIND_BLOCK_PROPOSAL,
997 KIND_PARTIAL_SIGNATURE,
998 KIND_SEALED_BLOCK,
999 KIND_PEGOUT_PSBT,
1000 KIND_PEGOUT_SIGNED,
1001 KIND_TRANSACTION,
1002 ],
1003 600,
1004 log,
1005 );
1006 let mut second = tokio::time::interval(Duration::from_secs(1));
1007 let poll = settings
1008 .parent
1009 .as_ref()
1010 .map(|p| p.poll.max(1))
1011 .unwrap_or(60);
1012 let mut parent_tick = tokio::time::interval(Duration::from_secs(poll));
1013 let mut announce_tick = tokio::time::interval(Duration::from_secs(3));
1014 let relays = settings.relays.clone();
1015
1016 let publish = |ev: Event, what: String| {
1017 let relays = relays.clone();
1018 tokio::spawn(async move {
1019 let r = publish_all(&relays, &ev, Duration::from_secs(8)).await;
1020 log(format!("{what} reached {} relay(s)", ok_count(&r)));
1021 });
1022 };
1023 let handle = |node: &mut Node<F>, actions: Vec<Action>| {
1024 for a in actions {
1025 match a {
1026 Action::Log(s) => log(s),
1027 Action::Publish(ev) => {
1028 let what = match ev.kind {
1029 KIND_BLOCK_PROPOSAL => format!(
1030 "round: proposal h{} {}…",
1031 sidestr_nostr::tags::first(&ev.tags, "h").unwrap_or("?"),
1032 &ev.id[..12]
1033 ),
1034 KIND_PARTIAL_SIGNATURE => format!(
1035 "round: partial h{}",
1036 sidestr_nostr::tags::first(&ev.tags, "h").unwrap_or("?")
1037 ),
1038 _ => format!(
1039 "round: sealed h{}",
1040 sidestr_nostr::tags::first(&ev.tags, "h").unwrap_or("?")
1041 ),
1042 };
1043 publish(ev, what);
1044 }
1045 Action::Sealed(_) => {
1046 node.last_block = unix_now();
1047 }
1048 }
1049 }
1050 };
1051 let handle_pegout = |node: &mut Node<F>, actions: Vec<PegoutAction>| {
1052 for a in actions {
1053 match a {
1054 PegoutAction::Log(s) => log(s),
1055 PegoutAction::Publish(ev) => {
1056 let what = format!(
1057 "peg-out round: kind {} for {}…",
1058 ev.kind,
1059 sidestr_nostr::tags::first(&ev.tags, "d")
1060 .map(|d| &d[..d.len().min(16)])
1061 .unwrap_or("?")
1062 );
1063 publish(ev, what);
1064 }
1065 PegoutAction::Broadcast(f) => {
1066 let Some(rpc) = node.parent.clone() else {
1067 continue;
1068 };
1069 let hex = hex::encode(bitcoin::consensus::encode::serialize(&f.tx));
1070 match rpc.call("sendrawtransaction", serde_json::json!([hex])) {
1071 Ok(_) => {
1072 if let Some(p) = node.pegout.as_mut() {
1073 p.mark_paid(&f.burn, f.record);
1074 write_json(&node.settings.dir.join("pegouts.json"), p.ledger());
1075 }
1076 }
1077 Err(e) => log(format!("peg-out round: sendrawtransaction: {e}")),
1078 }
1079 }
1080 }
1081 }
1082 };
1083
1084 loop {
1085 tokio::select! {
1086 _ = second.tick() => {
1087 let now = unix_now();
1088 let wait = if node.chain.state().mempool().count() > 0 { settings.tx_interval } else { settings.interval };
1089 let due = now.saturating_sub(node.last_block) >= wait;
1090 let actions = node.round.tick(unix_now_ms(), &mut node.chain, due);
1091 handle(&mut node, actions);
1092 while let Ok(q) = queries.try_recv() { node.answer(q); }
1093 }
1094 Some(q) = queries.recv() => { node.answer(q); }
1095 Some((_, ev)) = events.recv() => {
1096 let now = unix_now_ms();
1097 match ev.kind {
1098 KIND_TRANSACTION => {
1099 if node.txs.accept(&ev).is_some() {
1100 if let Ok(t) = sidestr_nostr::tx::parse_transaction(&ev, Some(&node.doc.id)) {
1101 match hex::decode(&t.tx_hex).map_err(|e| e.to_string()).and_then(|b| node.chain.submit(&b).map_err(|e| e.to_string())) {
1102 Ok(s) => if !s.dup { log(format!("tx {}… accepted from a relay, fee {}", &s.txid.to_string()[..16], s.fee)) },
1103 Err(e) => log(format!("tx from a relay refused: {e}")),
1104 }
1105 }
1106 }
1107 }
1108 KIND_PEGOUT_PSBT | KIND_PEGOUT_SIGNED => {
1109 let burns = node.chain.state().pegouts();
1110 if let Some(p) = node.pegout.as_mut() {
1111 let actions = p.on_event(now, &ev, &burns);
1112 handle_pegout(&mut node, actions);
1113 }
1114 }
1115 _ => {
1116 let actions = node.round.on_event(now, &mut node.chain, &ev);
1117 handle(&mut node, actions);
1118 }
1119 }
1120 }
1121 _ = parent_tick.tick() => {
1122 if node.parent.is_some() {
1123 node.peg_tick();
1124 node.reconcile();
1125 let coins = node.peg_coins();
1126 let burns = node.chain.state().pegouts();
1127 if let Some(p) = node.pegout.as_mut() {
1128 let actions = p.tick(unix_now_ms(), &burns, &coins);
1129 handle_pegout(&mut node, actions);
1130 }
1131 }
1132 }
1133 _ = announce_tick.tick() => {
1134 if !relays.is_empty() && !settings.mirrors.is_empty() {
1135 let tip = node.chain.state().tip();
1136 let now = unix_now();
1137 if node.announced != Some(tip.height) && now >= node.announce_retry_at {
1138 let headers = node.headers_hex();
1139 match TipTemplate::new(node.doc.id.clone(), tip.height, headers.clone(), settings.mirrors.clone()).and_then(|t| match &node.chain_hash { Some(h) => t.with_chain_hash(h), None => Ok(t) }).and_then(|t| sign_tip_with_peg(&node.key, &t, Some(&node.doc.challenge), now)) {
1140 Ok(ev) => {
1141 let r = publish_all(&relays, &ev, Duration::from_secs(8)).await;
1142 let ok = ok_count(&r);
1143 if ok > 0 { node.announced = Some(tip.height); } else { node.announce_retry_at = now + 60; }
1144 log(format!("announced tip {} {}… (kind 33333, {} headers, {} mirror(s)) to {ok}/{} relay(s)", tip.height, &tip.hash.to_string()[..16], headers.len(), settings.mirrors.len(), relays.len()));
1145 }
1146 Err(e) => log(format!("announce: {e}")),
1147 }
1148 }
1149 }
1150 }
1151 }
1152 }
1153}
1154
1155pub async fn run(settings: Settings) -> Result<()> {
1158 let doc = ChainDocument::from_json(
1159 &std::fs::read_to_string(&settings.chain)
1160 .map_err(|e| Error::Federation(format!("{}: {e}", settings.chain.display())))?,
1161 )?;
1162 doc.validate()?;
1163 if doc.signers.is_none() {
1164 return Err(Error::Federation(
1165 "not a federated chain: the document has no signers; a level-1 chain is produced by `siding produce`".into(),
1166 ));
1167 }
1168 match doc.family()? {
1169 sidestr_core::parents::Family::Stock => {
1170 run_as::<sidestr_core::block::Stock>(doc, settings).await
1171 }
1172 sidestr_core::parents::Family::Blake2b => {
1173 run_as::<sidestr_header::Blake2bV2>(doc, settings).await
1174 }
1175 }
1176}
1177
1178#[cfg(test)]
1179mod tests {
1180 use super::*;
1181
1182 #[test]
1183 fn ranges_parse_as_a_mirror_sends_them() {
1184 assert_eq!(parse_range("bytes=8-15", 100), Some((8, 15)));
1185 assert_eq!(parse_range("bytes=90-", 100), Some((90, 99)));
1186 assert_eq!(parse_range("bytes=90-100", 100), None);
1187 assert_eq!(parse_range("items=1-2", 100), None);
1188 }
1189
1190 #[test]
1191 fn pegin_records_round_trip() {
1192 let f = FoundPegin {
1193 txid: "ab".repeat(32),
1194 vout: 1,
1195 amount: 5,
1196 script: bitcoin::ScriptBuf::from_bytes(vec![0x51, 0x20]),
1197 height: 9,
1198 parent_address: None,
1199 };
1200 let r = PeginRecord::from(&f);
1201 assert_eq!(r.found(), f);
1202 }
1203
1204 #[test]
1207 fn the_chain_hash_is_read_from_the_event_beside_the_document() {
1208 use sidestr_nostr::chain::{parse_chain_event, sign_chain_event};
1209 use sidestr_nostr::event::SecretKeySigner;
1210 let fixtures = concat!(
1211 env!("CARGO_MANIFEST_DIR"),
1212 "/../sidestr-core/fixtures/fedtest"
1213 );
1214 let text = std::fs::read_to_string(format!("{fixtures}/chain.json")).unwrap();
1215 let key = std::fs::read_to_string(format!("{fixtures}/signer2.key")).unwrap();
1216 let signer = SecretKeySigner::from_hex(&key).unwrap();
1217 let dir =
1218 std::env::temp_dir().join(format!("sidestr-round-chain-hash-{}", std::process::id()));
1219 std::fs::create_dir_all(&dir).unwrap();
1220 let chain = dir.join("chain.json");
1221 std::fs::write(&chain, &text).unwrap();
1222 let (h, at) = read_chain_hash(&chain, None, "sidestr:fedtest").unwrap();
1224 assert_eq!((h, at.clone()), (None, dir.join("chain-event.json")));
1225 let ev = sign_chain_event(&signer, &text, 1_790_100_000).unwrap();
1226 std::fs::write(&at, serde_json::to_string_pretty(&ev).unwrap()).unwrap();
1227 let (h, _) = read_chain_hash(&chain, None, "sidestr:fedtest").unwrap();
1228 assert_eq!(h.as_deref(), Some(ev.id.as_str()));
1229 assert_eq!(parse_chain_event(&ev).unwrap().hash, ev.id);
1230 let e = read_chain_hash(&chain, None, "sidestr:other").unwrap_err();
1232 assert!(
1233 e.to_string()
1234 .contains("is the event of sidestr:fedtest, not sidestr:other"),
1235 "{e}"
1236 );
1237 let mut forged = ev.clone();
1239 forged.content = forged.content.replacen("fedtest", "fedtesu", 1);
1240 let elsewhere = dir.join("forged.json");
1241 std::fs::write(&elsewhere, serde_json::to_string(&forged).unwrap()).unwrap();
1242 assert!(read_chain_hash(&chain, Some(&elsewhere), "sidestr:fedtest").is_err());
1243 let stranger = SecretKeySigner::from_bytes(&[0x31; 32]).unwrap();
1245 let theirs = sign_chain_event(&stranger, &text, 1_790_100_000).unwrap();
1246 std::fs::write(&elsewhere, serde_json::to_string(&theirs).unwrap()).unwrap();
1247 let e = read_chain_hash(&chain, Some(&elsewhere), "sidestr:fedtest").unwrap_err();
1248 assert!(
1249 e.to_string().contains("not one of the document's signers"),
1250 "{e}"
1251 );
1252 let _ = std::fs::remove_dir_all(&dir);
1253 }
1254}