Skip to content

Commit fed80c6

Browse files
committed
[mobile][desktop] Preserve sync progress across internal retries
1 parent 78bab6e commit fed80c6

4 files changed

Lines changed: 233 additions & 3 deletions

File tree

integration_test/regtest_sync_startup_stall_recovery_test.dart

Lines changed: 96 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,12 @@
11
import 'package:flutter/widgets.dart';
2+
import 'package:flutter_riverpod/flutter_riverpod.dart';
23
import 'package:flutter_test/flutter_test.dart';
34
import 'package:integration_test/integration_test.dart';
45
import 'package:zcash_wallet/app.dart';
56
import 'package:zcash_wallet/src/core/config/rpc_endpoint_config.dart';
67
import 'package:zcash_wallet/src/core/storage/app_secure_store.dart';
78
import 'package:zcash_wallet/src/core/storage/wallet_paths.dart';
9+
import 'package:zcash_wallet/src/providers/sync_provider.dart';
810
import 'package:zcash_wallet/src/rust/api/sync.dart' as rust_sync;
911

1012
import 'support/desktop_regtest_flow.dart';
@@ -17,6 +19,14 @@ const _networkFailure =
1719
const _endpointFailure =
1820
'Cannot reach the configured Zcash endpoint. Check your endpoint settings.';
1921
const _genericFailure = 'Sync failed. Retry sync to continue.';
22+
const _proxyListenPort = int.fromEnvironment(
23+
'ZCASH_E2E_PROXY_PORT',
24+
defaultValue: 19068,
25+
);
26+
const _lightwalletdTargetPort = int.fromEnvironment(
27+
'ZCASH_E2E_LIGHTWALLETD_PORT',
28+
defaultValue: 9067,
29+
);
2030

2131
void main() {
2232
IntegrationTestWidgetsFlutterBinding.ensureInitialized();
@@ -29,7 +39,11 @@ void main() {
2939
'recovers from a stalled startup UTXO stream instead of staying at zero percent',
3040
(tester) async {
3141
await cleanupDesktopRegtestWallet();
32-
final proxy = RegtestLightwalletdProxy(log: e2eLog);
42+
final proxy = RegtestLightwalletdProxy(
43+
listenPort: _proxyListenPort,
44+
targetPort: _lightwalletdTargetPort,
45+
log: e2eLog,
46+
);
3347
await proxy.start();
3448
// Keep later scanning minimal while retaining a persisted completed tip.
3549
proxy.setSlowHeight(1);
@@ -117,4 +131,85 @@ void main() {
117131
},
118132
timeout: const Timeout(Duration(minutes: 3)),
119133
);
134+
135+
testWidgets(
136+
'keeps logical sync progress monotonic across an internal Rust retry',
137+
(tester) async {
138+
await cleanupDesktopRegtestWallet();
139+
final proxy = RegtestLightwalletdProxy(
140+
listenPort: _proxyListenPort,
141+
targetPort: _lightwalletdTargetPort,
142+
log: e2eLog,
143+
);
144+
proxy.failBlockRangeCalls(const [2, 3]);
145+
await proxy.start();
146+
addTearDown(() async {
147+
try {
148+
await cleanupDesktopRegtestWallet();
149+
} finally {
150+
await proxy.stop();
151+
}
152+
});
153+
154+
final storage = AppSecureStore.instance;
155+
await storage.writePlain(kRpcEndpointUrlKey, proxy.url);
156+
await storage.writePlain(
157+
kRpcEndpointPresetKey,
158+
kCustomRpcEndpointPresetId,
159+
);
160+
161+
await tester.pumpWidget(await buildBootstrappedZcashWalletApp());
162+
final container = ProviderScope.containerOf(
163+
tester.element(
164+
find.byKey(const ValueKey('welcome_import_wallet_button')),
165+
),
166+
);
167+
final observedPercentages = <double>[];
168+
final subscription = container.listen(syncProvider, (previous, next) {
169+
final sync = next.value;
170+
if (sync != null &&
171+
(sync.isSyncing || sync.isSyncComplete) &&
172+
(observedPercentages.isEmpty ||
173+
(observedPercentages.last - sync.percentage).abs() >
174+
0.000001)) {
175+
observedPercentages.add(sync.percentage);
176+
}
177+
});
178+
addTearDown(subscription.close);
179+
180+
await importDesktopRegtestWallet(tester);
181+
await pumpUntil(
182+
tester,
183+
() {
184+
final sync = container.read(syncProvider).value;
185+
return proxy.failedBlockRangeCount == 2 &&
186+
!rust_sync.isSyncRunning() &&
187+
sync?.isSyncComplete == true;
188+
},
189+
description: 'completed sync after an internal block-range retry',
190+
timeout: const Duration(minutes: 2),
191+
);
192+
193+
expect(proxy.failedBlockRangeCount, 2);
194+
expect(proxy.blockRangeCallCount, greaterThanOrEqualTo(4));
195+
expect(
196+
observedPercentages,
197+
contains(allOf(greaterThan(0.0), lessThan(1.0))),
198+
);
199+
for (var i = 1; i < observedPercentages.length; i++) {
200+
expect(
201+
observedPercentages[i] + 0.000001,
202+
greaterThanOrEqualTo(observedPercentages[i - 1]),
203+
reason:
204+
'A retry inside one logical Rust sync must not move progress '
205+
'backward: $observedPercentages',
206+
);
207+
}
208+
e2eLog(
209+
'internal retry preserved monotonic progress: '
210+
'$observedPercentages',
211+
);
212+
},
213+
timeout: const Timeout(Duration(minutes: 3)),
214+
);
120215
}

integration_test/support/regtest_lightwalletd_proxy.dart

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,13 +36,18 @@ class RegtestLightwalletdProxy
3636
int? _slowHeight;
3737
int _sendTransactionFailuresRemaining = 0;
3838
int _failedSendTransactionCount = 0;
39+
final Set<int> _blockRangeFailureCalls = <int>{};
40+
int _blockRangeCallCount = 0;
41+
int _failedBlockRangeCount = 0;
3942
bool _stallNextAddressUtxosAfterHeaders = false;
4043
Completer<void>? _addressUtxosStallRelease;
4144
int _addressUtxosStreamCallCount = 0;
4245
bool _serveEmptyGenesisTreeState = false;
4346

4447
String get url => 'http://127.0.0.1:$listenPort';
4548
int get failedSendTransactionCount => _failedSendTransactionCount;
49+
int get blockRangeCallCount => _blockRangeCallCount;
50+
int get failedBlockRangeCount => _failedBlockRangeCount;
4651
int get addressUtxosStreamCallCount => _addressUtxosStreamCallCount;
4752

4853
Future<void> start() async {
@@ -85,6 +90,22 @@ class RegtestLightwalletdProxy
8590
_log('primary proxy will fail next $count SendTransaction call(s)');
8691
}
8792

93+
void failBlockRangeCalls(Iterable<int> callNumbers) {
94+
final calls = callNumbers.toSet();
95+
if (calls.any((call) => call <= _blockRangeCallCount)) {
96+
throw ArgumentError.value(
97+
callNumbers,
98+
'callNumbers',
99+
'must contain only future GetBlockRange call numbers',
100+
);
101+
}
102+
_blockRangeFailureCalls.addAll(calls);
103+
_log(
104+
'primary proxy will fail GetBlockRange call(s) '
105+
'${calls.toList()..sort()}',
106+
);
107+
}
108+
88109
void stallNextAddressUtxosStreamAfterHeaders() {
89110
_stallNextAddressUtxosAfterHeaders = true;
90111
_addressUtxosStallRelease = Completer<void>();
@@ -166,6 +187,17 @@ class RegtestLightwalletdProxy
166187
service.BlockRange request,
167188
) {
168189
_throwIfDown();
190+
_blockRangeCallCount += 1;
191+
if (_blockRangeFailureCalls.remove(_blockRangeCallCount)) {
192+
_failedBlockRangeCount += 1;
193+
_log(
194+
'primary proxy forced GetBlockRange failure '
195+
'(call=$_blockRangeCallCount, failed=$_failedBlockRangeCount)',
196+
);
197+
return Stream.error(
198+
grpc.GrpcError.unavailable('forced GetBlockRange failure'),
199+
);
200+
}
169201
return _client.getBlockRange(request);
170202
}
171203

rust/src/wallet/sync_engine/mod.rs

Lines changed: 103 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
use std::collections::{BTreeSet, HashSet};
22
use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
3-
use std::sync::Arc;
3+
use std::sync::{Arc, Mutex};
44

55
use nonempty::NonEmpty;
66
use rusqlite::{params, OptionalExtension};
@@ -76,6 +76,39 @@ pub struct SyncProgressEvent {
7676
pub phase: String,
7777
}
7878

79+
/// Keeps retry attempts inside one `run_sync_inner` call on a single logical
80+
/// progress scale.
81+
///
82+
/// `run_sync_impl` recalculates its work denominator from the DB at the start
83+
/// of every attempt. That is correct for resuming from the persisted scan
84+
/// frontier, but the attempt-local percentage would otherwise jump backward to
85+
/// zero. A fresh `RetryProgressMapper` is created for every logical sync call,
86+
/// so a genuinely new session (including account-growth catch-up) still starts
87+
/// from its own zero-based progress.
88+
#[derive(Debug, Default)]
89+
struct RetryProgressMapper {
90+
retry_floor: f64,
91+
last_emitted_percentage: f64,
92+
}
93+
94+
impl RetryProgressMapper {
95+
fn begin_retry(&mut self) {
96+
self.retry_floor = self.last_emitted_percentage.clamp(0.0, 1.0);
97+
}
98+
99+
fn map_percentage(&self, attempt_percentage: f64) -> f64 {
100+
let attempt_percentage = attempt_percentage.clamp(0.0, 1.0);
101+
self.retry_floor + ((1.0 - self.retry_floor) * attempt_percentage)
102+
}
103+
104+
fn map_event(&mut self, mut event: SyncProgressEvent) -> SyncProgressEvent {
105+
event.percentage = self.map_percentage(event.percentage);
106+
event.display_target_percentage = self.map_percentage(event.display_target_percentage);
107+
self.last_emitted_percentage = event.percentage;
108+
event
109+
}
110+
}
111+
79112
#[cfg(any(target_os = "macos", target_os = "windows", target_os = "linux"))]
80113
const BATCH_SIZE_FOREGROUND: u32 = 2000;
81114
#[cfg(not(any(target_os = "macos", target_os = "windows", target_os = "linux")))]
@@ -1556,6 +1589,7 @@ pub async fn run_sync_inner(
15561589
) -> Result<(), String> {
15571590
const MAX_RETRIES: u32 = 3;
15581591
let mut last_err = String::new();
1592+
let retry_progress = Mutex::new(RetryProgressMapper::default());
15591593
*SYNC_START.lock().unwrap() = Some(std::time::Instant::now());
15601594

15611595
for attempt in 0..=MAX_RETRIES {
@@ -1582,8 +1616,19 @@ pub async fn run_sync_inner(
15821616
return Ok(());
15831617
}
15841618
}
1619+
retry_progress
1620+
.lock()
1621+
.unwrap_or_else(|poisoned| poisoned.into_inner())
1622+
.begin_retry();
15851623
}
15861624

1625+
let mapped_progress_fn = |event| {
1626+
let mapped = retry_progress
1627+
.lock()
1628+
.unwrap_or_else(|poisoned| poisoned.into_inner())
1629+
.map_event(event);
1630+
progress_fn(mapped);
1631+
};
15871632
match run_sync_impl(
15881633
db_data_path,
15891634
lightwalletd_url,
@@ -1592,7 +1637,7 @@ pub async fn run_sync_inner(
15921637
running_mode,
15931638
desired_mode,
15941639
allow_resubmit,
1595-
&progress_fn,
1640+
&mapped_progress_fn,
15961641
)
15971642
.await
15981643
{
@@ -2904,6 +2949,20 @@ mod tests {
29042949
);
29052950
}
29062951

2952+
fn progress_event(percentage: f64, display_target_percentage: f64) -> SyncProgressEvent {
2953+
SyncProgressEvent {
2954+
scanned_height: 0,
2955+
chain_tip_height: 0,
2956+
percentage,
2957+
display_target_percentage,
2958+
display_target_blocks: 0,
2959+
is_syncing: true,
2960+
is_complete: false,
2961+
has_new_tx: false,
2962+
phase: "scan".into(),
2963+
}
2964+
}
2965+
29072966
#[test]
29082967
fn work_progress_matches_remaining_block_ratio() {
29092968
let mode = ProgressDisplayMode::Work;
@@ -2915,6 +2974,48 @@ mod tests {
29152974
);
29162975
}
29172976

2977+
#[test]
2978+
fn retry_progress_resumes_from_previous_attempt_floor() {
2979+
let mut mapper = RetryProgressMapper::default();
2980+
let first_attempt = mapper.map_event(progress_event(225.0 / 1_200.0, 250.0 / 1_200.0));
2981+
assert_pct(first_attempt.percentage, 225.0 / 1_200.0);
2982+
2983+
mapper.begin_retry();
2984+
let retry_start = mapper.map_event(progress_event(0.0, 25.0 / 975.0));
2985+
assert_pct(retry_start.percentage, 225.0 / 1_200.0);
2986+
assert_pct(retry_start.display_target_percentage, 250.0 / 1_200.0);
2987+
2988+
let retry_batch_complete = mapper.map_event(progress_event(25.0 / 975.0, 50.0 / 975.0));
2989+
assert_pct(retry_batch_complete.percentage, 250.0 / 1_200.0);
2990+
assert_pct(
2991+
retry_batch_complete.display_target_percentage,
2992+
275.0 / 1_200.0,
2993+
);
2994+
}
2995+
2996+
#[test]
2997+
fn new_logical_sync_does_not_inherit_previous_retry_floor() {
2998+
let mut previous_sync = RetryProgressMapper::default();
2999+
previous_sync.map_event(progress_event(0.5, 0.6));
3000+
3001+
let mut new_sync = RetryProgressMapper::default();
3002+
let event = new_sync.map_event(progress_event(0.0, 0.1));
3003+
3004+
assert_pct(event.percentage, 0.0);
3005+
assert_pct(event.display_target_percentage, 0.1);
3006+
}
3007+
3008+
#[test]
3009+
fn first_attempt_preserves_existing_phase_progress_semantics() {
3010+
let mut mapper = RetryProgressMapper::default();
3011+
3012+
assert_pct(mapper.map_event(progress_event(0.99, 1.0)).percentage, 0.99);
3013+
assert_pct(
3014+
mapper.map_event(progress_event(0.95, 0.975)).percentage,
3015+
0.95,
3016+
);
3017+
}
3018+
29183019
#[test]
29193020
fn chain_window_progress_uses_session_window_not_absolute_height() {
29203021
let mode = ProgressDisplayMode::ChainWindow {

scripts/e2e/flutter-macos-regtest-sync-startup-stall-recovery.sh

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,8 @@ set -euo pipefail
44
ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)"
55
FLUTTER_DEVICE="${FLUTTER_DEVICE:-macos}"
66
RESET_REGTEST="${RESET_REGTEST:-1}"
7+
export ZCASH_E2E_SYNC_BATCH_SIZE="${ZCASH_E2E_SYNC_BATCH_SIZE:-25}"
8+
export ZCASH_E2E_SYNC_BATCH_DELAY_MS="${ZCASH_E2E_SYNC_BATCH_DELAY_MS:-300}"
79

810
require_cmd() {
911
if ! command -v "$1" >/dev/null 2>&1; then

0 commit comments

Comments
 (0)