import '../domain/domain.dart'; import 'ports.dart'; final class PushSyncItemInput { const PushSyncItemInput({ required this.resourceType, required this.clientId, required this.schemaVersion, required this.clientUpdatedAt, required this.deletedAt, required this.payloadJson, this.validationError, }); const PushSyncItemInput.invalid({ required String message, this.resourceType, this.clientId, }) : schemaVersion = null, clientUpdatedAt = null, deletedAt = null, payloadJson = null, validationError = message; final String? resourceType; final String? clientId; final int? schemaVersion; final DateTime? clientUpdatedAt; final DateTime? deletedAt; final Map? payloadJson; final String? validationError; } final class PushSyncResultItem { const PushSyncResultItem({ required this.resourceType, required this.clientId, this.serverId, required this.status, this.serverUpdatedAt, this.message, }); final String? resourceType; final String? clientId; final String? serverId; final String status; final DateTime? serverUpdatedAt; final String? message; } final class PushSyncResult { const PushSyncResult({required this.serverCursor, required this.results}); final DateTime serverCursor; final List results; } final class PullSyncResult { const PullSyncResult({required this.serverCursor, required this.items}); final DateTime serverCursor; final List items; } final class ExchangeSyncResult { const ExchangeSyncResult({required this.push, required this.pull}); final PushSyncResult push; final PullSyncResult pull; } final class PushSyncUseCase { const PushSyncUseCase({ required this.resources, required this.clock, required this.ids, }); final SyncedResourceRepository resources; final Clock clock; final IdGenerator ids; Future execute({ required String ownerUserId, required String? deviceId, required List items, }) async { final results = []; final normalizedDeviceId = _blankToNull(deviceId); for (final item in items) { try { final validationError = item.validationError; if (validationError != null) { throw ValidationException(validationError); } if (normalizedDeviceId == null) { throw const ValidationException('deviceId must not be blank.'); } final resource = SyncedResource( serverId: ids.newId(), ownerUserId: ownerUserId, resourceType: SyncedResourceType.parse( _required(item.resourceType, 'resourceType'), ), clientId: _required(item.clientId, 'clientId'), payloadJson: item.payloadJson ?? _missing('payload'), schemaVersion: item.schemaVersion ?? _missing('schemaVersion'), clientUpdatedAt: item.clientUpdatedAt ?? _missing('clientUpdatedAt'), serverUpdatedAt: clock.now(), deletedAt: item.deletedAt, originDeviceId: normalizedDeviceId, ); final written = await resources.upsertWithLww(resource); results.add( PushSyncResultItem( resourceType: written.resource.resourceType.wireName, clientId: written.resource.clientId, serverId: written.resource.serverId, status: written.status == SyncWriteStatus.accepted ? 'accepted' : 'ignoredOlder', serverUpdatedAt: written.resource.serverUpdatedAt, ), ); } on ValidationException catch (error) { results.add( PushSyncResultItem( resourceType: item.resourceType, clientId: item.clientId, status: 'error', message: error.message, ), ); } on FormatException catch (error) { results.add( PushSyncResultItem( resourceType: item.resourceType, clientId: item.clientId, status: 'error', message: error.message, ), ); } } return PushSyncResult( serverCursor: await _currentCursor(resources, ownerUserId, clock), results: results, ); } } final class PullSyncUseCase { const PullSyncUseCase({required this.resources, required this.clock}); final SyncedResourceRepository resources; final Clock clock; Future execute({ required String ownerUserId, DateTime? since, }) async { final items = await resources.findAllForUserSince( ownerUserId: ownerUserId, since: since?.toUtc(), ); return PullSyncResult( serverCursor: _cursorFromItems(items, clock.now()), items: items, ); } } final class ExchangeSyncUseCase { const ExchangeSyncUseCase({required this.push, required this.pull}); final PushSyncUseCase push; final PullSyncUseCase pull; Future execute({ required String ownerUserId, required String? deviceId, required List items, DateTime? since, }) async { final pushResult = await push.execute( ownerUserId: ownerUserId, deviceId: deviceId, items: items, ); final pullResult = await pull.execute( ownerUserId: ownerUserId, since: since, ); return ExchangeSyncResult(push: pushResult, pull: pullResult); } } Future _currentCursor( SyncedResourceRepository resources, String ownerUserId, Clock clock, ) async { final allItems = await resources.findAllForUserSince( ownerUserId: ownerUserId, ); return _cursorFromItems(allItems, clock.now()); } DateTime _cursorFromItems(List items, DateTime fallback) { if (items.isEmpty) { return fallback.toUtc(); } return items .map((item) => item.serverUpdatedAt) .reduce((left, right) => left.isAfter(right) ? left : right) .toUtc(); } String _required(String? value, String label) { final normalized = _blankToNull(value); if (normalized == null) { throw ValidationException('$label must not be blank.'); } return normalized; } Never _missing(String label) { throw ValidationException('$label is required.'); } String? _blankToNull(String? value) { final trimmed = value?.trim(); return trimmed == null || trimmed.isEmpty ? null : trimmed; }