Files
GameTime/server/lib/infrastructure/postgres/share_repository.dart
Blomios 917777e18b chore(wip): consolidation intermédiaire multi-tickets (sprints Statistiques, UI, Bug resolution, Serveur-client)
Regroupe l'état de travail en cours réalisé dans un même worktree sur
plusieurs tickets/sprints (#85, #136, #145, #155-160, #162-164),
mélangeant des tickets QA et inProgress. Ne constitue pas une feature
terminée : commit de sauvegarde avant triage/split par ticket en
branches feature/* dédiées. Exclut les dossiers d'environnement de
build locaux et le heap dump parasite (.gitignore mis à jour).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-28 16:48:54 +02:00

355 lines
11 KiB
Dart

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<void> insertShare({
required Share share,
required List<ShareRecipient> 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<List<ShareInboxItem>> 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<String, Object?>),
];
}
@override
Future<Share?> 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<String, Object?>);
}
@override
Future<ShareRecipient?> 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<String, Object?>,
);
}
@override
Future<void> 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<List<SyncedResource>> acceptShare({
required String recipientId,
required DateTime respondedAt,
required List<SyncedResource> resources,
}) {
return connection.runTx((session) async {
final written = <SyncedResource>[];
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<void> 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<String, Object?> values) {
return ShareInboxItem(
share: _shareFromValues(values),
recipient: _recipientFromValues(values),
);
}
Share _shareFromValues(Map<String, Object?> 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<String, Object?> 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<SyncedResource> _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<String, Object?>,
);
}
SyncedResource _resourceFromValues(Map<String, Object?> 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<String, Object?> _payloadValue(Object? value) {
if (value is Map<String, Object?>) {
return value;
}
if (value is Map) {
return Map<String, Object?>.from(value);
}
if (value is String) {
final decoded = jsonDecode(value);
if (decoded is Map<String, Object?>) {
return decoded;
}
if (decoded is Map) {
return Map<String, Object?>.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();
}