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
157 changes: 157 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -81,6 +81,8 @@ where
Ok(updated)
}

/// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
/// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
pub(crate) async fn insert_or_update(&self, object: SO) -> Result<bool, Error> {
let _guard = self.mutation_lock.lock().await;

Expand DownExpand Up@@ -170,6 +172,36 @@ where
Ok(DataStoreUpdateResult::Updated)
}

/// Atomically transforms the entry for `id` through `f` and persists the result.
///
/// `f` receives the current entry (`None` when absent) and returns the new state to write;
/// returning `None` leaves the store untouched. The read, the closure, and the write share
/// one critical section of the mutation lock, so no concurrent writer can land in between —
/// unlike a separate [`Self::get`] followed by an insert or update.
///
/// The closure runs on a clone of the entry with the in-memory map lock released, so it may
/// freely read this store or others (reads see the pre-mutation state) without ordering map
/// locks against each other. Keep it cheap and non-blocking.
///
/// Returns the written object, or `None` when the closure declined to write.
pub(crate) async fn mutate<F: FnOnce(Option<&SO>) -> Option<SO>>(
&self, id: &SO::Id, f: F,
) -> Result<Option<SO>, Error> {
let _guard = self.mutation_lock.lock().await;

let current = self.objects.lock().expect("lock").get(id).cloned();
let new_object = match f(current.as_ref()) {
Some(new_object) => new_object,
None => return Ok(None),
};
debug_assert!(new_object.id() == *id, "mutate closure must not change the object's id");

self.persist(&new_object).await?;
let mut locked_objects = self.objects.lock().expect("lock");
locked_objects.insert(new_object.id(), new_object.clone());
Ok(Some(new_object))
}

/// Returns in-memory objects matching `f`.
///
/// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
Expand DownExpand Up@@ -403,6 +435,131 @@ mod tests {
assert_eq!(Ok(true), data_store.insert_or_update(new_iou_object).await);
}

#[tokio::test]
async fn mutate_inserts_when_absent() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let primary_namespace = "datastore_test_primary".to_string();
let secondary_namespace = "datastore_test_secondary".to_string();
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
Vec::new(),
primary_namespace.clone(),
secondary_namespace.clone(),
Arc::clone(&store),
logger,
);

let id = TestObjectId { id: [42u8; 4] };
let object = TestObject { id, data: [23u8; 3] };
let result = data_store
.mutate(&id, |existing| {
assert!(existing.is_none());
Some(object)
})
.await;
assert_eq!(Ok(Some(object)), result);

assert_eq!(Some(object), data_store.get(&id));
let store_key = id.encode_to_hex_str();
assert!(KVStore::read(&*store, &primary_namespace, &secondary_namespace, &store_key)
.await
.is_ok());
}

#[tokio::test]
async fn mutate_transforms_existing_entry() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// The closure sees the current entry and derives the new state from it.
let result = data_store
.mutate(&id, |existing| {
let mut new_object = *existing.unwrap();
new_object.data[0] += 1;
Some(new_object)
})
.await;
let expected = TestObject { id, data: [24u8, 23u8, 23u8] };
assert_eq!(Ok(Some(expected)), result);
assert_eq!(Some(expected), data_store.get(&id));
}

#[tokio::test]
async fn mutate_runs_the_closure_without_the_map_lock() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// Closures gate cross-store decisions on reads of other stores, which lock their own
// in-memory maps. Holding this store's map lock across the closure would order it
// before theirs and invite lock-order inversions, so the closure must run with the map
// lock released.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
assert!(data_store.objects.try_lock().is_ok());
None
})
.await;
assert_eq!(Ok(None), result);
}

#[tokio::test]
async fn mutate_persists_nothing_when_closure_declines() {
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

// Returning `None` must not attempt a write (the store fails all writes) nor touch memory.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
None
})
.await;
assert_eq!(Ok(None), result);
assert_eq!(Some(existing_object), data_store.get(&id));
}

#[tokio::test]
async fn mutate_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id: existing_id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

let changed = TestObject { id: existing_id, data: [24u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&existing_id, |_| Some(changed)).await
);
assert_eq!(Some(existing_object), data_store.get(&existing_id));

let new_id = TestObjectId { id: [55u8; 4] };
let new_object = TestObject { id: new_id, data: [34u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&new_id, |_| Some(new_object)).await
);
assert!(data_store.get(&new_id).is_none());
}

#[tokio::test]
async fn insert_or_update_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
Expand Down
75 changes: 75 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -242,4 +242,79 @@ mod tests {
"current txid must not remain in its own conflict list"
);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;

use crate::payment::store::PaymentDetailsUpdate;

let txid = test_txid(7);
let payment_id = PaymentId(txid.to_byte_array());

// A pending entry wallet sync has already mirrored a confirmation into (via
// `apply_funding_status_update_locked`) before classification ran.
let confirmed_details = PaymentDetails::new(
payment_id,
PaymentKind::Onchain {
txid,
status: ConfirmationStatus::Confirmed {
block_hash: BlockHash::from_byte_array([8u8; 32]),
height: 100,
timestamp: 1,
},
tx_type: None,
},
Some(2_000_000),
Some(999),
PaymentDirection::Outbound,
PaymentStatus::Pending,
);
let mirrored = PendingPaymentDetails::new(confirmed_details, Vec::new(), Vec::new());

// A fresh classification is always Unconfirmed and carries the candidate history; its
// figures are the active candidate's.
let fresh = pending_onchain_payment(payment_id, txid);
let candidates = vec![FundingTxCandidate {
txid,
amount_msat: fresh.amount_msat,
fee_paid_msat: fresh.fee_paid_msat,
}];

// The old fresh-insert path merged the full fresh record, downgrading the mirrored
// confirmation.
let mut downgraded = mirrored.clone();
let full_update =
PendingPaymentDetails::new(fresh.clone(), Vec::new(), candidates.clone()).to_update();
assert!(downgraded.update(full_update));
assert!(
matches!(
downgraded.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. }
),
"a full merge of a fresh classification downgrades a mirrored confirmation",
);

// The narrow classification update merges the candidates while preserving the
// confirmation state wallet sync owns. It names the confirmed txid, so its
// contribution-derived figures replace the mirrored wallet-view ones.
let mut merged = mirrored.clone();
let narrow_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)),
conflicting_txids: None,
candidates: candidates.clone(),
};
assert!(merged.update(narrow_update));
assert!(
matches!(
merged.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. }
),
"a narrow classification update must not downgrade a mirrored confirmation",
);
assert_eq!(merged.candidates, candidates);
assert_eq!(merged.details.amount_msat, Some(1_000));
assert_eq!(merged.details.fee_paid_msat, Some(100));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
157 changes: 157 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -81,6 +81,8 @@ where
Ok(updated)
}

/// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
/// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
pub(crate) async fn insert_or_update(&self, object: SO) -> Result<bool, Error> {
let _guard = self.mutation_lock.lock().await;

Expand DownExpand Up@@ -170,6 +172,36 @@ where
Ok(DataStoreUpdateResult::Updated)
}

/// Atomically transforms the entry for `id` through `f` and persists the result.
///
/// `f` receives the current entry (`None` when absent) and returns the new state to write;
/// returning `None` leaves the store untouched. The read, the closure, and the write share
/// one critical section of the mutation lock, so no concurrent writer can land in between —
/// unlike a separate [`Self::get`] followed by an insert or update.
///
/// The closure runs on a clone of the entry with the in-memory map lock released, so it may
/// freely read this store or others (reads see the pre-mutation state) without ordering map
/// locks against each other. Keep it cheap and non-blocking.
///
/// Returns the written object, or `None` when the closure declined to write.
pub(crate) async fn mutate<F: FnOnce(Option<&SO>) -> Option<SO>>(
&self, id: &SO::Id, f: F,
) -> Result<Option<SO>, Error> {
let _guard = self.mutation_lock.lock().await;

let current = self.objects.lock().expect("lock").get(id).cloned();
let new_object = match f(current.as_ref()) {
Some(new_object) => new_object,
None => return Ok(None),
};
debug_assert!(new_object.id() == *id, "mutate closure must not change the object's id");

self.persist(&new_object).await?;
let mut locked_objects = self.objects.lock().expect("lock");
locked_objects.insert(new_object.id(), new_object.clone());
Ok(Some(new_object))
}

/// Returns in-memory objects matching `f`.
///
/// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
Expand DownExpand Up@@ -403,6 +435,131 @@ mod tests {
assert_eq!(Ok(true), data_store.insert_or_update(new_iou_object).await);
}

#[tokio::test]
async fn mutate_inserts_when_absent() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let primary_namespace = "datastore_test_primary".to_string();
let secondary_namespace = "datastore_test_secondary".to_string();
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
Vec::new(),
primary_namespace.clone(),
secondary_namespace.clone(),
Arc::clone(&store),
logger,
);

let id = TestObjectId { id: [42u8; 4] };
let object = TestObject { id, data: [23u8; 3] };
let result = data_store
.mutate(&id, |existing| {
assert!(existing.is_none());
Some(object)
})
.await;
assert_eq!(Ok(Some(object)), result);

assert_eq!(Some(object), data_store.get(&id));
let store_key = id.encode_to_hex_str();
assert!(KVStore::read(&*store, &primary_namespace, &secondary_namespace, &store_key)
.await
.is_ok());
}

#[tokio::test]
async fn mutate_transforms_existing_entry() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// The closure sees the current entry and derives the new state from it.
let result = data_store
.mutate(&id, |existing| {
let mut new_object = *existing.unwrap();
new_object.data[0] += 1;
Some(new_object)
})
.await;
let expected = TestObject { id, data: [24u8, 23u8, 23u8] };
assert_eq!(Ok(Some(expected)), result);
assert_eq!(Some(expected), data_store.get(&id));
}

#[tokio::test]
async fn mutate_runs_the_closure_without_the_map_lock() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// Closures gate cross-store decisions on reads of other stores, which lock their own
// in-memory maps. Holding this store's map lock across the closure would order it
// before theirs and invite lock-order inversions, so the closure must run with the map
// lock released.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
assert!(data_store.objects.try_lock().is_ok());
None
})
.await;
assert_eq!(Ok(None), result);
}

#[tokio::test]
async fn mutate_persists_nothing_when_closure_declines() {
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

// Returning `None` must not attempt a write (the store fails all writes) nor touch memory.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
None
})
.await;
assert_eq!(Ok(None), result);
assert_eq!(Some(existing_object), data_store.get(&id));
}

#[tokio::test]
async fn mutate_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id: existing_id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

let changed = TestObject { id: existing_id, data: [24u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&existing_id, |_| Some(changed)).await
);
assert_eq!(Some(existing_object), data_store.get(&existing_id));

let new_id = TestObjectId { id: [55u8; 4] };
let new_object = TestObject { id: new_id, data: [34u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&new_id, |_| Some(new_object)).await
);
assert!(data_store.get(&new_id).is_none());
}

#[tokio::test]
async fn insert_or_update_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
Expand Down
75 changes: 75 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -242,4 +242,79 @@ mod tests {
"current txid must not remain in its own conflict list"
);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;

use crate::payment::store::PaymentDetailsUpdate;

let txid = test_txid(7);
let payment_id = PaymentId(txid.to_byte_array());

// A pending entry wallet sync has already mirrored a confirmation into (via
// `apply_funding_status_update_locked`) before classification ran.
let confirmed_details = PaymentDetails::new(
payment_id,
PaymentKind::Onchain {
txid,
status: ConfirmationStatus::Confirmed {
block_hash: BlockHash::from_byte_array([8u8; 32]),
height: 100,
timestamp: 1,
},
tx_type: None,
},
Some(2_000_000),
Some(999),
PaymentDirection::Outbound,
PaymentStatus::Pending,
);
let mirrored = PendingPaymentDetails::new(confirmed_details, Vec::new(), Vec::new());

// A fresh classification is always Unconfirmed and carries the candidate history; its
// figures are the active candidate's.
let fresh = pending_onchain_payment(payment_id, txid);
let candidates = vec![FundingTxCandidate {
txid,
amount_msat: fresh.amount_msat,
fee_paid_msat: fresh.fee_paid_msat,
}];

// The old fresh-insert path merged the full fresh record, downgrading the mirrored
// confirmation.
let mut downgraded = mirrored.clone();
let full_update =
PendingPaymentDetails::new(fresh.clone(), Vec::new(), candidates.clone()).to_update();
assert!(downgraded.update(full_update));
assert!(
matches!(
downgraded.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. }
),
"a full merge of a fresh classification downgrades a mirrored confirmation",
);

// The narrow classification update merges the candidates while preserving the
// confirmation state wallet sync owns. It names the confirmed txid, so its
// contribution-derived figures replace the mirrored wallet-view ones.
let mut merged = mirrored.clone();
let narrow_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)),
conflicting_txids: None,
candidates: candidates.clone(),
};
assert!(merged.update(narrow_update));
assert!(
matches!(
merged.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. }
),
"a narrow classification update must not downgrade a mirrored confirmation",
);
assert_eq!(merged.candidates, candidates);
assert_eq!(merged.details.amount_msat, Some(1_000));
assert_eq!(merged.details.fee_paid_msat, Some(100));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
157 changes: 157 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -81,6 +81,8 @@ where
Ok(updated)
}

/// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
/// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
pub(crate) async fn insert_or_update(&self, object: SO) -> Result<bool, Error> {
let _guard = self.mutation_lock.lock().await;

Expand DownExpand Up@@ -170,6 +172,36 @@ where
Ok(DataStoreUpdateResult::Updated)
}

/// Atomically transforms the entry for `id` through `f` and persists the result.
///
/// `f` receives the current entry (`None` when absent) and returns the new state to write;
/// returning `None` leaves the store untouched. The read, the closure, and the write share
/// one critical section of the mutation lock, so no concurrent writer can land in between —
/// unlike a separate [`Self::get`] followed by an insert or update.
///
/// The closure runs on a clone of the entry with the in-memory map lock released, so it may
/// freely read this store or others (reads see the pre-mutation state) without ordering map
/// locks against each other. Keep it cheap and non-blocking.
///
/// Returns the written object, or `None` when the closure declined to write.
pub(crate) async fn mutate<F: FnOnce(Option<&SO>) -> Option<SO>>(
&self, id: &SO::Id, f: F,
) -> Result<Option<SO>, Error> {
let _guard = self.mutation_lock.lock().await;

let current = self.objects.lock().expect("lock").get(id).cloned();
let new_object = match f(current.as_ref()) {
Some(new_object) => new_object,
None => return Ok(None),
};
debug_assert!(new_object.id() == *id, "mutate closure must not change the object's id");

self.persist(&new_object).await?;
let mut locked_objects = self.objects.lock().expect("lock");
locked_objects.insert(new_object.id(), new_object.clone());
Ok(Some(new_object))
}

/// Returns in-memory objects matching `f`.
///
/// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
Expand DownExpand Up@@ -403,6 +435,131 @@ mod tests {
assert_eq!(Ok(true), data_store.insert_or_update(new_iou_object).await);
}

#[tokio::test]
async fn mutate_inserts_when_absent() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let primary_namespace = "datastore_test_primary".to_string();
let secondary_namespace = "datastore_test_secondary".to_string();
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
Vec::new(),
primary_namespace.clone(),
secondary_namespace.clone(),
Arc::clone(&store),
logger,
);

let id = TestObjectId { id: [42u8; 4] };
let object = TestObject { id, data: [23u8; 3] };
let result = data_store
.mutate(&id, |existing| {
assert!(existing.is_none());
Some(object)
})
.await;
assert_eq!(Ok(Some(object)), result);

assert_eq!(Some(object), data_store.get(&id));
let store_key = id.encode_to_hex_str();
assert!(KVStore::read(&*store, &primary_namespace, &secondary_namespace, &store_key)
.await
.is_ok());
}

#[tokio::test]
async fn mutate_transforms_existing_entry() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// The closure sees the current entry and derives the new state from it.
let result = data_store
.mutate(&id, |existing| {
let mut new_object = *existing.unwrap();
new_object.data[0] += 1;
Some(new_object)
})
.await;
let expected = TestObject { id, data: [24u8, 23u8, 23u8] };
assert_eq!(Ok(Some(expected)), result);
assert_eq!(Some(expected), data_store.get(&id));
}

#[tokio::test]
async fn mutate_runs_the_closure_without_the_map_lock() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// Closures gate cross-store decisions on reads of other stores, which lock their own
// in-memory maps. Holding this store's map lock across the closure would order it
// before theirs and invite lock-order inversions, so the closure must run with the map
// lock released.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
assert!(data_store.objects.try_lock().is_ok());
None
})
.await;
assert_eq!(Ok(None), result);
}

#[tokio::test]
async fn mutate_persists_nothing_when_closure_declines() {
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

// Returning `None` must not attempt a write (the store fails all writes) nor touch memory.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
None
})
.await;
assert_eq!(Ok(None), result);
assert_eq!(Some(existing_object), data_store.get(&id));
}

#[tokio::test]
async fn mutate_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id: existing_id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

let changed = TestObject { id: existing_id, data: [24u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&existing_id, |_| Some(changed)).await
);
assert_eq!(Some(existing_object), data_store.get(&existing_id));

let new_id = TestObjectId { id: [55u8; 4] };
let new_object = TestObject { id: new_id, data: [34u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&new_id, |_| Some(new_object)).await
);
assert!(data_store.get(&new_id).is_none());
}

#[tokio::test]
async fn insert_or_update_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
Expand Down
75 changes: 75 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -242,4 +242,79 @@ mod tests {
"current txid must not remain in its own conflict list"
);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;

use crate::payment::store::PaymentDetailsUpdate;

let txid = test_txid(7);
let payment_id = PaymentId(txid.to_byte_array());

// A pending entry wallet sync has already mirrored a confirmation into (via
// `apply_funding_status_update_locked`) before classification ran.
let confirmed_details = PaymentDetails::new(
payment_id,
PaymentKind::Onchain {
txid,
status: ConfirmationStatus::Confirmed {
block_hash: BlockHash::from_byte_array([8u8; 32]),
height: 100,
timestamp: 1,
},
tx_type: None,
},
Some(2_000_000),
Some(999),
PaymentDirection::Outbound,
PaymentStatus::Pending,
);
let mirrored = PendingPaymentDetails::new(confirmed_details, Vec::new(), Vec::new());

// A fresh classification is always Unconfirmed and carries the candidate history; its
// figures are the active candidate's.
let fresh = pending_onchain_payment(payment_id, txid);
let candidates = vec![FundingTxCandidate {
txid,
amount_msat: fresh.amount_msat,
fee_paid_msat: fresh.fee_paid_msat,
}];

// The old fresh-insert path merged the full fresh record, downgrading the mirrored
// confirmation.
let mut downgraded = mirrored.clone();
let full_update =
PendingPaymentDetails::new(fresh.clone(), Vec::new(), candidates.clone()).to_update();
assert!(downgraded.update(full_update));
assert!(
matches!(
downgraded.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. }
),
"a full merge of a fresh classification downgrades a mirrored confirmation",
);

// The narrow classification update merges the candidates while preserving the
// confirmation state wallet sync owns. It names the confirmed txid, so its
// contribution-derived figures replace the mirrored wallet-view ones.
let mut merged = mirrored.clone();
let narrow_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)),
conflicting_txids: None,
candidates: candidates.clone(),
};
assert!(merged.update(narrow_update));
assert!(
matches!(
merged.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. }
),
"a narrow classification update must not downgrade a mirrored confirmation",
);
assert_eq!(merged.candidates, candidates);
assert_eq!(merged.details.amount_msat, Some(1_000));
assert_eq!(merged.details.fee_paid_msat, Some(100));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
157 changes: 157 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -81,6 +81,8 @@ where
Ok(updated)
}

/// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
/// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
pub(crate) async fn insert_or_update(&self, object: SO) -> Result<bool, Error> {
let _guard = self.mutation_lock.lock().await;

Expand DownExpand Up@@ -170,6 +172,36 @@ where
Ok(DataStoreUpdateResult::Updated)
}

/// Atomically transforms the entry for `id` through `f` and persists the result.
///
/// `f` receives the current entry (`None` when absent) and returns the new state to write;
/// returning `None` leaves the store untouched. The read, the closure, and the write share
/// one critical section of the mutation lock, so no concurrent writer can land in between —
/// unlike a separate [`Self::get`] followed by an insert or update.
///
/// The closure runs on a clone of the entry with the in-memory map lock released, so it may
/// freely read this store or others (reads see the pre-mutation state) without ordering map
/// locks against each other. Keep it cheap and non-blocking.
///
/// Returns the written object, or `None` when the closure declined to write.
pub(crate) async fn mutate<F: FnOnce(Option<&SO>) -> Option<SO>>(
&self, id: &SO::Id, f: F,
) -> Result<Option<SO>, Error> {
let _guard = self.mutation_lock.lock().await;

let current = self.objects.lock().expect("lock").get(id).cloned();
let new_object = match f(current.as_ref()) {
Some(new_object) => new_object,
None => return Ok(None),
};
debug_assert!(new_object.id() == *id, "mutate closure must not change the object's id");

self.persist(&new_object).await?;
let mut locked_objects = self.objects.lock().expect("lock");
locked_objects.insert(new_object.id(), new_object.clone());
Ok(Some(new_object))
}

/// Returns in-memory objects matching `f`.
///
/// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
Expand DownExpand Up@@ -403,6 +435,131 @@ mod tests {
assert_eq!(Ok(true), data_store.insert_or_update(new_iou_object).await);
}

#[tokio::test]
async fn mutate_inserts_when_absent() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let primary_namespace = "datastore_test_primary".to_string();
let secondary_namespace = "datastore_test_secondary".to_string();
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
Vec::new(),
primary_namespace.clone(),
secondary_namespace.clone(),
Arc::clone(&store),
logger,
);

let id = TestObjectId { id: [42u8; 4] };
let object = TestObject { id, data: [23u8; 3] };
let result = data_store
.mutate(&id, |existing| {
assert!(existing.is_none());
Some(object)
})
.await;
assert_eq!(Ok(Some(object)), result);

assert_eq!(Some(object), data_store.get(&id));
let store_key = id.encode_to_hex_str();
assert!(KVStore::read(&*store, &primary_namespace, &secondary_namespace, &store_key)
.await
.is_ok());
}

#[tokio::test]
async fn mutate_transforms_existing_entry() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// The closure sees the current entry and derives the new state from it.
let result = data_store
.mutate(&id, |existing| {
let mut new_object = *existing.unwrap();
new_object.data[0] += 1;
Some(new_object)
})
.await;
let expected = TestObject { id, data: [24u8, 23u8, 23u8] };
assert_eq!(Ok(Some(expected)), result);
assert_eq!(Some(expected), data_store.get(&id));
}

#[tokio::test]
async fn mutate_runs_the_closure_without_the_map_lock() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// Closures gate cross-store decisions on reads of other stores, which lock their own
// in-memory maps. Holding this store's map lock across the closure would order it
// before theirs and invite lock-order inversions, so the closure must run with the map
// lock released.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
assert!(data_store.objects.try_lock().is_ok());
None
})
.await;
assert_eq!(Ok(None), result);
}

#[tokio::test]
async fn mutate_persists_nothing_when_closure_declines() {
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

// Returning `None` must not attempt a write (the store fails all writes) nor touch memory.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
None
})
.await;
assert_eq!(Ok(None), result);
assert_eq!(Some(existing_object), data_store.get(&id));
}

#[tokio::test]
async fn mutate_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id: existing_id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

let changed = TestObject { id: existing_id, data: [24u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&existing_id, |_| Some(changed)).await
);
assert_eq!(Some(existing_object), data_store.get(&existing_id));

let new_id = TestObjectId { id: [55u8; 4] };
let new_object = TestObject { id: new_id, data: [34u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&new_id, |_| Some(new_object)).await
);
assert!(data_store.get(&new_id).is_none());
}

#[tokio::test]
async fn insert_or_update_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
Expand Down
75 changes: 75 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -242,4 +242,79 @@ mod tests {
"current txid must not remain in its own conflict list"
);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;

use crate::payment::store::PaymentDetailsUpdate;

let txid = test_txid(7);
let payment_id = PaymentId(txid.to_byte_array());

// A pending entry wallet sync has already mirrored a confirmation into (via
// `apply_funding_status_update_locked`) before classification ran.
let confirmed_details = PaymentDetails::new(
payment_id,
PaymentKind::Onchain {
txid,
status: ConfirmationStatus::Confirmed {
block_hash: BlockHash::from_byte_array([8u8; 32]),
height: 100,
timestamp: 1,
},
tx_type: None,
},
Some(2_000_000),
Some(999),
PaymentDirection::Outbound,
PaymentStatus::Pending,
);
let mirrored = PendingPaymentDetails::new(confirmed_details, Vec::new(), Vec::new());

// A fresh classification is always Unconfirmed and carries the candidate history; its
// figures are the active candidate's.
let fresh = pending_onchain_payment(payment_id, txid);
let candidates = vec![FundingTxCandidate {
txid,
amount_msat: fresh.amount_msat,
fee_paid_msat: fresh.fee_paid_msat,
}];

// The old fresh-insert path merged the full fresh record, downgrading the mirrored
// confirmation.
let mut downgraded = mirrored.clone();
let full_update =
PendingPaymentDetails::new(fresh.clone(), Vec::new(), candidates.clone()).to_update();
assert!(downgraded.update(full_update));
assert!(
matches!(
downgraded.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. }
),
"a full merge of a fresh classification downgrades a mirrored confirmation",
);

// The narrow classification update merges the candidates while preserving the
// confirmation state wallet sync owns. It names the confirmed txid, so its
// contribution-derived figures replace the mirrored wallet-view ones.
let mut merged = mirrored.clone();
let narrow_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)),
conflicting_txids: None,
candidates: candidates.clone(),
};
assert!(merged.update(narrow_update));
assert!(
matches!(
merged.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. }
),
"a narrow classification update must not downgrade a mirrored confirmation",
);
assert_eq!(merged.candidates, candidates);
assert_eq!(merged.details.amount_msat, Some(1_000));
assert_eq!(merged.details.fee_paid_msat, Some(100));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
157 changes: 157 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -81,6 +81,8 @@ where
Ok(updated)
}

/// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
/// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
pub(crate) async fn insert_or_update(&self, object: SO) -> Result<bool, Error> {
let _guard = self.mutation_lock.lock().await;

Expand DownExpand Up@@ -170,6 +172,36 @@ where
Ok(DataStoreUpdateResult::Updated)
}

/// Atomically transforms the entry for `id` through `f` and persists the result.
///
/// `f` receives the current entry (`None` when absent) and returns the new state to write;
/// returning `None` leaves the store untouched. The read, the closure, and the write share
/// one critical section of the mutation lock, so no concurrent writer can land in between —
/// unlike a separate [`Self::get`] followed by an insert or update.
///
/// The closure runs on a clone of the entry with the in-memory map lock released, so it may
/// freely read this store or others (reads see the pre-mutation state) without ordering map
/// locks against each other. Keep it cheap and non-blocking.
///
/// Returns the written object, or `None` when the closure declined to write.
pub(crate) async fn mutate<F: FnOnce(Option<&SO>) -> Option<SO>>(
&self, id: &SO::Id, f: F,
) -> Result<Option<SO>, Error> {
let _guard = self.mutation_lock.lock().await;

let current = self.objects.lock().expect("lock").get(id).cloned();
let new_object = match f(current.as_ref()) {
Some(new_object) => new_object,
None => return Ok(None),
};
debug_assert!(new_object.id() == *id, "mutate closure must not change the object's id");

self.persist(&new_object).await?;
let mut locked_objects = self.objects.lock().expect("lock");
locked_objects.insert(new_object.id(), new_object.clone());
Ok(Some(new_object))
}

/// Returns in-memory objects matching `f`.
///
/// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
Expand DownExpand Up@@ -403,6 +435,131 @@ mod tests {
assert_eq!(Ok(true), data_store.insert_or_update(new_iou_object).await);
}

#[tokio::test]
async fn mutate_inserts_when_absent() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let primary_namespace = "datastore_test_primary".to_string();
let secondary_namespace = "datastore_test_secondary".to_string();
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
Vec::new(),
primary_namespace.clone(),
secondary_namespace.clone(),
Arc::clone(&store),
logger,
);

let id = TestObjectId { id: [42u8; 4] };
let object = TestObject { id, data: [23u8; 3] };
let result = data_store
.mutate(&id, |existing| {
assert!(existing.is_none());
Some(object)
})
.await;
assert_eq!(Ok(Some(object)), result);

assert_eq!(Some(object), data_store.get(&id));
let store_key = id.encode_to_hex_str();
assert!(KVStore::read(&*store, &primary_namespace, &secondary_namespace, &store_key)
.await
.is_ok());
}

#[tokio::test]
async fn mutate_transforms_existing_entry() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// The closure sees the current entry and derives the new state from it.
let result = data_store
.mutate(&id, |existing| {
let mut new_object = *existing.unwrap();
new_object.data[0] += 1;
Some(new_object)
})
.await;
let expected = TestObject { id, data: [24u8, 23u8, 23u8] };
assert_eq!(Ok(Some(expected)), result);
assert_eq!(Some(expected), data_store.get(&id));
}

#[tokio::test]
async fn mutate_runs_the_closure_without_the_map_lock() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// Closures gate cross-store decisions on reads of other stores, which lock their own
// in-memory maps. Holding this store's map lock across the closure would order it
// before theirs and invite lock-order inversions, so the closure must run with the map
// lock released.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
assert!(data_store.objects.try_lock().is_ok());
None
})
.await;
assert_eq!(Ok(None), result);
}

#[tokio::test]
async fn mutate_persists_nothing_when_closure_declines() {
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

// Returning `None` must not attempt a write (the store fails all writes) nor touch memory.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
None
})
.await;
assert_eq!(Ok(None), result);
assert_eq!(Some(existing_object), data_store.get(&id));
}

#[tokio::test]
async fn mutate_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id: existing_id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

let changed = TestObject { id: existing_id, data: [24u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&existing_id, |_| Some(changed)).await
);
assert_eq!(Some(existing_object), data_store.get(&existing_id));

let new_id = TestObjectId { id: [55u8; 4] };
let new_object = TestObject { id: new_id, data: [34u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&new_id, |_| Some(new_object)).await
);
assert!(data_store.get(&new_id).is_none());
}

#[tokio::test]
async fn insert_or_update_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
Expand Down
75 changes: 75 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -242,4 +242,79 @@ mod tests {
"current txid must not remain in its own conflict list"
);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;

use crate::payment::store::PaymentDetailsUpdate;

let txid = test_txid(7);
let payment_id = PaymentId(txid.to_byte_array());

// A pending entry wallet sync has already mirrored a confirmation into (via
// `apply_funding_status_update_locked`) before classification ran.
let confirmed_details = PaymentDetails::new(
payment_id,
PaymentKind::Onchain {
txid,
status: ConfirmationStatus::Confirmed {
block_hash: BlockHash::from_byte_array([8u8; 32]),
height: 100,
timestamp: 1,
},
tx_type: None,
},
Some(2_000_000),
Some(999),
PaymentDirection::Outbound,
PaymentStatus::Pending,
);
let mirrored = PendingPaymentDetails::new(confirmed_details, Vec::new(), Vec::new());

// A fresh classification is always Unconfirmed and carries the candidate history; its
// figures are the active candidate's.
let fresh = pending_onchain_payment(payment_id, txid);
let candidates = vec![FundingTxCandidate {
txid,
amount_msat: fresh.amount_msat,
fee_paid_msat: fresh.fee_paid_msat,
}];

// The old fresh-insert path merged the full fresh record, downgrading the mirrored
// confirmation.
let mut downgraded = mirrored.clone();
let full_update =
PendingPaymentDetails::new(fresh.clone(), Vec::new(), candidates.clone()).to_update();
assert!(downgraded.update(full_update));
assert!(
matches!(
downgraded.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. }
),
"a full merge of a fresh classification downgrades a mirrored confirmation",
);

// The narrow classification update merges the candidates while preserving the
// confirmation state wallet sync owns. It names the confirmed txid, so its
// contribution-derived figures replace the mirrored wallet-view ones.
let mut merged = mirrored.clone();
let narrow_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)),
conflicting_txids: None,
candidates: candidates.clone(),
};
assert!(merged.update(narrow_update));
assert!(
matches!(
merged.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. }
),
"a narrow classification update must not downgrade a mirrored confirmation",
);
assert_eq!(merged.candidates, candidates);
assert_eq!(merged.details.amount_msat, Some(1_000));
assert_eq!(merged.details.fee_paid_msat, Some(100));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
157 changes: 157 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -81,6 +81,8 @@ where
Ok(updated)
}

/// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
/// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
pub(crate) async fn insert_or_update(&self, object: SO) -> Result<bool, Error> {
let _guard = self.mutation_lock.lock().await;

Expand DownExpand Up@@ -170,6 +172,36 @@ where
Ok(DataStoreUpdateResult::Updated)
}

/// Atomically transforms the entry for `id` through `f` and persists the result.
///
/// `f` receives the current entry (`None` when absent) and returns the new state to write;
/// returning `None` leaves the store untouched. The read, the closure, and the write share
/// one critical section of the mutation lock, so no concurrent writer can land in between —
/// unlike a separate [`Self::get`] followed by an insert or update.
///
/// The closure runs on a clone of the entry with the in-memory map lock released, so it may
/// freely read this store or others (reads see the pre-mutation state) without ordering map
/// locks against each other. Keep it cheap and non-blocking.
///
/// Returns the written object, or `None` when the closure declined to write.
pub(crate) async fn mutate<F: FnOnce(Option<&SO>) -> Option<SO>>(
&self, id: &SO::Id, f: F,
) -> Result<Option<SO>, Error> {
let _guard = self.mutation_lock.lock().await;

let current = self.objects.lock().expect("lock").get(id).cloned();
let new_object = match f(current.as_ref()) {
Some(new_object) => new_object,
None => return Ok(None),
};
debug_assert!(new_object.id() == *id, "mutate closure must not change the object's id");

self.persist(&new_object).await?;
let mut locked_objects = self.objects.lock().expect("lock");
locked_objects.insert(new_object.id(), new_object.clone());
Ok(Some(new_object))
}

/// Returns in-memory objects matching `f`.
///
/// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
Expand DownExpand Up@@ -403,6 +435,131 @@ mod tests {
assert_eq!(Ok(true), data_store.insert_or_update(new_iou_object).await);
}

#[tokio::test]
async fn mutate_inserts_when_absent() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let primary_namespace = "datastore_test_primary".to_string();
let secondary_namespace = "datastore_test_secondary".to_string();
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
Vec::new(),
primary_namespace.clone(),
secondary_namespace.clone(),
Arc::clone(&store),
logger,
);

let id = TestObjectId { id: [42u8; 4] };
let object = TestObject { id, data: [23u8; 3] };
let result = data_store
.mutate(&id, |existing| {
assert!(existing.is_none());
Some(object)
})
.await;
assert_eq!(Ok(Some(object)), result);

assert_eq!(Some(object), data_store.get(&id));
let store_key = id.encode_to_hex_str();
assert!(KVStore::read(&*store, &primary_namespace, &secondary_namespace, &store_key)
.await
.is_ok());
}

#[tokio::test]
async fn mutate_transforms_existing_entry() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// The closure sees the current entry and derives the new state from it.
let result = data_store
.mutate(&id, |existing| {
let mut new_object = *existing.unwrap();
new_object.data[0] += 1;
Some(new_object)
})
.await;
let expected = TestObject { id, data: [24u8, 23u8, 23u8] };
assert_eq!(Ok(Some(expected)), result);
assert_eq!(Some(expected), data_store.get(&id));
}

#[tokio::test]
async fn mutate_runs_the_closure_without_the_map_lock() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// Closures gate cross-store decisions on reads of other stores, which lock their own
// in-memory maps. Holding this store's map lock across the closure would order it
// before theirs and invite lock-order inversions, so the closure must run with the map
// lock released.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
assert!(data_store.objects.try_lock().is_ok());
None
})
.await;
assert_eq!(Ok(None), result);
}

#[tokio::test]
async fn mutate_persists_nothing_when_closure_declines() {
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

// Returning `None` must not attempt a write (the store fails all writes) nor touch memory.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
None
})
.await;
assert_eq!(Ok(None), result);
assert_eq!(Some(existing_object), data_store.get(&id));
}

#[tokio::test]
async fn mutate_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id: existing_id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

let changed = TestObject { id: existing_id, data: [24u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&existing_id, |_| Some(changed)).await
);
assert_eq!(Some(existing_object), data_store.get(&existing_id));

let new_id = TestObjectId { id: [55u8; 4] };
let new_object = TestObject { id: new_id, data: [34u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&new_id, |_| Some(new_object)).await
);
assert!(data_store.get(&new_id).is_none());
}

#[tokio::test]
async fn insert_or_update_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
Expand Down
75 changes: 75 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -242,4 +242,79 @@ mod tests {
"current txid must not remain in its own conflict list"
);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;

use crate::payment::store::PaymentDetailsUpdate;

let txid = test_txid(7);
let payment_id = PaymentId(txid.to_byte_array());

// A pending entry wallet sync has already mirrored a confirmation into (via
// `apply_funding_status_update_locked`) before classification ran.
let confirmed_details = PaymentDetails::new(
payment_id,
PaymentKind::Onchain {
txid,
status: ConfirmationStatus::Confirmed {
block_hash: BlockHash::from_byte_array([8u8; 32]),
height: 100,
timestamp: 1,
},
tx_type: None,
},
Some(2_000_000),
Some(999),
PaymentDirection::Outbound,
PaymentStatus::Pending,
);
let mirrored = PendingPaymentDetails::new(confirmed_details, Vec::new(), Vec::new());

// A fresh classification is always Unconfirmed and carries the candidate history; its
// figures are the active candidate's.
let fresh = pending_onchain_payment(payment_id, txid);
let candidates = vec![FundingTxCandidate {
txid,
amount_msat: fresh.amount_msat,
fee_paid_msat: fresh.fee_paid_msat,
}];

// The old fresh-insert path merged the full fresh record, downgrading the mirrored
// confirmation.
let mut downgraded = mirrored.clone();
let full_update =
PendingPaymentDetails::new(fresh.clone(), Vec::new(), candidates.clone()).to_update();
assert!(downgraded.update(full_update));
assert!(
matches!(
downgraded.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. }
),
"a full merge of a fresh classification downgrades a mirrored confirmation",
);

// The narrow classification update merges the candidates while preserving the
// confirmation state wallet sync owns. It names the confirmed txid, so its
// contribution-derived figures replace the mirrored wallet-view ones.
let mut merged = mirrored.clone();
let narrow_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)),
conflicting_txids: None,
candidates: candidates.clone(),
};
assert!(merged.update(narrow_update));
assert!(
matches!(
merged.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. }
),
"a narrow classification update must not downgrade a mirrored confirmation",
);
assert_eq!(merged.candidates, candidates);
assert_eq!(merged.details.amount_msat, Some(1_000));
assert_eq!(merged.details.fee_paid_msat, Some(100));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
157 changes: 157 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -81,6 +81,8 @@ where
Ok(updated)
}

/// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
/// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
pub(crate) async fn insert_or_update(&self, object: SO) -> Result<bool, Error> {
let _guard = self.mutation_lock.lock().await;

Expand DownExpand Up@@ -170,6 +172,36 @@ where
Ok(DataStoreUpdateResult::Updated)
}

/// Atomically transforms the entry for `id` through `f` and persists the result.
///
/// `f` receives the current entry (`None` when absent) and returns the new state to write;
/// returning `None` leaves the store untouched. The read, the closure, and the write share
/// one critical section of the mutation lock, so no concurrent writer can land in between —
/// unlike a separate [`Self::get`] followed by an insert or update.
///
/// The closure runs on a clone of the entry with the in-memory map lock released, so it may
/// freely read this store or others (reads see the pre-mutation state) without ordering map
/// locks against each other. Keep it cheap and non-blocking.
///
/// Returns the written object, or `None` when the closure declined to write.
pub(crate) async fn mutate<F: FnOnce(Option<&SO>) -> Option<SO>>(
&self, id: &SO::Id, f: F,
) -> Result<Option<SO>, Error> {
let _guard = self.mutation_lock.lock().await;

let current = self.objects.lock().expect("lock").get(id).cloned();
let new_object = match f(current.as_ref()) {
Some(new_object) => new_object,
None => return Ok(None),
};
debug_assert!(new_object.id() == *id, "mutate closure must not change the object's id");

self.persist(&new_object).await?;
let mut locked_objects = self.objects.lock().expect("lock");
locked_objects.insert(new_object.id(), new_object.clone());
Ok(Some(new_object))
}

/// Returns in-memory objects matching `f`.
///
/// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
Expand DownExpand Up@@ -403,6 +435,131 @@ mod tests {
assert_eq!(Ok(true), data_store.insert_or_update(new_iou_object).await);
}

#[tokio::test]
async fn mutate_inserts_when_absent() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let primary_namespace = "datastore_test_primary".to_string();
let secondary_namespace = "datastore_test_secondary".to_string();
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
Vec::new(),
primary_namespace.clone(),
secondary_namespace.clone(),
Arc::clone(&store),
logger,
);

let id = TestObjectId { id: [42u8; 4] };
let object = TestObject { id, data: [23u8; 3] };
let result = data_store
.mutate(&id, |existing| {
assert!(existing.is_none());
Some(object)
})
.await;
assert_eq!(Ok(Some(object)), result);

assert_eq!(Some(object), data_store.get(&id));
let store_key = id.encode_to_hex_str();
assert!(KVStore::read(&*store, &primary_namespace, &secondary_namespace, &store_key)
.await
.is_ok());
}

#[tokio::test]
async fn mutate_transforms_existing_entry() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// The closure sees the current entry and derives the new state from it.
let result = data_store
.mutate(&id, |existing| {
let mut new_object = *existing.unwrap();
new_object.data[0] += 1;
Some(new_object)
})
.await;
let expected = TestObject { id, data: [24u8, 23u8, 23u8] };
assert_eq!(Ok(Some(expected)), result);
assert_eq!(Some(expected), data_store.get(&id));
}

#[tokio::test]
async fn mutate_runs_the_closure_without_the_map_lock() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// Closures gate cross-store decisions on reads of other stores, which lock their own
// in-memory maps. Holding this store's map lock across the closure would order it
// before theirs and invite lock-order inversions, so the closure must run with the map
// lock released.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
assert!(data_store.objects.try_lock().is_ok());
None
})
.await;
assert_eq!(Ok(None), result);
}

#[tokio::test]
async fn mutate_persists_nothing_when_closure_declines() {
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

// Returning `None` must not attempt a write (the store fails all writes) nor touch memory.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
None
})
.await;
assert_eq!(Ok(None), result);
assert_eq!(Some(existing_object), data_store.get(&id));
}

#[tokio::test]
async fn mutate_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id: existing_id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

let changed = TestObject { id: existing_id, data: [24u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&existing_id, |_| Some(changed)).await
);
assert_eq!(Some(existing_object), data_store.get(&existing_id));

let new_id = TestObjectId { id: [55u8; 4] };
let new_object = TestObject { id: new_id, data: [34u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&new_id, |_| Some(new_object)).await
);
assert!(data_store.get(&new_id).is_none());
}

#[tokio::test]
async fn insert_or_update_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
Expand Down
75 changes: 75 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -242,4 +242,79 @@ mod tests {
"current txid must not remain in its own conflict list"
);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;

use crate::payment::store::PaymentDetailsUpdate;

let txid = test_txid(7);
let payment_id = PaymentId(txid.to_byte_array());

// A pending entry wallet sync has already mirrored a confirmation into (via
// `apply_funding_status_update_locked`) before classification ran.
let confirmed_details = PaymentDetails::new(
payment_id,
PaymentKind::Onchain {
txid,
status: ConfirmationStatus::Confirmed {
block_hash: BlockHash::from_byte_array([8u8; 32]),
height: 100,
timestamp: 1,
},
tx_type: None,
},
Some(2_000_000),
Some(999),
PaymentDirection::Outbound,
PaymentStatus::Pending,
);
let mirrored = PendingPaymentDetails::new(confirmed_details, Vec::new(), Vec::new());

// A fresh classification is always Unconfirmed and carries the candidate history; its
// figures are the active candidate's.
let fresh = pending_onchain_payment(payment_id, txid);
let candidates = vec![FundingTxCandidate {
txid,
amount_msat: fresh.amount_msat,
fee_paid_msat: fresh.fee_paid_msat,
}];

// The old fresh-insert path merged the full fresh record, downgrading the mirrored
// confirmation.
let mut downgraded = mirrored.clone();
let full_update =
PendingPaymentDetails::new(fresh.clone(), Vec::new(), candidates.clone()).to_update();
assert!(downgraded.update(full_update));
assert!(
matches!(
downgraded.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. }
),
"a full merge of a fresh classification downgrades a mirrored confirmation",
);

// The narrow classification update merges the candidates while preserving the
// confirmation state wallet sync owns. It names the confirmed txid, so its
// contribution-derived figures replace the mirrored wallet-view ones.
let mut merged = mirrored.clone();
let narrow_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)),
conflicting_txids: None,
candidates: candidates.clone(),
};
assert!(merged.update(narrow_update));
assert!(
matches!(
merged.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. }
),
"a narrow classification update must not downgrade a mirrored confirmation",
);
assert_eq!(merged.candidates, candidates);
assert_eq!(merged.details.amount_msat, Some(1_000));
assert_eq!(merged.details.fee_paid_msat, Some(100));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
157 changes: 157 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -81,6 +81,8 @@ where
Ok(updated)
}

/// Like [`Self::insert`], but when an entry with the object's id already exists, merges the
/// object's full update ([`StorableObject::to_update`]) into it instead of replacing it.
pub(crate) async fn insert_or_update(&self, object: SO) -> Result<bool, Error> {
let _guard = self.mutation_lock.lock().await;

Expand DownExpand Up@@ -170,6 +172,36 @@ where
Ok(DataStoreUpdateResult::Updated)
}

/// Atomically transforms the entry for `id` through `f` and persists the result.
///
/// `f` receives the current entry (`None` when absent) and returns the new state to write;
/// returning `None` leaves the store untouched. The read, the closure, and the write share
/// one critical section of the mutation lock, so no concurrent writer can land in between —
/// unlike a separate [`Self::get`] followed by an insert or update.
///
/// The closure runs on a clone of the entry with the in-memory map lock released, so it may
/// freely read this store or others (reads see the pre-mutation state) without ordering map
/// locks against each other. Keep it cheap and non-blocking.
///
/// Returns the written object, or `None` when the closure declined to write.
pub(crate) async fn mutate<F: FnOnce(Option<&SO>) -> Option<SO>>(
&self, id: &SO::Id, f: F,
) -> Result<Option<SO>, Error> {
let _guard = self.mutation_lock.lock().await;

let current = self.objects.lock().expect("lock").get(id).cloned();
let new_object = match f(current.as_ref()) {
Some(new_object) => new_object,
None => return Ok(None),
};
debug_assert!(new_object.id() == *id, "mutate closure must not change the object's id");

self.persist(&new_object).await?;
let mut locked_objects = self.objects.lock().expect("lock");
locked_objects.insert(new_object.id(), new_object.clone());
Ok(Some(new_object))
}

/// Returns in-memory objects matching `f`.
///
/// The async mutation lock serializes writers, but this synchronous reader cannot wait on it.
Expand DownExpand Up@@ -403,6 +435,131 @@ mod tests {
assert_eq!(Ok(true), data_store.insert_or_update(new_iou_object).await);
}

#[tokio::test]
async fn mutate_inserts_when_absent() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let primary_namespace = "datastore_test_primary".to_string();
let secondary_namespace = "datastore_test_secondary".to_string();
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
Vec::new(),
primary_namespace.clone(),
secondary_namespace.clone(),
Arc::clone(&store),
logger,
);

let id = TestObjectId { id: [42u8; 4] };
let object = TestObject { id, data: [23u8; 3] };
let result = data_store
.mutate(&id, |existing| {
assert!(existing.is_none());
Some(object)
})
.await;
assert_eq!(Ok(Some(object)), result);

assert_eq!(Some(object), data_store.get(&id));
let store_key = id.encode_to_hex_str();
assert!(KVStore::read(&*store, &primary_namespace, &secondary_namespace, &store_key)
.await
.is_ok());
}

#[tokio::test]
async fn mutate_transforms_existing_entry() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// The closure sees the current entry and derives the new state from it.
let result = data_store
.mutate(&id, |existing| {
let mut new_object = *existing.unwrap();
new_object.data[0] += 1;
Some(new_object)
})
.await;
let expected = TestObject { id, data: [24u8, 23u8, 23u8] };
assert_eq!(Ok(Some(expected)), result);
assert_eq!(Some(expected), data_store.get(&id));
}

#[tokio::test]
async fn mutate_runs_the_closure_without_the_map_lock() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
let logger = Arc::new(TestLogger::new());
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store: DataStore<TestObject, Arc<TestLogger>> = DataStore::new(
vec![existing_object],
"datastore_test_primary".to_string(),
"datastore_test_secondary".to_string(),
store,
logger,
);

// Closures gate cross-store decisions on reads of other stores, which lock their own
// in-memory maps. Holding this store's map lock across the closure would order it
// before theirs and invite lock-order inversions, so the closure must run with the map
// lock released.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
assert!(data_store.objects.try_lock().is_ok());
None
})
.await;
assert_eq!(Ok(None), result);
}

#[tokio::test]
async fn mutate_persists_nothing_when_closure_declines() {
let id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

// Returning `None` must not attempt a write (the store fails all writes) nor touch memory.
let result = data_store
.mutate(&id, |existing| {
assert_eq!(Some(&existing_object), existing);
None
})
.await;
assert_eq!(Ok(None), result);
assert_eq!(Some(existing_object), data_store.get(&id));
}

#[tokio::test]
async fn mutate_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
let existing_object = TestObject { id: existing_id, data: [23u8; 3] };
let data_store = new_failing_data_store(vec![existing_object]);

let changed = TestObject { id: existing_id, data: [24u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&existing_id, |_| Some(changed)).await
);
assert_eq!(Some(existing_object), data_store.get(&existing_id));

let new_id = TestObjectId { id: [55u8; 4] };
let new_object = TestObject { id: new_id, data: [34u8; 3] };
assert_eq!(
Err(Error::PersistenceFailed),
data_store.mutate(&new_id, |_| Some(new_object)).await
);
assert!(data_store.get(&new_id).is_none());
}

#[tokio::test]
async fn insert_or_update_does_not_mutate_memory_if_persist_fails() {
let existing_id = TestObjectId { id: [42u8; 4] };
Expand Down
75 changes: 75 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -242,4 +242,79 @@ mod tests {
"current txid must not remain in its own conflict list"
);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;

use crate::payment::store::PaymentDetailsUpdate;

let txid = test_txid(7);
let payment_id = PaymentId(txid.to_byte_array());

// A pending entry wallet sync has already mirrored a confirmation into (via
// `apply_funding_status_update_locked`) before classification ran.
let confirmed_details = PaymentDetails::new(
payment_id,
PaymentKind::Onchain {
txid,
status: ConfirmationStatus::Confirmed {
block_hash: BlockHash::from_byte_array([8u8; 32]),
height: 100,
timestamp: 1,
},
tx_type: None,
},
Some(2_000_000),
Some(999),
PaymentDirection::Outbound,
PaymentStatus::Pending,
);
let mirrored = PendingPaymentDetails::new(confirmed_details, Vec::new(), Vec::new());

// A fresh classification is always Unconfirmed and carries the candidate history; its
// figures are the active candidate's.
let fresh = pending_onchain_payment(payment_id, txid);
let candidates = vec![FundingTxCandidate {
txid,
amount_msat: fresh.amount_msat,
fee_paid_msat: fresh.fee_paid_msat,
}];

// The old fresh-insert path merged the full fresh record, downgrading the mirrored
// confirmation.
let mut downgraded = mirrored.clone();
let full_update =
PendingPaymentDetails::new(fresh.clone(), Vec::new(), candidates.clone()).to_update();
assert!(downgraded.update(full_update));
assert!(
matches!(
downgraded.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. }
),
"a full merge of a fresh classification downgrades a mirrored confirmation",
);

// The narrow classification update merges the candidates while preserving the
// confirmation state wallet sync owns. It names the confirmed txid, so its
// contribution-derived figures replace the mirrored wallet-view ones.
let mut merged = mirrored.clone();
let narrow_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)),
conflicting_txids: None,
candidates: candidates.clone(),
};
assert!(merged.update(narrow_update));
assert!(
matches!(
merged.details.kind,
PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. }
),
"a narrow classification update must not downgrade a mirrored confirmation",
);
assert_eq!(merged.candidates, candidates);
assert_eq!(merged.details.amount_msat, Some(1_000));
assert_eq!(merged.details.fee_paid_msat, Some(100));
}
}
Loading
Loading