Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1,345 changes: 626 additions & 719 deletions libs/Cargo.lock

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion libs/sdk-core/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ tokio-stream = { version = "0.1.15", features = ["sync"] }
serde_with = "3.3.0"
ryu = "1.0.18"
ldk-node = { git = "https://github.com/lightningdevkit/ldk-node", rev = "v0.7.0-rc.0" }
vss-client = { version = "0.3.1", default-features = false }
vss-client-ng = { git = "https://github.com/lightningdevkit/vss-client", rev = "5934149a0ba75867eb364661bfe15e2d36ce8b88", default-features = false }
r2d2 = "0.8"
r2d2_sqlite = "0.24"

Expand Down
7 changes: 4 additions & 3 deletions libs/sdk-core/src/ldk/backup_transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ use crate::error::{SdkError, SdkResult};
use crate::ldk::store::{VersionedStore, VssStore};
use crate::ldk::store_builder;
use crate::ldk::store_builder::CustomRetryPolicy;
use crate::node_api::NodeResult;
use crate::Config;

pub(crate) struct LdkBackupTransport {
Expand All @@ -12,9 +13,9 @@ pub(crate) struct LdkBackupTransport {
impl LdkBackupTransport {
const KEY: &str = "backup";

pub fn new(config: &Config, seed: &[u8]) -> Self {
let store = store_builder::build_vss_store(config, seed, "backups");
Self { store }
pub fn new(config: &Config, seed: &[u8]) -> NodeResult<Self> {
let store = store_builder::build_vss_store(config, seed, "backups")?;
Ok(Self { store })
}
}

Expand Down
2 changes: 1 addition & 1 deletion libs/sdk-core/src/ldk/node_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ impl Ldk {

builder.set_liquidity_source_lsps2(lsp_id, lsp_address, None);

let vss_store = build_vss_store(&config, &seed, "ldk_node");
let vss_store = build_vss_store(&config, &seed, "ldk_node")?;

// It is not possible to use oneshot here, because `oneshot::Sender::send()`
// consumes itself, not allowing to call `closed()` method after.
Expand Down
2 changes: 1 addition & 1 deletion libs/sdk-core/src/ldk/store/versioned_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ impl std::error::Error for Error {}
/// When updating a key, the version must match the current version in the store,
/// otherwise a conflict error is returned.
/// This trait provides a simplified abstraction over
/// <https://docs.rs/vss-client/latest/vss_client/client/struct.VssClient.html>.
/// <https://docs.rs/vss-client-ng/latest/vss_client_ng/client/struct.VssClient.html>.
#[async_trait]
pub trait VersionedStore {
/// Retrieves a value and its version from the store.
Expand Down
116 changes: 102 additions & 14 deletions libs/sdk-core/src/ldk/store/vss_store.rs
Original file line number Diff line number Diff line change
@@ -1,22 +1,77 @@
use bitcoin::hashes::{sha256, Hash, HashEngine, Hmac, HmacEngine};
use rand::RngCore;
use sdk_common::ensure_sdk;
use tonic::async_trait;
use vss_client::client::VssClient;
use vss_client::error::VssError;
use vss_client::types::{
use vss_client_ng::client::VssClient;
use vss_client_ng::error::VssError;
use vss_client_ng::prost::Message;
use vss_client_ng::types::{
DeleteObjectRequest, GetObjectRequest, GetObjectResponse, KeyValue, ListKeyVersionsRequest,
PutObjectRequest,
PutObjectRequest, Storable,
};
use vss_client::util::retry::RetryPolicy;
use vss_client_ng::util::key_obfuscator::KeyObfuscator;
use vss_client_ng::util::retry::RetryPolicy;
use vss_client_ng::util::storable_builder::{EntropySource, StorableBuilder};

use crate::ldk::store::versioned_store::{Error, VersionedStore};

pub struct VssStore<P: RetryPolicy<E = VssError> + Send + Sync> {
client: VssClient<P>,
store_id: String,
storable_builder: StorableBuilder<RandEntropySource>,
key_obfuscator: KeyObfuscator,
data_encryption_key: [u8; 32],
}

impl<P: RetryPolicy<E = VssError> + Send + Sync> VssStore<P> {
pub fn new(client: VssClient<P>, store_id: String) -> Self {
Self { client, store_id }
pub fn new(client: VssClient<P>, store_id: String, vss_seed: [u8; 32]) -> Self {
let (data_encryption_key, obfuscation_master_key) =
derive_data_encryption_and_obfuscation_keys(&vss_seed);
let key_obfuscator = KeyObfuscator::new(obfuscation_master_key);
let storable_builder = StorableBuilder::new(RandEntropySource);

Self {
client,
store_id,
storable_builder,
key_obfuscator,
data_encryption_key,
}
}

fn obfuscate_key(&self, key: &str) -> String {
self.key_obfuscator.obfuscate(key)
}

fn deobfuscate_key(&self, key: &str) -> Result<String, Error> {
self.key_obfuscator
.deobfuscate(key)
.map_err(|e| Error::Internal(format!("Failed to deobfuscate key: {e}")))
}

/// Persist the payload and commit to the version and the key to cross-check
/// against the VSS metadata returned on reads.
fn construct_storable(&self, key: &str, value: Vec<u8>, version: i64) -> Vec<u8> {
self.storable_builder
.build(
value,
version + 1, // VSS server will increment the version.
&self.data_encryption_key,
key.as_bytes(),
)
.encode_to_vec()
}

fn deconstruct_storable(&self, key: &str, bytes: &[u8]) -> Result<(Vec<u8>, i64), Error> {
let key = self.obfuscate_key(key);
let storable = Storable::decode(bytes).map_err(|e| {
Error::Internal(format!(
"Failed to decode encrypted value for key `{key}`: {e}"
))
})?;
self.storable_builder
.deconstruct(storable, &self.data_encryption_key, key.as_bytes())
.map_err(|e| Error::Internal(format!("Failed to decrypt value for key `{key}`: {e}")))
}
}

Expand All @@ -25,11 +80,19 @@ impl<P: RetryPolicy<E = VssError> + Send + Sync> VersionedStore for VssStore<P>
async fn get(&self, key: String) -> Result<Option<(Vec<u8>, i64)>, Error> {
let request = GetObjectRequest {
store_id: self.store_id.clone(),
key,
key: self.obfuscate_key(&key),
};

match self.client.get_object(&request).await {
Ok(GetObjectResponse { value: Some(kv) }) => Ok(Some((kv.value, kv.version))),
Ok(GetObjectResponse { value: Some(kv) }) => {
let (value, stored_version) = self.deconstruct_storable(&key, &kv.value)?;
ensure_sdk!(stored_version == kv.version,
Error::Internal(format!(
"Version mismatch for key `{key}`: decrypted version={stored_version} but metadata version={}",
kv.version
)));
Ok(Some((value, kv.version)))
}
Ok(GetObjectResponse { value: None }) => Ok(None),
Err(VssError::NoSuchKeyError(_)) => Ok(None),
Err(e) => Err(e.into()),
Expand All @@ -38,9 +101,9 @@ impl<P: RetryPolicy<E = VssError> + Send + Sync> VersionedStore for VssStore<P>

async fn put(&self, key: String, value: Vec<u8>, version: i64) -> Result<(), Error> {
let key_value = KeyValue {
key,
key: self.obfuscate_key(&key),
version,
value,
value: self.construct_storable(&key, value, version),
};
let request = PutObjectRequest {
store_id: self.store_id.clone(),
Expand All @@ -54,7 +117,7 @@ impl<P: RetryPolicy<E = VssError> + Send + Sync> VersionedStore for VssStore<P>

async fn delete(&self, key: String) -> Result<(), Error> {
let key_value = KeyValue {
key,
key: self.obfuscate_key(&key),
version: -1,
value: Vec::new(),
};
Expand Down Expand Up @@ -90,12 +153,37 @@ impl<P: RetryPolicy<E = VssError> + Send + Sync> VersionedStore for VssStore<P>

let versions = versions
.into_iter()
.map(|kv| (kv.key, kv.version))
.collect();
.map(|kv| {
let key = self.deobfuscate_key(&kv.key)?;
Ok((key, kv.version))
})
.collect::<Result<Vec<_>, Error>>()?;
Ok(versions)
}
}

// Copied from https://github.com/lightningdevkit/ldk-node/blob/37045f4708a0721f14bcebe704803418c7c15203/src/io/vss_store.rs#L670
fn derive_data_encryption_and_obfuscation_keys(vss_seed: &[u8; 32]) -> ([u8; 32], [u8; 32]) {
let hkdf = |initial_key_material: &[u8], salt: &[u8]| -> [u8; 32] {
let mut engine = HmacEngine::<sha256::Hash>::new(salt);
engine.input(initial_key_material);
Hmac::from_engine(engine).to_byte_array()
};

let prk = hkdf(vss_seed, b"pseudo_random_key");
let k1 = hkdf(&prk, b"data_encryption_key");
let k2 = hkdf(&prk, &[&k1[..], b"obfuscation_key"].concat());
(k1, k2)
}

struct RandEntropySource;

impl EntropySource for RandEntropySource {
fn fill_bytes(&self, buffer: &mut [u8]) {
rand::thread_rng().fill_bytes(buffer);
}
}

impl From<VssError> for Error {
fn from(err: VssError) -> Self {
match err {
Expand Down
28 changes: 21 additions & 7 deletions libs/sdk-core/src/ldk/store_builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@ use std::path::Path;
use std::sync::Arc;
use std::time::{Duration, SystemTime};

use bitcoin::bip32::{ChildNumber, Xpriv};
use bitcoin::key::Secp256k1;
use hex::ToHex;
use r2d2::Pool;
use r2d2_sqlite::SqliteConnectionManager;
Expand All @@ -14,9 +16,9 @@ use sdk_common::bitcoin::hashes::Hash;
use sdk_common::prelude::Network;
use tokio::runtime::Handle;
use tokio::sync::mpsc;
use vss_client::client::VssClient;
use vss_client::error::VssError;
use vss_client::util::retry::{
use vss_client_ng::client::VssClient;
use vss_client_ng::error::VssError;
use vss_client_ng::util::retry::{
ExponentialBackoffRetryPolicy, FilteredRetryPolicy, JitteredRetryPolicy,
MaxAttemptsRetryPolicy, MaxTotalDelayRetryPolicy, RetryPolicy,
};
Expand All @@ -36,17 +38,29 @@ pub(crate) type CustomRetryPolicy = FilteredRetryPolicy<
pub(crate) type LockingStore = crate::ldk::store::LockingStore<VssStore<CustomRetryPolicy>>;
pub(crate) type MirroringStore = crate::ldk::store::MirroringStore<Arc<LockingStore>, LockingStore>;

const VSS_HARDENED_CHILD_INDEX: u32 = 877;

pub(crate) fn build_vss_store(
config: &Config,
seed: &[u8],
store_id: &str,
) -> VssStore<CustomRetryPolicy> {
) -> NodeResult<VssStore<CustomRetryPolicy>> {
let bitcoin_network: crate::bitcoin::Network = config.network.into();
let xprv = Xpriv::new_master(bitcoin_network, seed)?.derive_priv(
&Secp256k1::new(),
&[ChildNumber::Hardened {
index: VSS_HARDENED_CHILD_INDEX,
}],
)?;

let seed = xprv.private_key.secret_bytes();
let seed_hash = Sha256::hash(&seed);
let store_id = match config.network {
Network::Regtest => {
// Regtest instance of VSS does not implement authentication,
// that is why the hash of the seed is used to avoid collisions.
let seed_hash = Sha256::hash(seed).encode_hex::<String>();
format!("{seed_hash}/{store_id}")
let seed_hash_hex = seed_hash.encode_hex::<String>();
format!("{seed_hash_hex}/{store_id}")
}
_ => store_id.to_string(),
};
Expand All @@ -65,7 +79,7 @@ pub(crate) fn build_vss_store(
}) as _);

let vss_client = VssClient::new(config.vss_url.clone(), retry_policy);
VssStore::new(vss_client, store_id)
Ok(VssStore::new(vss_client, store_id, seed))
}

pub(crate) async fn build_mirroring_store(
Expand Down
2 changes: 1 addition & 1 deletion libs/sdk-core/src/node_builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ pub async fn build_node(
restore_only: Option<bool>,
persister: Arc<SqliteStorage>,
) -> NodeResult<NodeImpls> {
let backup_transport = Arc::new(LdkBackupTransport::new(&config, &seed));
let backup_transport = Arc::new(LdkBackupTransport::new(&config, &seed)?);
let ldk = Ldk::build(config, &seed, restore_only).await?;
let ldk = Arc::new(ldk);
let lsp: Option<Arc<dyn LspAPI>> = Some(ldk.clone());
Expand Down
Loading
Loading