Skip to content

Commit cfe63f0

Browse files
committed
fix(sync): update array sync module for new wire types and Encryption API
- Add `producer_id`, `epoch`, and `seq` fields (all zero) to `ArrayDeltaMsg` constructions in tests for inbound delta handling. - Add `applied_seq` and `AckStatus::Applied` to `ArrayAckMsg` construction in `ack_sender`. - Update internal test calls to `PagedbStorageDefault::open` to pass `Encryption::Plaintext` to match the new storage open signature.
1 parent 1499bc8 commit cfe63f0

6 files changed

Lines changed: 81 additions & 8 deletions

File tree

nodedb-lite/src/sync/array/ack_sender.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,8 @@ async fn send_acks<S: StorageEngine>(
8585
array: array.clone(),
8686
replica_id: replica_id.as_u64(),
8787
ack_hlc_bytes: ack_hlc.to_bytes(),
88+
applied_seq: 0,
89+
status: nodedb_types::sync::wire::AckStatus::Applied,
8890
};
8991

9092
let frame = match nodedb_types::sync::wire::SyncFrame::try_encode(

nodedb-lite/src/sync/array/catchup.rs

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -234,14 +234,28 @@ mod tests {
234234
let target_hlc = hlc(42_000);
235235

236236
{
237-
let storage = Arc::new(PagedbStorageDefault::open(&path).await.unwrap());
237+
let storage = Arc::new(
238+
PagedbStorageDefault::open(
239+
&path,
240+
crate::storage::encryption::Encryption::Plaintext,
241+
)
242+
.await
243+
.unwrap(),
244+
);
238245
let tracker = CatchupTracker::load(Arc::clone(&storage)).await.unwrap();
239246
tracker.record("arr", target_hlc).await.unwrap();
240247
assert_eq!(tracker.last_seen("arr"), target_hlc);
241248
}
242249

243250
{
244-
let storage = Arc::new(PagedbStorageDefault::open(&path).await.unwrap());
251+
let storage = Arc::new(
252+
PagedbStorageDefault::open(
253+
&path,
254+
crate::storage::encryption::Encryption::Plaintext,
255+
)
256+
.await
257+
.unwrap(),
258+
);
245259
let tracker = CatchupTracker::load(storage).await.unwrap();
246260
assert_eq!(
247261
tracker.last_seen("arr"),

nodedb-lite/src/sync/array/inbound/delta.rs

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,9 @@ mod tests {
9595
let msg = ArrayDeltaMsg {
9696
array: "arr".into(),
9797
op_payload: payload,
98+
producer_id: 0,
99+
epoch: 0,
100+
seq: 0,
98101
};
99102
let outcome = inbound.handle_delta(&msg).unwrap();
100103
assert_eq!(outcome, InboundOutcome::Applied);
@@ -133,6 +136,9 @@ mod tests {
133136
let msg = ArrayDeltaMsg {
134137
array: "arr".into(),
135138
op_payload: payload.clone(),
139+
producer_id: 0,
140+
epoch: 0,
141+
seq: 0,
136142
};
137143

138144
// First application — should be Applied.
@@ -143,6 +149,9 @@ mod tests {
143149
let msg2 = ArrayDeltaMsg {
144150
array: "arr".into(),
145151
op_payload: payload,
152+
producer_id: 0,
153+
epoch: 0,
154+
seq: 0,
146155
};
147156
let o2 = inbound.handle_delta(&msg2).unwrap();
148157
assert_eq!(o2, InboundOutcome::Idempotent);
@@ -157,6 +166,9 @@ mod tests {
157166
let msg = ArrayDeltaMsg {
158167
array: "unknown_arr".into(),
159168
op_payload: payload,
169+
producer_id: 0,
170+
epoch: 0,
171+
seq: 0,
160172
};
161173
let outcome = inbound.handle_delta(&msg).unwrap();
162174
assert!(
@@ -201,6 +213,9 @@ mod tests {
201213
let msg = ArrayDeltaMsg {
202214
array: "arr".into(),
203215
op_payload: payload,
216+
producer_id: 0,
217+
epoch: 0,
218+
seq: 0,
204219
};
205220
let outcome = inbound.handle_delta(&msg).unwrap();
206221
assert!(

nodedb-lite/src/sync/array/op_log_store.rs

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -446,15 +446,29 @@ mod tests {
446446
let path = dir.path().join("op_log_test.pagedb");
447447

448448
{
449-
let storage = Arc::new(PagedbStorageDefault::open(&path).await.unwrap());
449+
let storage = Arc::new(
450+
PagedbStorageDefault::open(
451+
&path,
452+
crate::storage::encryption::Encryption::Plaintext,
453+
)
454+
.await
455+
.unwrap(),
456+
);
450457
let log = KvOpLogStore::new(Arc::clone(&storage));
451458
log.append(&make_op("arr", 10)).unwrap();
452459
log.append(&make_op("arr", 20)).unwrap();
453460
}
454461

455462
// Reopen the same file.
456463
{
457-
let storage = Arc::new(PagedbStorageDefault::open(&path).await.unwrap());
464+
let storage = Arc::new(
465+
PagedbStorageDefault::open(
466+
&path,
467+
crate::storage::encryption::Encryption::Plaintext,
468+
)
469+
.await
470+
.unwrap(),
471+
);
458472
let log = KvOpLogStore::new(storage);
459473
assert_eq!(log.len().unwrap(), 2);
460474
let ops: Vec<_> = log

nodedb-lite/src/sync/array/pending.rs

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -266,14 +266,28 @@ mod tests {
266266
let path = dir.path().join("pending_test.pagedb");
267267

268268
{
269-
let storage = Arc::new(PagedbStorageDefault::open(&path).await.unwrap());
269+
let storage = Arc::new(
270+
PagedbStorageDefault::open(
271+
&path,
272+
crate::storage::encryption::Encryption::Plaintext,
273+
)
274+
.await
275+
.unwrap(),
276+
);
270277
let q = PendingQueue::new(storage);
271278
q.enqueue(&make_op(5)).await.unwrap();
272279
q.enqueue(&make_op(15)).await.unwrap();
273280
}
274281

275282
{
276-
let storage = Arc::new(PagedbStorageDefault::open(&path).await.unwrap());
283+
let storage = Arc::new(
284+
PagedbStorageDefault::open(
285+
&path,
286+
crate::storage::encryption::Encryption::Plaintext,
287+
)
288+
.await
289+
.unwrap(),
290+
);
277291
let q = PendingQueue::new(storage);
278292
assert_eq!(q.len().await.unwrap(), 2);
279293
let ops = q.drain_batch(usize::MAX).await.unwrap();

nodedb-lite/src/sync/array/schema_registry.rs

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -302,14 +302,28 @@ mod tests {
302302

303303
let schema_hlc;
304304
{
305-
let storage = Arc::new(PagedbStorageDefault::open(&path).await.unwrap());
305+
let storage = Arc::new(
306+
PagedbStorageDefault::open(
307+
&path,
308+
crate::storage::encryption::Encryption::Plaintext,
309+
)
310+
.await
311+
.unwrap(),
312+
);
306313
let replica = Arc::new(ReplicaState::load_or_init(&*storage).await.unwrap());
307314
let reg = SchemaRegistry::new(Arc::clone(&storage), Arc::clone(&replica));
308315
schema_hlc = reg.put_schema("arr", &simple_schema("arr")).await.unwrap();
309316
}
310317

311318
{
312-
let storage = Arc::new(PagedbStorageDefault::open(&path).await.unwrap());
319+
let storage = Arc::new(
320+
PagedbStorageDefault::open(
321+
&path,
322+
crate::storage::encryption::Encryption::Plaintext,
323+
)
324+
.await
325+
.unwrap(),
326+
);
313327
let replica = Arc::new(ReplicaState::load_or_init(&*storage).await.unwrap());
314328
let reg = SchemaRegistry::load(Arc::clone(&storage), replica)
315329
.await

0 commit comments

Comments
 (0)