Files
nearle_pos/lib/presentation/sync/providers/sync_controller.dart
2026-08-06 19:29:23 +05:30

291 lines
9.7 KiB
Dart

import 'package:flutter_riverpod/flutter_riverpod.dart';
import '../../../app/providers.dart';
import '../../../core/constants/app_constants.dart';
import '../../../data/remote/mqtt_order_transport.dart';
import '../../../data/sync/health_reporter.dart';
import '../../../data/sync/presence_reporter.dart';
import '../../modules/providers/printer_settings.dart';
import '../../../domain/entities/shift_report.dart';
import '../../../domain/entities/sync_event.dart';
import '../../../domain/repositories/sync_repository.dart';
import '../../auth/providers/auth_controller.dart';
import '../../pos/providers/catalog_providers.dart';
/// Bumped after every import so catalogue-backed providers refetch.
final catalogueVersionProvider = StateProvider<int>((ref) => 0);
/// Bumped after every sale or sync so order-backed providers refetch.
final orderVersionProvider = StateProvider<int>((ref) => 0);
/// Whether products exist on this terminal. Billing is gated on it.
final catalogueReadyProvider = Provider<bool>((ref) {
ref.watch(catalogueVersionProvider);
return ref.watch(syncRepositoryProvider).hasCatalogue;
});
final lastImportAtProvider = Provider<DateTime?>((ref) {
ref.watch(catalogueVersionProvider);
return ref.watch(syncRepositoryProvider).lastImportAt;
});
// ------------------------------------------------------- Morning: import
sealed class ImportState {
const ImportState();
}
class ImportIdle extends ImportState {
const ImportIdle();
}
class ImportRunning extends ImportState {
const ImportRunning(this.progress, this.stage);
final double progress;
final String stage;
}
class ImportDone extends ImportState {
const ImportDone(this.event);
final SyncEvent event;
}
class ImportFailed extends ImportState {
const ImportFailed(this.message);
final String message;
}
class CatalogueImportController extends StateNotifier<ImportState> {
CatalogueImportController(this._ref) : super(const ImportIdle());
final Ref _ref;
Future<bool> run() async {
if (state is ImportRunning) return false;
state = const ImportRunning(0, 'Starting…');
final event = await _ref.read(syncRepositoryProvider).importCatalogue(
onProgress: (progress, stage) {
if (mounted) state = ImportRunning(progress, stage);
},
);
if (event.status == SyncStatus.synced) {
_ref.read(catalogueVersionProvider.notifier).state++;
_ref.invalidate(allProductsProvider);
_ref.invalidate(visibleProductsProvider);
_ref.invalidate(categoryCountsProvider);
_ref.invalidate(lowStockProductsProvider);
state = ImportDone(event);
return true;
}
state = ImportFailed(event.error ?? 'Import failed.');
return false;
}
/// Drops a stale success or failure banner.
void reset() => state = const ImportIdle();
}
final catalogueImportProvider =
StateNotifierProvider<CatalogueImportController, ImportState>(
(ref) => CatalogueImportController(ref),
);
// ---------------------------------------------------- Business hours: read
/// Bills still held on this terminal at sync_status = 0.
final unsyncedCountProvider = FutureProvider<int>((ref) {
ref.watch(orderVersionProvider);
return ref.watch(syncRepositoryProvider).unsyncedCount();
});
/// Everything this terminal traded today, across every operator.
final todayReportProvider = FutureProvider<ShiftReport>((ref) {
ref.watch(orderVersionProvider);
final session = ref.watch(cashierSessionProvider);
final user = ref.watch(currentUserProvider);
return ref.watch(syncRepositoryProvider).todayReport(
terminalId: session.terminalId,
cashierName: user?.name ?? session.name,
);
});
/// Only the bills the signed-in operator rang.
///
/// This is the figure a cashier counts their drawer against at the end of a
/// shift, so it must not include anyone else's sales.
final myShiftReportProvider = FutureProvider<ShiftReport>((ref) {
ref.watch(orderVersionProvider);
final session = ref.watch(cashierSessionProvider);
final user = ref.watch(currentUserProvider);
return ref.watch(syncRepositoryProvider).todayReport(
terminalId: session.terminalId,
cashierName: user?.name ?? session.name,
scopeToCashier: true,
);
});
/// Per-order sync state for the events log.
final orderSyncRowsProvider = FutureProvider<List<OrderSyncRow>>((ref) {
ref.watch(orderVersionProvider);
return ref.watch(syncRepositoryProvider).orderSyncRows();
});
final syncEventsProvider = Provider<List<SyncEvent>>((ref) {
ref.watch(orderVersionProvider);
ref.watch(catalogueVersionProvider);
return ref.watch(syncRepositoryProvider).events;
});
// ------------------------------------------------------ End of day: upload
sealed class OrderSyncState {
const OrderSyncState();
}
class SyncIdle extends OrderSyncState {
const SyncIdle();
}
class SyncRunning extends OrderSyncState {
const SyncRunning(this.progress, this.stage);
final double progress;
final String stage;
}
class SyncFinished extends OrderSyncState {
const SyncFinished(this.outcome);
final SyncOutcome outcome;
}
class OrderSyncController extends StateNotifier<OrderSyncState> {
OrderSyncController(this._ref) : super(const SyncIdle());
final Ref _ref;
bool get isRunning => state is SyncRunning;
/// Uploads every bill at sync_status = 0 and flips the accepted ones to 1.
///
/// Goes through the engine rather than straight to the repository, so a
/// cashier pressing sync while a background drain is already mid-flight
/// joins it instead of starting a second pass over the same rows. It also
/// clears a halt: pressing the button is how you retry after the back office
/// has been fixed.
Future<SyncOutcome> run() async {
if (isRunning) {
return const SyncOutcome(attempted: 0, uploaded: 0);
}
state = const SyncRunning(0, 'Starting…');
final outcome = await _ref.read(syncEngineProvider).syncNow(
onProgress: (progress, stage) {
if (mounted) state = SyncRunning(progress, stage);
},
);
_ref.read(orderVersionProvider.notifier).state++;
if (mounted) state = SyncFinished(outcome);
return outcome;
}
void reset() => state = const SyncIdle();
}
final orderSyncProvider =
StateNotifierProvider<OrderSyncController, OrderSyncState>(
(ref) => OrderSyncController(ref),
);
// ------------------------------------------------------- Background drain
/// Brings the queue-and-drain machinery up, once, when the shell mounts.
///
/// Deliberately not gated on sign-in: a terminal that boots holding yesterday's
/// bills should be emptying its queue before anyone reaches the till.
///
/// Overridden to a no-op in widget tests, which have no network stack and
/// cannot drive real disk I/O on a fake clock.
final syncBootstrapProvider = FutureProvider<void>((ref) async {
// Restore the route this terminal was pointed at. Without this the settings
// are written on Save and then silently ignored on the next launch, which
// reads exactly like they never saved.
final store = ref.read(localStoreProvider);
if (store.isReady) {
ref.read(syncConfigProvider.notifier).state =
await store.syncConfig.load(ref.read(syncConfigProvider));
}
await ref.read(connectivityServiceProvider).start();
final engine = ref.read(syncEngineProvider);
// A background drain moves bills out of the pending set, so the tallies and
// shift totals on screen are stale the moment one finishes.
var wasSyncing = false;
final subscription = engine.states.listen((state) {
if (wasSyncing && !state.isSyncing) {
ref.read(orderVersionProvider.notifier).state++;
}
wasSyncing = state.isSyncing;
});
ref.onDispose(subscription.cancel);
await engine.start();
// The retained presence record needs a Last Will to pair with, so it genuinely
// only exists on the broker. The heartbeat below does not, and is started for
// every route.
final transport = ref.read(orderTransportProvider);
if (transport is MqttOrderTransport) {
final reporter = PresenceReporter(
transport: transport,
terminal: ref.read(terminalIdentityProvider),
config: ref.read(syncConfigProvider),
engine: engine,
appVersion: AppConstants.appVersion,
catalogueRevision: () async =>
ref.read(syncRepositoryProvider).catalogueRevision,
);
ref.onDispose(reporter.dispose);
await reporter.start();
}
// The 30-second heartbeat the head-office board reads. Separate from the
// retained presence record above: that one is paired with the Last Will and
// answers "is this till alive", while this carries queue depth, today's
// trading and hardware state — what tells a till that is merely quiet from
// one that has stopped uploading.
//
// Outside the MQTT check on purpose. It used to be inside, which meant a shop
// on the HTTP route uploaded every bill correctly and never appeared on the
// board at all — with nothing logged, because nothing had failed. Both routes
// can carry a heartbeat now, and the transport decides how.
{
final health = HealthReporter(
transport: transport,
terminal: ref.read(terminalIdentityProvider),
config: ref.read(syncConfigProvider),
engine: engine,
repository: ref.read(syncRepositoryProvider),
appVersion: AppConstants.appVersion,
// Read on each beat, so re-pointing the printer in Settings takes effect
// without a restart.
printerEndpoint: () {
final printer = ref.read(printerSettingsProvider);
final host = printer.drawerHost;
if (host == null || host.isEmpty) return null;
return (host: host, port: printer.drawerPort);
},
);
ref.onDispose(health.dispose);
await health.start();
}
});