lightning-transaction-sync 0.2.7

Utilities for syncing LDK via the transaction-based `Confirm` interface.
Documentation
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// http://opensource.org/licenses/MIT>, at your option. You may not use this file except in
// accordance with one or both of these licenses.

use bitcoin::block::Header;
use bitcoin::{BlockHash, OutPoint, ScriptBuf, Transaction, Txid};
use lightning::chain::channelmonitor::ANTI_REORG_DELAY;
use lightning::chain::{Confirm, WatchedOutput};

use std::collections::HashMap;
use std::ops::Deref;

// Represents the current state.
pub(crate) struct SyncState {
	// Transactions that were previously processed, but must not be forgotten
	// yet since they still need to be monitored for confirmation on-chain,
	// mapped to the list of `script_pubkey`s we might use to check for their
	// confirmation status. Note a watcher may re-register a transaction with a
	// bogus `script_pubkey` which must not override a previously-registered one.
	pub watched_transactions: HashMap<Txid, Vec<ScriptBuf>>,
	// Outputs that were previously processed, but must not be forgotten yet as
	// as we still need to monitor any spends on-chain.
	pub watched_outputs: HashMap<OutPoint, WatchedOutput>,
	// Outputs for which we previously saw a spend on-chain but kept around until the spends reach
	// sufficient depth.
	pub outputs_spends_pending_threshold_conf: Vec<(Txid, u32, OutPoint, WatchedOutput)>,
	// The tip hash observed during our last sync.
	pub last_sync_hash: Option<BlockHash>,
	// Indicates whether we need to resync, e.g., after encountering an error.
	pub pending_sync: bool,
}

impl SyncState {
	pub fn new() -> Self {
		Self {
			watched_transactions: HashMap::new(),
			watched_outputs: HashMap::new(),
			outputs_spends_pending_threshold_conf: Vec::new(),
			last_sync_hash: None,
			pending_sync: false,
		}
	}
	pub fn sync_unconfirmed_transactions<C: Deref>(
		&mut self, confirmables: &Vec<C>, unconfirmed_txs: Vec<Txid>,
	) where
		C::Target: Confirm,
	{
		for txid in unconfirmed_txs {
			for c in confirmables {
				c.transaction_unconfirmed(&txid);
			}

			self.watched_transactions.entry(txid).or_default();

			// If a previously-confirmed output spend is unconfirmed, re-add the watched output to
			// the tracking map.
			self.outputs_spends_pending_threshold_conf.retain(
				|(conf_txid, _, prev_outpoint, output)| {
					if txid == *conf_txid {
						self.watched_outputs.insert(*prev_outpoint, output.clone());
						false
					} else {
						true
					}
				},
			)
		}
	}

	pub fn sync_confirmed_transactions<C: Deref>(
		&mut self, confirmables: &Vec<C>, confirmed_txs: Vec<ConfirmedTx>,
	) where
		C::Target: Confirm,
	{
		for ctx in confirmed_txs {
			for c in confirmables {
				c.transactions_confirmed(
					&ctx.block_header,
					&[(ctx.pos, &ctx.tx)],
					ctx.block_height,
				);
			}

			self.watched_transactions.remove(&ctx.txid);

			for input in &ctx.tx.input {
				if let Some(output) = self.watched_outputs.remove(&input.previous_output) {
					let spent = (ctx.txid, ctx.block_height, input.previous_output, output);
					self.outputs_spends_pending_threshold_conf.push(spent);
				}
			}
		}
	}

	pub fn prune_output_spends(&mut self, cur_height: u32) {
		self.outputs_spends_pending_threshold_conf
			.retain(|(_, conf_height, _, _)| cur_height < conf_height + ANTI_REORG_DELAY - 1);
	}
}

// A queue that is to be filled by `Filter` and drained during the next syncing round.
pub(crate) struct FilterQueue {
	// Transactions that were registered via the `Filter` interface and have to be processed,
	// mapped to the `script_pubkey`s they were registered with. A transaction may be
	// registered multiple times with different `script_pubkey`s, in which case we keep all of
	// them, as some may be bogus.
	pub transactions: HashMap<Txid, Vec<ScriptBuf>>,
	// Outputs that were registered via the `Filter` interface and have to be processed.
	pub outputs: HashMap<OutPoint, WatchedOutput>,
}

impl FilterQueue {
	pub fn new() -> Self {
		Self { transactions: HashMap::new(), outputs: HashMap::new() }
	}

	// Processes the transaction and output queues and adds them to the given [`SyncState`].
	//
	// Returns `true` if new items had been registered.
	pub fn process_queues(&mut self, sync_state: &mut SyncState) -> bool {
		let mut pending_registrations = false;

		if !self.transactions.is_empty() {
			pending_registrations = true;

			for (txid, script_pubkeys) in self.transactions.drain() {
				let watched = sync_state.watched_transactions.entry(txid).or_default();
				for script_pubkey in script_pubkeys {
					if !watched.contains(&script_pubkey) {
						watched.push(script_pubkey);
					}
				}
			}
		}

		if !self.outputs.is_empty() {
			pending_registrations = true;

			sync_state.watched_outputs.extend(self.outputs.drain());
		}
		pending_registrations
	}
}

#[derive(Debug)]
pub(crate) struct ConfirmedTx {
	pub tx: Transaction,
	pub txid: Txid,
	pub block_header: Header,
	pub block_height: u32,
	pub pos: usize,
}