import 'dart:convert'; import 'package:postgres/postgres.dart'; import '../../application/application.dart'; import '../../domain/domain.dart'; final class PostgresShareRepository implements ShareRepository { const PostgresShareRepository(this.connection); final Connection connection; @override Future insertShare({ required Share share, required List recipients, }) async { await connection.execute( Sql.named(''' INSERT INTO shares ( id, sender_user_id, share_kind, pack_name, resource_type, payload_json, created_at, revoked_at ) VALUES ( @id::uuid, @sender_user_id::uuid, @share_kind, @pack_name, @resource_type, @payload_json::jsonb, @created_at, @revoked_at ) '''), parameters: { 'id': share.id, 'sender_user_id': share.senderUserId, 'share_kind': share.kind.wireName, 'pack_name': share.packName, 'resource_type': share.resourceType?.wireName, 'payload_json': jsonEncode(share.payloadJson), 'created_at': share.createdAt, 'revoked_at': share.revokedAt, }, ); for (final recipient in recipients) { await connection.execute( Sql.named(''' INSERT INTO share_recipients ( id, share_id, recipient_user_id, status, responded_at ) VALUES ( @id::uuid, @share_id::uuid, @recipient_user_id::uuid, @status, @responded_at ) ON CONFLICT (share_id, recipient_user_id) DO NOTHING '''), parameters: { 'id': recipient.id, 'share_id': recipient.shareId, 'recipient_user_id': recipient.recipientUserId, 'status': recipient.status.wireName, 'responded_at': recipient.respondedAt, }, ); } } @override Future> listInbox(String recipientUserId) async { final result = await connection.execute( Sql.named(''' SELECT s.id AS share_id, s.sender_user_id, s.share_kind, s.pack_name, s.resource_type, s.payload_json, s.created_at, s.revoked_at, sr.id AS recipient_id, sr.recipient_user_id, sr.status, sr.responded_at FROM share_recipients sr INNER JOIN shares s ON s.id = sr.share_id WHERE sr.recipient_user_id = @recipient_user_id::uuid ORDER BY s.created_at DESC, s.id DESC '''), parameters: {'recipient_user_id': recipientUserId}, ); return [ for (final row in result) _inboxItemFromValues(row.toColumnMap() as Map), ]; } @override Future findShareById(String shareId) async { final result = await connection.execute( Sql.named(''' SELECT id, sender_user_id, share_kind, pack_name, resource_type, payload_json, created_at, revoked_at FROM shares WHERE id = @id::uuid LIMIT 1 '''), parameters: {'id': shareId}, ); return result.isEmpty ? null : _shareFromValues(result.single.toColumnMap() as Map); } @override Future findRecipient({ required String shareId, required String recipientUserId, }) async { final result = await connection.execute( Sql.named(''' SELECT id, share_id, recipient_user_id, status, responded_at FROM share_recipients WHERE share_id = @share_id::uuid AND recipient_user_id = @recipient_user_id::uuid LIMIT 1 '''), parameters: {'share_id': shareId, 'recipient_user_id': recipientUserId}, ); return result.isEmpty ? null : _recipientFromValues( result.single.toColumnMap() as Map, ); } @override Future updateRecipientStatus({ required String recipientId, required ShareRecipientStatus status, required DateTime respondedAt, }) async { await connection.execute( Sql.named(''' UPDATE share_recipients SET status = @status, responded_at = @responded_at WHERE id = @id::uuid '''), parameters: { 'id': recipientId, 'status': status.wireName, 'responded_at': respondedAt.toUtc(), }, ); } @override Future> acceptShare({ required String recipientId, required DateTime respondedAt, required List resources, }) { return connection.runTx((session) async { final written = []; for (final resource in resources) { written.add(await _upsertResource(session, resource)); } await session.execute( Sql.named(''' UPDATE share_recipients SET status = @status, responded_at = @responded_at WHERE id = @id::uuid '''), parameters: { 'id': recipientId, 'status': ShareRecipientStatus.accepted.wireName, 'responded_at': respondedAt.toUtc(), }, ); return written; }); } @override Future revokeShare({ required String shareId, required DateTime revokedAt, }) async { await connection.execute( Sql.named(''' UPDATE shares SET revoked_at = @revoked_at WHERE id = @id::uuid AND revoked_at IS NULL '''), parameters: {'id': shareId, 'revoked_at': revokedAt.toUtc()}, ); } } ShareInboxItem _inboxItemFromValues(Map values) { return ShareInboxItem( share: _shareFromValues(values), recipient: _recipientFromValues(values), ); } Share _shareFromValues(Map values) { return Share( id: _stringValue(values['share_id'] ?? values['id']), senderUserId: _stringValue(values['sender_user_id']), kind: ShareKind.parse( _stringValue(values['share_kind'] ?? ShareKind.single.wireName), ), packName: values['pack_name'] as String?, resourceType: _nullableResourceType(values['resource_type']), payloadJson: _payloadValue(values['payload_json']), createdAt: _dateTimeValue(values['created_at']), revokedAt: _nullableDateTimeValue(values['revoked_at']), ); } ShareRecipient _recipientFromValues(Map values) { return ShareRecipient( id: _stringValue(values['recipient_id'] ?? values['id']), shareId: _stringValue(values['share_id']), recipientUserId: _stringValue(values['recipient_user_id']), status: ShareRecipientStatus.parse(_stringValue(values['status'])), respondedAt: _nullableDateTimeValue(values['responded_at']), ); } Future _upsertResource( Session session, SyncedResource resource, ) async { final result = await session.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 ) 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 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, }, ); return _resourceFromValues( result.single.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?, ); } SyncedResourceType? _nullableResourceType(Object? value) { if (value == null) { return null; } return SyncedResourceType.parse(_stringValue(value)); } 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(); }