import 'dart:convert'; import 'package:postgres/postgres.dart'; import '../../application/application.dart'; import '../../domain/domain.dart'; final class PostgresSyncedResourceRepository implements SyncedResourceRepository { const PostgresSyncedResourceRepository(this.connection); final Connection connection; @override Future upsertWithLww(SyncedResource resource) async { final result = await connection.execute( Sql.named(''' WITH upserted AS ( INSERT INTO synced_resources ( server_id, owner_user_id, resource_type, client_id, payload_json, schema_version, client_updated_at, deleted_at, origin_device_id ) VALUES ( @server_id::uuid, @owner_user_id::uuid, @resource_type, @client_id, @payload_json::jsonb, @schema_version, @client_updated_at, @deleted_at, @origin_device_id ) ON CONFLICT (owner_user_id, resource_type, client_id) DO UPDATE SET payload_json = EXCLUDED.payload_json, schema_version = EXCLUDED.schema_version, client_updated_at = EXCLUDED.client_updated_at, deleted_at = EXCLUDED.deleted_at, origin_device_id = EXCLUDED.origin_device_id WHERE synced_resources.client_updated_at < EXCLUDED.client_updated_at RETURNING server_id, owner_user_id, resource_type, client_id, payload_json, schema_version, client_updated_at, server_updated_at, deleted_at, origin_device_id, 'accepted'::text AS write_status ) SELECT * FROM upserted UNION ALL SELECT existing.server_id, existing.owner_user_id, existing.resource_type, existing.client_id, existing.payload_json, existing.schema_version, existing.client_updated_at, existing.server_updated_at, existing.deleted_at, existing.origin_device_id, 'ignoredOlder'::text AS write_status FROM synced_resources existing WHERE existing.owner_user_id = @owner_user_id::uuid AND existing.resource_type = @resource_type AND existing.client_id = @client_id AND NOT EXISTS (SELECT 1 FROM upserted) LIMIT 1 '''), parameters: { 'server_id': resource.serverId, 'owner_user_id': resource.ownerUserId, 'resource_type': resource.resourceType.wireName, 'client_id': resource.clientId, 'payload_json': jsonEncode(resource.payloadJson), 'schema_version': resource.schemaVersion, 'client_updated_at': resource.clientUpdatedAt, 'deleted_at': resource.deletedAt, 'origin_device_id': resource.originDeviceId, }, ); final row = result.single; final values = row.toColumnMap() as Map; return SyncWriteResult( status: values['write_status'] == 'accepted' ? SyncWriteStatus.accepted : SyncWriteStatus.ignoredOlder, resource: _resourceFromValues(values), ); } @override Future> findAllForUserSince({ required String ownerUserId, DateTime? since, }) async { final result = since == null ? await connection.execute( Sql.named(''' SELECT server_id, owner_user_id, resource_type, client_id, payload_json, schema_version, client_updated_at, server_updated_at, deleted_at, origin_device_id FROM synced_resources WHERE owner_user_id = @owner_user_id::uuid ORDER BY server_updated_at ASC, server_id ASC '''), parameters: {'owner_user_id': ownerUserId}, ) : await connection.execute( Sql.named(''' SELECT server_id, owner_user_id, resource_type, client_id, payload_json, schema_version, client_updated_at, server_updated_at, deleted_at, origin_device_id FROM synced_resources WHERE owner_user_id = @owner_user_id::uuid AND server_updated_at > @since ORDER BY server_updated_at ASC, server_id ASC '''), parameters: {'owner_user_id': ownerUserId, 'since': since.toUtc()}, ); return [ for (final row in result) _resourceFromValues(row.toColumnMap() as Map), ]; } } SyncedResource _resourceFromValues(Map values) { return SyncedResource( serverId: _stringValue(values['server_id']), ownerUserId: _stringValue(values['owner_user_id']), resourceType: SyncedResourceType.parse( _stringValue(values['resource_type']), ), clientId: _stringValue(values['client_id']), payloadJson: _payloadValue(values['payload_json']), schemaVersion: values['schema_version'] as int, clientUpdatedAt: _dateTimeValue(values['client_updated_at']), serverUpdatedAt: _dateTimeValue(values['server_updated_at']), deletedAt: _nullableDateTimeValue(values['deleted_at']), originDeviceId: values['origin_device_id'] as String?, ); } Map _payloadValue(Object? value) { if (value is Map) { return value; } if (value is Map) { return Map.from(value); } if (value is String) { final decoded = jsonDecode(value); if (decoded is Map) { return decoded; } if (decoded is Map) { return Map.from(decoded); } } throw const FormatException('Expected JSON object payload.'); } String _stringValue(Object? value) { if (value == null) { throw const FormatException('Expected non-null string value.'); } return value.toString(); } DateTime _dateTimeValue(Object? value) { final result = _nullableDateTimeValue(value); if (result == null) { throw const FormatException('Expected non-null DateTime value.'); } return result; } DateTime? _nullableDateTimeValue(Object? value) { if (value == null) { return null; } if (value is DateTime) { return value.toUtc(); } return DateTime.parse(value.toString()).toUtc(); }