49 lines
1.8 KiB
Dart
49 lines
1.8 KiB
Dart
import 'cloud/cloud_storage_provider.dart' show CloudLockFile;
|
|
|
|
/// Ticket-based mutex using timestamped lock files as the coordination
|
|
/// point: create a lock named after our own timestamp, then wait until no
|
|
/// *other* lock older than ours remains — reaping ("deleting") any it
|
|
/// finds older than [staleAge] along the way, as a backstop against a
|
|
/// device that crashed/went offline before releasing its own lock.
|
|
///
|
|
/// Pure orchestration over injected operations (rather than a concrete
|
|
/// cloud storage dependency) so the state machine is unit-testable without
|
|
/// a real network connection, and works identically no matter which
|
|
/// [CloudStorageProvider] it's wired to.
|
|
Future<String> acquireLock({
|
|
required String username,
|
|
required Future<String> Function(String lockName) createLock,
|
|
required Future<List<CloudLockFile>> Function() listLocks,
|
|
required Future<void> Function(String lockId) deleteLock,
|
|
DateTime Function()? nowUtc,
|
|
Future<void> Function(Duration)? delay,
|
|
Duration staleAge = const Duration(minutes: 10),
|
|
Duration pollInterval = const Duration(seconds: 1),
|
|
}) async {
|
|
final now = nowUtc ?? () => DateTime.now().toUtc();
|
|
final wait = delay ?? Future.delayed;
|
|
|
|
final epochMillis = now().millisecondsSinceEpoch;
|
|
final lockName = '$username-$epochMillis.lock';
|
|
final lockId = await createLock(lockName);
|
|
|
|
while (true) {
|
|
final locks = await listLocks();
|
|
final others = locks.where((l) => l.id != lockId);
|
|
final olderOthers =
|
|
others.where((l) => l.createdAtUtc.millisecondsSinceEpoch < epochMillis).toList();
|
|
|
|
if (olderOthers.isEmpty) break;
|
|
|
|
final current = now();
|
|
for (final stale in olderOthers) {
|
|
if (current.difference(stale.createdAtUtc) > staleAge) {
|
|
await deleteLock(stale.id);
|
|
}
|
|
}
|
|
|
|
await wait(pollInterval);
|
|
}
|
|
|
|
return lockId;
|
|
}
|