diff --git a/ios/Runner.xcodeproj/project.pbxproj b/ios/Runner.xcodeproj/project.pbxproj index 6ee65a5f..404012c8 100644 --- a/ios/Runner.xcodeproj/project.pbxproj +++ b/ios/Runner.xcodeproj/project.pbxproj @@ -9,7 +9,6 @@ /* Begin PBXBuildFile section */ 1498D2341E8E89220040F4C2 /* GeneratedPluginRegistrant.m in Sources */ = {isa = PBXBuildFile; fileRef = 1498D2331E8E89220040F4C2 /* GeneratedPluginRegistrant.m */; }; 22D5B0E9A1860BEEC6084203 /* Pods_Runner.framework in Frameworks */ = {isa = PBXBuildFile; fileRef = 44DF6CB9CFAE3FB41F9BF8BF /* Pods_Runner.framework */; }; - 331C808B294A63AB00263BE5 /* RunnerTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 331C807B294A618700263BE5 /* RunnerTests.swift */; }; 3B3967161E833CAA004F5970 /* AppFrameworkInfo.plist in Resources */ = {isa = PBXBuildFile; fileRef = 3B3967151E833CAA004F5970 /* AppFrameworkInfo.plist */; }; 59C723F32EA9C6BF002F18BF /* AppIcon.icon in Resources */ = {isa = PBXBuildFile; fileRef = 59C723F22EA9C6BF002F18BF /* AppIcon.icon */; }; 74858FAF1ED2DC5600515810 /* AppDelegate.swift in Sources */ = {isa = PBXBuildFile; fileRef = 74858FAE1ED2DC5600515810 /* AppDelegate.swift */; }; @@ -51,7 +50,6 @@ 1498D2321E8E86230040F4C2 /* GeneratedPluginRegistrant.h */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.c.h; path = GeneratedPluginRegistrant.h; sourceTree = ""; }; 1498D2331E8E89220040F4C2 /* GeneratedPluginRegistrant.m */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = sourcecode.c.objc; path = GeneratedPluginRegistrant.m; sourceTree = ""; }; 18D5FD77F116FC12573839C2 /* Pods-Runner.release.xcconfig */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = text.xcconfig; name = "Pods-Runner.release.xcconfig"; path = "Target Support Files/Pods-Runner/Pods-Runner.release.xcconfig"; sourceTree = ""; }; - 331C807B294A618700263BE5 /* RunnerTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = RunnerTests.swift; sourceTree = ""; }; 331C8081294A63A400263BE5 /* RunnerTests.xctest */ = {isa = PBXFileReference; explicitFileType = wrapper.cfbundle; includeInIndex = 0; path = RunnerTests.xctest; sourceTree = BUILT_PRODUCTS_DIR; }; 3B3967151E833CAA004F5970 /* AppFrameworkInfo.plist */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = text.plist.xml; name = AppFrameworkInfo.plist; path = Flutter/AppFrameworkInfo.plist; sourceTree = ""; }; 44DF6CB9CFAE3FB41F9BF8BF /* Pods_Runner.framework */ = {isa = PBXFileReference; explicitFileType = wrapper.framework; includeInIndex = 0; path = Pods_Runner.framework; sourceTree = BUILT_PRODUCTS_DIR; }; @@ -98,7 +96,6 @@ 331C8082294A63A400263BE5 /* RunnerTests */ = { isa = PBXGroup; children = ( - 331C807B294A618700263BE5 /* RunnerTests.swift */, ); path = RunnerTests; sourceTree = ""; @@ -370,7 +367,6 @@ isa = PBXSourcesBuildPhase; buildActionMask = 2147483647; files = ( - 331C808B294A63AB00263BE5 /* RunnerTests.swift in Sources */, ); runOnlyForDeploymentPostprocessing = 0; }; diff --git a/ios/RunnerTests/RunnerTests.swift b/ios/RunnerTests/RunnerTests.swift deleted file mode 100644 index 86a7c3b1..00000000 --- a/ios/RunnerTests/RunnerTests.swift +++ /dev/null @@ -1,12 +0,0 @@ -import Flutter -import UIKit -import XCTest - -class RunnerTests: XCTestCase { - - func testExample() { - // If you add code to the Runner application, consider adding tests here. - // See https://developer.apple.com/documentation/xctest for more information about using XCTest. - } - -} diff --git a/lib/src/core/network/atproto/data/repositories/feed_repository_impl.dart b/lib/src/core/network/atproto/data/repositories/feed_repository_impl.dart index f7bb02d4..393705bc 100644 --- a/lib/src/core/network/atproto/data/repositories/feed_repository_impl.dart +++ b/lib/src/core/network/atproto/data/repositories/feed_repository_impl.dart @@ -1,6 +1,3 @@ -import 'dart:async'; -import 'dart:convert'; -import 'dart:io'; import 'dart:typed_data'; import 'dart:ui' as ui; @@ -26,17 +23,12 @@ import 'package:bluesky_poptart/app/bsky/feed/repost.dart' as bsky_repost; import 'package:poptart_lex/com/atproto/repo/strong_ref.dart'; import 'package:poptart_lex/com/atproto/repo/upload_blob.dart' as repo_upload_blob; -import 'package:poptart_lex/com/atproto/server/get_service_auth.dart' - as server_get_service_auth; import 'package:poptart/poptart.dart'; import 'package:bluesky_poptart/app/bsky/richtext/facet.dart'; import 'package:get_it/get_it.dart'; -import 'package:http/http.dart' as http; import 'package:image/image.dart' as img; import 'package:image_picker/image_picker.dart'; -import 'package:path/path.dart' as path; -import 'package:spark/src/core/config/app_config.dart'; import 'package:spark/src/core/network/atproto/data/adapters/bsky/feed_adapter.dart'; import 'package:spark/src/core/network/atproto/data/models/feed_models.dart'; import 'package:spark/src/core/network/atproto/data/models/models.dart'; @@ -44,11 +36,11 @@ import 'package:spark/src/core/network/atproto/data/models/pref_models.dart'; import 'package:spark/src/core/network/atproto/data/models/record_write_adapters.dart'; import 'package:spark/src/core/network/atproto/data/repositories/feed_repository.dart'; import 'package:spark/src/core/network/atproto/data/repositories/sprk_repository.dart'; +import 'package:spark/src/core/network/atproto/data/services/video_upload_service.dart'; import 'package:spark/src/core/utils/bluesky_crosspost_text.dart'; import 'package:spark/src/core/utils/logging/log_service.dart'; import 'package:spark/src/core/utils/logging/logger.dart'; import 'package:spark/src/core/utils/share_urls.dart'; -import 'package:spark/src/core/utils/video_upload_exception.dart'; import 'package:sprk_poptart/so/sprk/feed/like.dart' as sprk_like; import 'package:sprk_poptart/so/sprk/feed/get_author_feed.dart' as sprk_get_author_feed; @@ -72,7 +64,6 @@ import 'package:sprk_poptart/so/sprk/feed/get_suggested_feeds.dart' import 'package:sprk_poptart/so/sprk/feed/get_timeline.dart' as sprk_get_timeline; import 'package:sprk_poptart/so/sprk/feed/repost.dart' as sprk_repost; -import 'package:sprk_poptart/so/sprk/sound/defs/audio_details.dart'; import 'package:sprk_poptart/so/sprk/feed/search_posts.dart' as sprk_search_posts; @@ -82,27 +73,19 @@ class FeedRepositoryImpl implements FeedRepository { this._client, { SparkLogger? logger, DateTime Function()? now, - http.Client? videoHttpClient, - Future Function(Duration)? videoProcessingDelay, - File Function(String)? videoFile, - Future Function(PoptartClient)? videoServiceAuthTokenRequest, + VideoUploadService? videoUploadService, }) : _logger = logger ?? GetIt.instance().getLogger('FeedRepository'), - _now = now ?? DateTime.now, - _videoHttpClient = videoHttpClient ?? http.Client(), - _videoProcessingDelay = videoProcessingDelay ?? Future.delayed, - _videoFile = videoFile ?? File.new, - _requestVideoServiceAuthTokenOverride = videoServiceAuthTokenRequest { + _now = now ?? DateTime.now { + _videoUploadService = + videoUploadService ?? + VideoUploadClient(_client.authRepository, logger: _logger, now: _now); _logger.v('FeedRepository initialized'); } final SprkRepository _client; final SparkLogger _logger; final DateTime Function() _now; - final http.Client _videoHttpClient; - final Future Function(Duration) _videoProcessingDelay; - final File Function(String) _videoFile; - final Future Function(PoptartClient)? - _requestVideoServiceAuthTokenOverride; + late final VideoUploadService _videoUploadService; /// Formats labeler DIDs into the atproto-accept-labelers header format /// Format: "did1,did2,did3" (comma-separated list) @@ -1252,353 +1235,13 @@ class FeedRepositoryImpl implements FeedRepository { Future uploadVideo( String videoPath, { void Function(double progress)? onUploadProgress, - }) async { - _logger.d('Uploading video from path: $videoPath'); - - return _client.executeWithRetry(() async { - if (!_client.authRepository.isAuthenticated) { - _logger.w('Not authenticated'); - throw Exception('Not authenticated'); - } - final authAtProto = _client.authRepository.atproto; - if (authAtProto == null || authAtProto.oAuthSession == null) { - throw Exception('AtProto not initialized'); - } - - // Handle file:// URL scheme - var cleanVideoPath = videoPath; - if (videoPath.startsWith('file://')) { - cleanVideoPath = videoPath.replaceFirst('file://', ''); - } - - // Validate the video file - final file = _videoFile(cleanVideoPath); - if (!file.existsSync()) { - throw Exception('Video file not found: $cleanVideoPath'); - } - - // Check if the video is in a compatible format - // Use BigInt to avoid overflow for very large files (>2GB) - final videoSizeBigInt = BigInt.from(await file.length()); - if (videoSizeBigInt == BigInt.zero) { - throw Exception('Video file is empty'); - } - - // Check for integer overflow (files > 2GB could overflow int on 32-bit) - if (videoSizeBigInt > BigInt.from(2 * 1024 * 1024 * 1024)) { - _logger.w( - 'Video file exceeds 2GB, may cause issues: $videoSizeBigInt bytes', - ); - throw VideoUploadException( - 'Video is too large. Maximum supported size is 2GB.', - statusCode: 413, - uploadSizeBytes: videoSizeBigInt.toInt(), - limitBytes: 2 * 1024 * 1024 * 1024, - ); - } - - final videoSizeBytes = videoSizeBigInt.toInt(); - _logger.i('Video file size: $videoSizeBytes bytes'); - final maxUploadSizeBytes = (AppConfig.maxUploadSizeMB * 1024 * 1024) - .round(); - if (maxUploadSizeBytes > 0 && videoSizeBytes > maxUploadSizeBytes) { - _logger.w( - 'Video file exceeds upload limit: $videoSizeBytes bytes ' - '(limit: $maxUploadSizeBytes bytes)', - ); - throw VideoUploadException( - 'Video is too large to upload.', - statusCode: 413, - uploadSizeBytes: videoSizeBytes, - limitBytes: maxUploadSizeBytes, - ); - } - - // Validate content length will fit in HTTP header (max ~2GB for int32) - if (videoSizeBytes > 2147483647) { - _logger.e('Video file too large for HTTP content-length header'); - throw VideoUploadException( - 'Video is too large to upload.', - statusCode: 413, - uploadSizeBytes: videoSizeBytes, - limitBytes: 2147483647, - ); - } - - var serviceToken = await _createVideoServiceAuthToken(); - final uploadRequest = - http.StreamedRequest( - 'POST', - Uri.parse( - '${AppConfig.videoServiceUrl}/xrpc/so.sprk.video.uploadVideo', - ), - ) - ..contentLength = videoSizeBytes - ..headers.addAll({ - 'Authorization': 'Bearer $serviceToken', - 'Content-Type': _getContentType(cleanVideoPath), - }); - - onUploadProgress?.call(0); - final uploadResponseFuture = _videoHttpClient.send(uploadRequest); - try { - await uploadRequest.sink.addStream( - _trackUploadProgress( - file.openRead(), - totalBytes: videoSizeBytes, - onUploadProgress: onUploadProgress, - ), - ); - } finally { - unawaited(uploadRequest.sink.close()); - } - var response = await http.Response.fromStream(await uploadResponseFuture); - - if (response.statusCode != 200) { - _logger.e( - 'Video upload failed: ${response.statusCode} ${response.body}', - ); - throw VideoUploadException( - _buildVideoUploadFailureMessage( - fallback: response.statusCode == 413 - ? 'Video is too large to upload.' - : 'Failed to upload video.', - detail: response.body, - ), - statusCode: response.statusCode, - uploadSizeBytes: videoSizeBytes, - limitBytes: maxUploadSizeBytes > 0 ? maxUploadSizeBytes : null, - responseBody: response.body, - ); - } - - // Parse the response - dynamic responseData = jsonDecode(response.body); - _logger.d('Video upload response: $responseData'); - - // Poll job status until it finishes (handles both QUEUED and PROCESSING) - var jobState = responseData['jobStatus']?['state'] as String?; - var attempts = 0; - var consecutivePollErrors = 0; - const maxAttempts = 120; // ~4 minutes at 2s interval - const maxConsecutivePollErrors = 3; // Allow 3 consecutive polling errors - while (jobState == 'JOB_STATE_QUEUED' || - jobState == 'JOB_STATE_PROCESSING') { - _logger.d('Video upload in progress, status: $jobState'); - // Small backoff to avoid hammering the service - await _videoProcessingDelay(const Duration(seconds: 2)); - attempts++; - if (attempts > maxAttempts) { - throw const VideoUploadException( - 'Video processing timed out. Please try again.', - ); - } - - try { - response = await _videoHttpClient.get( - Uri.parse( - '${AppConfig.videoServiceUrl}/xrpc/so.sprk.video.getJobStatus', - ).replace( - queryParameters: {'jobId': responseData['jobStatus']?['jobId']}, - ), - headers: { - 'Authorization': 'Bearer $serviceToken', - 'Content-Type': _getContentType(cleanVideoPath), - }, - ); - if (_isExpiredVideoServiceTokenResponse(response)) { - _logger.i( - 'Video service token expired while polling; minting a new token', - ); - serviceToken = await _createVideoServiceAuthToken( - refreshPdsSessionOnFailure: true, - ); - response = await _videoHttpClient.get( - Uri.parse( - '${AppConfig.videoServiceUrl}/xrpc/so.sprk.video.getJobStatus', - ).replace( - queryParameters: {'jobId': responseData['jobStatus']?['jobId']}, - ), - headers: { - 'Authorization': 'Bearer $serviceToken', - 'Content-Type': _getContentType(cleanVideoPath), - }, - ); - } - if (response.statusCode != 200) { - throw Exception( - 'Failed to check video upload status: ${response.statusCode} ' - '${response.body}', - ); - } - responseData = jsonDecode(response.body); - _logger.d('Video upload status response: $responseData'); - jobState = responseData['jobStatus']?['state'] as String?; - // Reset consecutive errors on success - consecutivePollErrors = 0; - } catch (e) { - // Network or parsing error during polling - log and retry - consecutivePollErrors++; - _logger.w( - 'Error polling video upload status on attempt ' - '$attempts/$maxAttempts ' - '(consecutive errors: $consecutivePollErrors/$maxConsecutivePollErrors): ' - '$e', - ); - - // Only fail if we've had too many consecutive errors - if (consecutivePollErrors >= maxConsecutivePollErrors) { - _logger.e( - 'Too many consecutive polling errors, giving up: $e', - error: e, - ); - throw VideoUploadException( - _buildVideoUploadFailureMessage( - fallback: 'Failed to check video processing status.', - detail: e.toString(), - ), - responseBody: e.toString(), - ); - } - - // Continue polling on transient errors - _logger.d('Retrying poll after error...'); - jobState = - 'JOB_STATE_PROCESSING'; // Assume still processing and retry - } - } - - if (responseData['jobStatus']?['state'] == 'JOB_STATE_FAILED') { - final failureMessage = _buildVideoUploadFailureMessage( - fallback: 'Video processing failed.', - detail: responseData['jobStatus'] ?? responseData, - ); - _logger.e( - 'Video processing job failed: $failureMessage', - error: responseData, - ); - throw VideoUploadException( - failureMessage, - responseBody: jsonEncode(responseData), - ); - } - - // Parse video blob - Map videoBlobData; - if (responseData case {'jobStatus': {'blob': final blobData}}) { - videoBlobData = blobData as Map; - } else if (responseData case {'blobRef': final blobRef}) { - videoBlobData = blobRef as Map; - } else { - throw Exception('Unexpected response format: $responseData'); - } - final videoBlob = Blob.fromJson(videoBlobData); - - // Parse audio blob if present - Blob? audioBlob; - AudioDetails? audioDetails; - if (responseData case {'jobStatus': {'audio': final audioData}}) { - final audio = audioData as Map; - if (audio['blob'] != null) { - audioBlob = Blob.fromJson(audio['blob'] as Map); - _logger.d('Extracted audio blob: ${audioBlob.size} bytes'); - } - if (audio['details'] != null) { - audioDetails = AudioDetails.fromJson( - audio['details'] as Map, - ); - } - } - - return VideoUploadResult( - videoBlob: videoBlob, - audioBlob: audioBlob, - audioDetails: audioDetails, - ); - }); - } - - Stream> _trackUploadProgress( - Stream> chunks, { - required int totalBytes, - void Function(double progress)? onUploadProgress, - }) async* { - var uploadedBytes = 0; - - await for (final chunk in chunks) { - uploadedBytes += chunk.length; - if (totalBytes > 0) { - onUploadProgress?.call( - (uploadedBytes / totalBytes).clamp(0, 1).toDouble(), - ); - } - yield chunk; - } - - onUploadProgress?.call(1); - } - - Future _createVideoServiceAuthToken({ - bool refreshPdsSessionOnFailure = false, - }) async { - final atproto = _client.authRepository.atproto; - if (atproto == null) { - throw Exception('AtProto not initialized'); - } - - try { - return await _requestVideoServiceAuthToken(atproto); - } catch (e) { - if (!refreshPdsSessionOnFailure) { - rethrow; - } - - _logger.i( - 'Refreshing PDS session before minting video service token', - error: e, - ); - final refreshed = await _client.authRepository.refreshToken(); - if (!refreshed) { - throw Exception('Session expired. Please log in again.'); - } - - final refreshedAtproto = _client.authRepository.atproto; - if (refreshedAtproto == null) { - throw Exception('AtProto not initialized after refresh'); - } - return _requestVideoServiceAuthToken(refreshedAtproto); - } - } - - Future _requestVideoServiceAuthToken(PoptartClient atproto) async { - final override = _requestVideoServiceAuthTokenOverride; - if (override != null) { - return override(atproto); - } - final serviceTokenRes = await atproto.call( - server_get_service_auth.comAtprotoServerGetServiceAuth, - parameters: server_get_service_auth.ServerGetServiceAuthInput( - aud: 'did:web:${atproto.service}', - lxm: 'com.atproto.repo.uploadBlob', - exp: - _now() - .toUtc() - .add(const Duration(minutes: 5)) - .millisecondsSinceEpoch ~/ - 1000, + }) { + return _client.executeWithRetry( + () => _videoUploadService.uploadVideo( + videoPath, + onUploadProgress: onUploadProgress, ), ); - - return serviceTokenRes.data.token; - } - - bool _isExpiredVideoServiceTokenResponse(http.Response response) { - if (response.statusCode != 401) { - return false; - } - - final body = response.body.toLowerCase(); - return body.contains('jwt has expired') || body.contains('invalidtoken'); } /// Crosspost images to Bluesky using adapter to handle Bluesky-specific model @@ -2123,127 +1766,4 @@ class FeedRepositoryImpl implements FeedRepository { ); return (posts: [], cursor: null); } - - /// Helper method to determine content type based on file extension - String _getContentType(String videoPath) { - final extension = path.extension(videoPath).toLowerCase(); - - switch (extension) { - case '.mp4': - return 'video/mp4'; - case '.mov': - return 'video/quicktime'; - case '.avi': - return 'video/x-msvideo'; - case '.webm': - return 'video/webm'; - default: - return 'video/mp4'; // Default to mp4 - } - } - - String _buildVideoUploadFailureMessage({ - required String fallback, - dynamic detail, - }) { - final normalizedDetail = _extractVideoUploadFailureDetail(detail); - if (normalizedDetail == null) { - return fallback; - } - - final normalizedFallback = fallback.trim(); - if (normalizedDetail.toLowerCase() == normalizedFallback.toLowerCase()) { - return normalizedFallback; - } - if (normalizedDetail.toLowerCase().startsWith( - normalizedFallback.toLowerCase(), - )) { - return normalizedDetail; - } - - final separator = normalizedFallback.endsWith('.') ? ' ' : ': '; - return '$normalizedFallback$separator$normalizedDetail'; - } - - String? _extractVideoUploadFailureDetail(dynamic value) { - if (value == null) { - return null; - } - - if (value is String) { - final trimmed = value.trim(); - if (trimmed.isEmpty) { - return null; - } - - try { - final decoded = jsonDecode(trimmed); - final decodedDetail = _extractVideoUploadFailureDetail(decoded); - if (decodedDetail != null) { - return decodedDetail; - } - } catch (_) { - // Fall back to the raw string when the response is not JSON. - } - - return _sanitizeVideoUploadFailureText(trimmed); - } - - if (value is Map) { - for (final key in const [ - 'message', - 'status', - 'detail', - 'reason', - 'description', - 'error', - ]) { - final nestedDetail = _extractVideoUploadFailureDetail(value[key]); - if (nestedDetail != null) { - return nestedDetail; - } - } - - final jobStatusDetail = _extractVideoUploadFailureDetail( - value['jobStatus'], - ); - if (jobStatusDetail != null) { - return jobStatusDetail; - } - - return null; - } - - if (value is Iterable) { - for (final item in value) { - final itemDetail = _extractVideoUploadFailureDetail(item); - if (itemDetail != null) { - return itemDetail; - } - } - return null; - } - - return _sanitizeVideoUploadFailureText(value.toString()); - } - - String? _sanitizeVideoUploadFailureText(String text) { - final sanitized = text - .replaceFirst( - RegExp(r'^(exception|error):\s*', caseSensitive: false), - '', - ) - .replaceAll(RegExp(r'\s+'), ' ') - .trim(); - - if (sanitized.isEmpty || - sanitized == '{}' || - sanitized == '[]' || - sanitized.startsWith(' value; +} + +class StoryRecordPage { + const StoryRecordPage({required this.records, this.cursor}); + + final List records; + final String? cursor; +} + /// Interface for Story-related API endpoints abstract class StoryRepository { /// Post a story to the user's feed @@ -33,4 +47,11 @@ abstract class StoryRepository { /// /// [storyUris] List of story URIs to fetch Future> getStoryViews(List storyUris); + + Future listStoryRecords({ + required String did, + String? cursor, + }); + + Future deleteStoryRecord(AtUri uri); } diff --git a/lib/src/core/network/atproto/data/repositories/story_repository_impl.dart b/lib/src/core/network/atproto/data/repositories/story_repository_impl.dart index f3acfb6e..b596f648 100644 --- a/lib/src/core/network/atproto/data/repositories/story_repository_impl.dart +++ b/lib/src/core/network/atproto/data/repositories/story_repository_impl.dart @@ -1,4 +1,6 @@ import 'package:poptart_lex/com/atproto/label/defs.dart'; +import 'package:poptart_lex/com/atproto/repo/list_records.dart' + as repo_list_records; import 'package:poptart_lex/com/atproto/repo/strong_ref.dart'; import 'package:poptart/poptart.dart'; import 'package:spark/src/core/network/atproto/data/models/feed_models.dart'; @@ -19,6 +21,40 @@ class StoryRepositoryImpl implements StoryRepository { final SprkRepository _client; final DateTime Function() _now; + @override + Future listStoryRecords({ + required String did, + String? cursor, + }) { + return _client.executeWithRetry(() async { + final atproto = _client.authRepository.atproto; + if (atproto == null) { + throw StateError('AtProto not initialized'); + } + final result = await atproto.call( + repo_list_records.comAtprotoRepoListRecords, + parameters: repo_list_records.RepoListRecordsInput( + repo: did, + collection: 'so.sprk.story.post', + cursor: cursor, + limit: 100, + ), + ); + return StoryRecordPage( + records: [ + for (final record in result.data.records) + StoryRecordEntry(uri: record.uri, value: record.value), + ], + cursor: result.data.cursor, + ); + }); + } + + @override + Future deleteStoryRecord(AtUri uri) { + return _client.repo.deleteRecord(uri: uri); + } + /// Fixes story media JSON to match generated sprk_poptart view models. void _fixMediaStructure(Map storyJson) { final media = storyJson['media']; diff --git a/lib/src/core/network/atproto/data/services/video_upload_service.dart b/lib/src/core/network/atproto/data/services/video_upload_service.dart new file mode 100644 index 00000000..2f0d0b54 --- /dev/null +++ b/lib/src/core/network/atproto/data/services/video_upload_service.dart @@ -0,0 +1,506 @@ +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; + +import 'package:http/http.dart' as http; +import 'package:path/path.dart' as path; +import 'package:poptart/poptart.dart'; +import 'package:poptart_lex/com/atproto/server/get_service_auth.dart' + as server_get_service_auth; +import 'package:spark/src/core/auth/data/repositories/auth_repository.dart'; +import 'package:spark/src/core/config/app_config.dart'; +import 'package:spark/src/core/network/atproto/data/models/models.dart'; +import 'package:spark/src/core/utils/logging/logger.dart'; +import 'package:spark/src/core/utils/video_upload_exception.dart'; +import 'package:sprk_poptart/so/sprk/sound/defs/audio_details.dart'; + +abstract interface class VideoUploadService { + Future uploadVideo( + String videoPath, { + void Function(double progress)? onUploadProgress, + }); +} + +class VideoUploadClient implements VideoUploadService { + factory VideoUploadClient( + AuthRepository authRepository, { + required SparkLogger logger, + DateTime Function()? now, + http.Client Function()? httpClientFactory, + Future Function(Duration)? processingDelay, + File Function(String)? file, + Future Function(PoptartClient)? serviceAuthTokenRequest, + }) { + return VideoUploadClient._( + authRepository, + logger, + now ?? DateTime.now, + httpClientFactory ?? http.Client.new, + processingDelay ?? Future.delayed, + file ?? File.new, + serviceAuthTokenRequest, + ); + } + + VideoUploadClient._( + this._authRepository, + this._logger, + this._now, + this._httpClientFactory, + this._processingDelay, + this._file, + this._serviceAuthTokenRequest, + ); + + final AuthRepository _authRepository; + final SparkLogger _logger; + final DateTime Function() _now; + final http.Client Function() _httpClientFactory; + final Future Function(Duration) _processingDelay; + final File Function(String) _file; + final Future Function(PoptartClient)? _serviceAuthTokenRequest; + + @override + Future uploadVideo( + String videoPath, { + void Function(double progress)? onUploadProgress, + }) async { + _logger.d('Uploading video from path: $videoPath'); + + if (!_authRepository.isAuthenticated) { + _logger.w('Not authenticated'); + throw Exception('Not authenticated'); + } + final authAtProto = _authRepository.atproto; + if (authAtProto == null || authAtProto.oAuthSession == null) { + throw Exception('AtProto not initialized'); + } + + final cleanVideoPath = videoPath.startsWith('file://') + ? videoPath.replaceFirst('file://', '') + : videoPath; + final videoFile = _file(cleanVideoPath); + if (!videoFile.existsSync()) { + throw Exception('Video file not found: $cleanVideoPath'); + } + + final videoSizeBigInt = BigInt.from(await videoFile.length()); + if (videoSizeBigInt == BigInt.zero) { + throw Exception('Video file is empty'); + } + + if (videoSizeBigInt > BigInt.from(2 * 1024 * 1024 * 1024)) { + _logger.w( + 'Video file exceeds 2GB, may cause issues: $videoSizeBigInt bytes', + ); + throw VideoUploadException( + 'Video is too large. Maximum supported size is 2GB.', + statusCode: 413, + uploadSizeBytes: videoSizeBigInt.toInt(), + limitBytes: 2 * 1024 * 1024 * 1024, + ); + } + + final videoSizeBytes = videoSizeBigInt.toInt(); + _logger.i('Video file size: $videoSizeBytes bytes'); + final maxUploadSizeBytes = (AppConfig.maxUploadSizeMB * 1024 * 1024) + .round(); + if (maxUploadSizeBytes > 0 && videoSizeBytes > maxUploadSizeBytes) { + _logger.w( + 'Video file exceeds upload limit: $videoSizeBytes bytes ' + '(limit: $maxUploadSizeBytes bytes)', + ); + throw VideoUploadException( + 'Video is too large to upload.', + statusCode: 413, + uploadSizeBytes: videoSizeBytes, + limitBytes: maxUploadSizeBytes, + ); + } + + if (videoSizeBytes > 2147483647) { + _logger.e('Video file too large for HTTP content-length header'); + throw VideoUploadException( + 'Video is too large to upload.', + statusCode: 413, + uploadSizeBytes: videoSizeBytes, + limitBytes: 2147483647, + ); + } + + var serviceToken = await _createServiceAuthToken(); + final httpClient = _httpClientFactory(); + try { + final uploadRequest = + http.StreamedRequest( + 'POST', + Uri.parse( + '${AppConfig.videoServiceUrl}/xrpc/so.sprk.video.uploadVideo', + ), + ) + ..contentLength = videoSizeBytes + ..headers.addAll({ + 'Authorization': 'Bearer $serviceToken', + 'Content-Type': _getContentType(cleanVideoPath), + }); + + onUploadProgress?.call(0); + final uploadResponseFuture = httpClient.send(uploadRequest); + try { + await uploadRequest.sink.addStream( + _trackUploadProgress( + videoFile.openRead(), + totalBytes: videoSizeBytes, + onUploadProgress: onUploadProgress, + ), + ); + } finally { + unawaited(uploadRequest.sink.close()); + } + var response = await http.Response.fromStream(await uploadResponseFuture); + + if (response.statusCode != 200) { + _logger.e( + 'Video upload failed: ${response.statusCode} ${response.body}', + ); + throw VideoUploadException( + _buildFailureMessage( + fallback: response.statusCode == 413 + ? 'Video is too large to upload.' + : 'Failed to upload video.', + detail: response.body, + ), + statusCode: response.statusCode, + uploadSizeBytes: videoSizeBytes, + limitBytes: maxUploadSizeBytes > 0 ? maxUploadSizeBytes : null, + responseBody: response.body, + ); + } + + dynamic responseData = jsonDecode(response.body); + _logger.d('Video upload response: $responseData'); + + var jobState = responseData['jobStatus']?['state'] as String?; + var attempts = 0; + var consecutivePollErrors = 0; + const maxAttempts = 120; + const maxConsecutivePollErrors = 3; + while (jobState == 'JOB_STATE_QUEUED' || + jobState == 'JOB_STATE_PROCESSING') { + _logger.d('Video upload in progress, status: $jobState'); + await _processingDelay(const Duration(seconds: 2)); + attempts++; + if (attempts > maxAttempts) { + throw const VideoUploadException( + 'Video processing timed out. Please try again.', + ); + } + + try { + response = await httpClient.get( + _jobStatusUri(responseData), + headers: _videoHeaders(serviceToken, cleanVideoPath), + ); + if (_isExpiredTokenResponse(response)) { + _logger.i( + 'Video service token expired while polling; minting a new token', + ); + serviceToken = await _createServiceAuthToken( + refreshPdsSessionOnFailure: true, + ); + response = await httpClient.get( + _jobStatusUri(responseData), + headers: _videoHeaders(serviceToken, cleanVideoPath), + ); + } + if (response.statusCode != 200) { + throw Exception( + 'Failed to check video upload status: ${response.statusCode} ' + '${response.body}', + ); + } + responseData = jsonDecode(response.body); + _logger.d('Video upload status response: $responseData'); + jobState = responseData['jobStatus']?['state'] as String?; + consecutivePollErrors = 0; + } catch (error) { + consecutivePollErrors++; + _logger.w( + 'Error polling video upload status on attempt ' + '$attempts/$maxAttempts ' + '(consecutive errors: ' + '$consecutivePollErrors/$maxConsecutivePollErrors): $error', + ); + + if (consecutivePollErrors >= maxConsecutivePollErrors) { + _logger.e( + 'Too many consecutive polling errors, giving up: $error', + error: error, + ); + throw VideoUploadException( + _buildFailureMessage( + fallback: 'Failed to check video processing status.', + detail: error.toString(), + ), + responseBody: error.toString(), + ); + } + + _logger.d('Retrying poll after error...'); + jobState = 'JOB_STATE_PROCESSING'; + } + } + + if (responseData['jobStatus']?['state'] == 'JOB_STATE_FAILED') { + final failureMessage = _buildFailureMessage( + fallback: 'Video processing failed.', + detail: responseData['jobStatus'] ?? responseData, + ); + _logger.e( + 'Video processing job failed: $failureMessage', + error: responseData, + ); + throw VideoUploadException( + failureMessage, + responseBody: jsonEncode(responseData), + ); + } + + final Map videoBlobData; + if (responseData case {'jobStatus': {'blob': final blobData}}) { + videoBlobData = blobData as Map; + } else if (responseData case {'blobRef': final blobRef}) { + videoBlobData = blobRef as Map; + } else { + throw Exception('Unexpected response format: $responseData'); + } + final videoBlob = Blob.fromJson(videoBlobData); + + Blob? audioBlob; + AudioDetails? audioDetails; + if (responseData case {'jobStatus': {'audio': final audioData}}) { + final audio = audioData as Map; + if (audio['blob'] != null) { + audioBlob = Blob.fromJson(audio['blob'] as Map); + _logger.d('Extracted audio blob: ${audioBlob.size} bytes'); + } + if (audio['details'] != null) { + audioDetails = AudioDetails.fromJson( + audio['details'] as Map, + ); + } + } + + return VideoUploadResult( + videoBlob: videoBlob, + audioBlob: audioBlob, + audioDetails: audioDetails, + ); + } finally { + httpClient.close(); + } + } + + Uri _jobStatusUri(dynamic responseData) { + return Uri.parse( + '${AppConfig.videoServiceUrl}/xrpc/so.sprk.video.getJobStatus', + ).replace( + queryParameters: { + 'jobId': responseData['jobStatus']?['jobId'] as String?, + }, + ); + } + + Map _videoHeaders(String serviceToken, String videoPath) { + return { + 'Authorization': 'Bearer $serviceToken', + 'Content-Type': _getContentType(videoPath), + }; + } + + Stream> _trackUploadProgress( + Stream> chunks, { + required int totalBytes, + void Function(double progress)? onUploadProgress, + }) async* { + var uploadedBytes = 0; + + await for (final chunk in chunks) { + uploadedBytes += chunk.length; + if (totalBytes > 0) { + onUploadProgress?.call( + (uploadedBytes / totalBytes).clamp(0, 1).toDouble(), + ); + } + yield chunk; + } + + onUploadProgress?.call(1); + } + + Future _createServiceAuthToken({ + bool refreshPdsSessionOnFailure = false, + }) async { + final atproto = _authRepository.atproto; + if (atproto == null) { + throw Exception('AtProto not initialized'); + } + + try { + return await _requestServiceAuthToken(atproto); + } catch (error) { + if (!refreshPdsSessionOnFailure) { + rethrow; + } + + _logger.i( + 'Refreshing PDS session before minting video service token', + error: error, + ); + final refreshed = await _authRepository.refreshToken(); + if (!refreshed) { + throw Exception('Session expired. Please log in again.'); + } + + final refreshedAtproto = _authRepository.atproto; + if (refreshedAtproto == null) { + throw Exception('AtProto not initialized after refresh'); + } + return _requestServiceAuthToken(refreshedAtproto); + } + } + + Future _requestServiceAuthToken(PoptartClient atproto) async { + final requestOverride = _serviceAuthTokenRequest; + if (requestOverride != null) { + return requestOverride(atproto); + } + final response = await atproto.call( + server_get_service_auth.comAtprotoServerGetServiceAuth, + parameters: server_get_service_auth.ServerGetServiceAuthInput( + aud: 'did:web:${atproto.service}', + lxm: 'com.atproto.repo.uploadBlob', + exp: + _now() + .toUtc() + .add(const Duration(minutes: 5)) + .millisecondsSinceEpoch ~/ + 1000, + ), + ); + + return response.data.token; + } + + bool _isExpiredTokenResponse(http.Response response) { + if (response.statusCode != 401) { + return false; + } + + final body = response.body.toLowerCase(); + return body.contains('jwt has expired') || body.contains('invalidtoken'); + } + + String _getContentType(String videoPath) { + return switch (path.extension(videoPath).toLowerCase()) { + '.mov' => 'video/quicktime', + '.avi' => 'video/x-msvideo', + '.webm' => 'video/webm', + _ => 'video/mp4', + }; + } + + String _buildFailureMessage({required String fallback, dynamic detail}) { + final normalizedDetail = _extractFailureDetail(detail); + if (normalizedDetail == null) { + return fallback; + } + + final normalizedFallback = fallback.trim(); + if (normalizedDetail.toLowerCase() == normalizedFallback.toLowerCase()) { + return normalizedFallback; + } + if (normalizedDetail.toLowerCase().startsWith( + normalizedFallback.toLowerCase(), + )) { + return normalizedDetail; + } + + final separator = normalizedFallback.endsWith('.') ? ' ' : ': '; + return '$normalizedFallback$separator$normalizedDetail'; + } + + String? _extractFailureDetail(dynamic value) { + if (value == null) { + return null; + } + + if (value is String) { + final trimmed = value.trim(); + if (trimmed.isEmpty) { + return null; + } + + try { + final decodedDetail = _extractFailureDetail(jsonDecode(trimmed)); + if (decodedDetail != null) { + return decodedDetail; + } + } catch (_) { + // Use the raw response when it is not JSON. + } + + return _sanitizeFailureText(trimmed); + } + + if (value is Map) { + for (final key in const [ + 'message', + 'status', + 'detail', + 'reason', + 'description', + 'error', + ]) { + final nestedDetail = _extractFailureDetail(value[key]); + if (nestedDetail != null) { + return nestedDetail; + } + } + + return _extractFailureDetail(value['jobStatus']); + } + + if (value is Iterable) { + for (final item in value) { + final itemDetail = _extractFailureDetail(item); + if (itemDetail != null) { + return itemDetail; + } + } + return null; + } + + return _sanitizeFailureText(value.toString()); + } + + String? _sanitizeFailureText(String text) { + final sanitized = text + .replaceFirst( + RegExp(r'^(exception|error):\s*', caseSensitive: false), + '', + ) + .replaceAll(RegExp(r'\s+'), ' ') + .trim(); + + if (sanitized.isEmpty || + sanitized == '{}' || + sanitized == '[]' || + sanitized.startsWith(' Function() action); - -final soundPickerSearchDebounceSchedulerProvider = - Provider((ref) { - return (delay, action) { - final timer = Timer(delay, () => unawaited(action())); - return timer.cancel; - }; - }); - @riverpod class SoundPickerSearch extends _$SoundPickerSearch { static const int _limit = 25; @@ -65,7 +54,7 @@ class SoundPickerSearch extends _$SoundPickerSearch { error: null, ); - _cancelDebounce = ref.read(soundPickerSearchDebounceSchedulerProvider)( + _cancelDebounce = ref.read(debounceSchedulerProvider)( const Duration(milliseconds: 350), () => _searchAudios(trimmedQuery, requestToken: requestToken, reset: true), diff --git a/lib/src/features/search/providers/search_debounce_scheduler.dart b/lib/src/core/providers/debounce_scheduler.dart similarity index 69% rename from lib/src/features/search/providers/search_debounce_scheduler.dart rename to lib/src/core/providers/debounce_scheduler.dart index 93164af6..39f2df9e 100644 --- a/lib/src/features/search/providers/search_debounce_scheduler.dart +++ b/lib/src/core/providers/debounce_scheduler.dart @@ -2,12 +2,10 @@ import 'dart:async'; import 'package:flutter_riverpod/flutter_riverpod.dart'; -typedef SearchDebounceScheduler = +typedef DebounceScheduler = void Function() Function(Duration delay, Future Function() action); -final searchDebounceSchedulerProvider = Provider(( - ref, -) { +final debounceSchedulerProvider = Provider((ref) { return (delay, action) { final timer = Timer(delay, () => unawaited(action())); return timer.cancel; diff --git a/lib/src/features/feed/providers/feed_provider.dart b/lib/src/features/feed/providers/feed_provider.dart index 0d57b448..ee3dea27 100644 --- a/lib/src/features/feed/providers/feed_provider.dart +++ b/lib/src/features/feed/providers/feed_provider.dart @@ -1,4 +1,3 @@ -import 'dart:async'; import 'dart:collection'; import 'package:flutter_riverpod/flutter_riverpod.dart' show Provider, Ref; @@ -59,7 +58,7 @@ final feedSettingsGatewayProvider = Provider((ref) { @Riverpod(keepAlive: true) class FeedNotifier extends _$FeedNotifier { bool _isLoadingInProgress = false; - bool _isFetching = false; + Future? _inFlightFetch; DateTime? _lastErrorTime; static const _errorCooldown = Duration(seconds: 10); late Feed _feed; @@ -68,7 +67,6 @@ class FeedNotifier extends _$FeedNotifier { late final SparkLogger _logger; late final DownloadManagerInterface _downloadManager; late final FeedSettingsGateway _settingsGateway; - Completer? _fetchCompletion; // Track active fetch operation for cancellation int _fetchGeneration = 0; @@ -374,29 +372,42 @@ class FeedNotifier extends _$FeedNotifier { bool replaceExisting = false, int? generation, }) async { - if (state.isEndOfNetworkFeed) { - return; - } - final activeGeneration = generation ?? _fetchGeneration; - if (_isFetching) { - final fetchCompletion = _fetchCompletion; - if (generation != null && fetchCompletion != null) { - await fetchCompletion.future; - if (ref.mounted && activeGeneration == _fetchGeneration) { - await _maybeFetchNextBatch( - limit: limit, - replaceExisting: replaceExisting, - generation: activeGeneration, - ); + while (!state.isEndOfNetworkFeed) { + final inFlightFetch = _inFlightFetch; + if (inFlightFetch != null) { + if (generation == null) { + return; + } + await inFlightFetch; + if (!ref.mounted || activeGeneration != _fetchGeneration) { + return; + } + continue; + } + + final fetch = _fetchNextBatch( + limit: limit, + replaceExisting: replaceExisting, + generation: activeGeneration, + ); + _inFlightFetch = fetch; + try { + await fetch; + } finally { + if (identical(_inFlightFetch, fetch)) { + _inFlightFetch = null; } } return; } + } - _isFetching = true; - final fetchCompletion = Completer(); - _fetchCompletion = fetchCompletion; + Future _fetchNextBatch({ + required int? limit, + required bool replaceExisting, + required int generation, + }) async { try { var attempts = 0; var consecutiveEmptyResults = 0; @@ -405,16 +416,13 @@ class FeedNotifier extends _$FeedNotifier { while (attempts < maxAttempts && !state.isEndOfNetworkFeed) { attempts++; - // Check if generation has changed (fetch was superseded) - if (activeGeneration != _fetchGeneration) { + if (generation != _fetchGeneration) { _logger.d('Fetch superseded by newer generation, cancelling'); return; } final (:count, :posts, :cursor) = await fetch(limit: limit); - - // Check again after await - if (activeGeneration != _fetchGeneration) { + if (generation != _fetchGeneration) { _logger.d('Fetch superseded after network call, discarding results'); return; } @@ -424,7 +432,7 @@ class FeedNotifier extends _$FeedNotifier { if (fetchedPosts.isEmpty) { if (fetchedCount == 0 || cursor == null) { await endOfNetworkFeed(); - if (ref.mounted && activeGeneration == _fetchGeneration) { + if (ref.mounted && generation == _fetchGeneration) { state = state.copyWith( loadingFirstLoad: false, isEndOfNetworkFeed: true, @@ -436,7 +444,7 @@ class FeedNotifier extends _$FeedNotifier { consecutiveEmptyResults++; if (consecutiveEmptyResults >= maxConsecutiveEmpty) { await endOfNetworkFeed(); - if (ref.mounted && activeGeneration == _fetchGeneration) { + if (ref.mounted && generation == _fetchGeneration) { state = state.copyWith( loadingFirstLoad: false, isEndOfNetworkFeed: true, @@ -445,7 +453,7 @@ class FeedNotifier extends _$FeedNotifier { break; } } - if (ref.mounted && activeGeneration == _fetchGeneration) { + if (ref.mounted && generation == _fetchGeneration) { state = state.copyWith(cursor: cursor, loadingFirstLoad: false); } continue; @@ -455,9 +463,9 @@ class FeedNotifier extends _$FeedNotifier { fetchedPosts, cursor: cursor, replaceExisting: replaceExisting, - generation: activeGeneration, + generation: generation, ); - if (activeGeneration != _fetchGeneration) return; + if (generation != _fetchGeneration) return; if (state.error) return; if (addedPosts) { if (cursor == null) await endOfNetworkFeed(); @@ -472,18 +480,12 @@ class FeedNotifier extends _$FeedNotifier { } } catch (e, stackTrace) { // Only update error state if this generation is still current - if (ref.mounted && activeGeneration == _fetchGeneration) { + if (ref.mounted && generation == _fetchGeneration) { _logger.e('Error prefetching feed: $e', stackTrace: stackTrace); _lastErrorTime = DateTime.now(); state = state.copyWith(error: true, loadingFirstLoad: false); } rethrow; - } finally { - _isFetching = false; - if (!fetchCompletion.isCompleted) fetchCompletion.complete(); - if (identical(_fetchCompletion, fetchCompletion)) { - _fetchCompletion = null; - } } } diff --git a/lib/src/features/search/providers/actor_typeahead_provider.dart b/lib/src/features/search/providers/actor_typeahead_provider.dart index c32219f5..c96b683d 100644 --- a/lib/src/features/search/providers/actor_typeahead_provider.dart +++ b/lib/src/features/search/providers/actor_typeahead_provider.dart @@ -1,10 +1,10 @@ import 'package:get_it/get_it.dart'; import 'package:riverpod_annotation/riverpod_annotation.dart'; import 'package:spark/src/core/network/atproto/data/repositories/actor_repository.dart'; +import 'package:spark/src/core/providers/debounce_scheduler.dart'; import 'package:spark/src/core/utils/logging/log_service.dart'; import 'package:spark/src/core/utils/logging/logger.dart'; import 'package:spark/src/features/search/providers/actor_typeahead_state.dart'; -import 'package:spark/src/features/search/providers/search_debounce_scheduler.dart'; part 'actor_typeahead_provider.g.dart'; @@ -40,7 +40,7 @@ class ActorTypeahead extends _$ActorTypeahead { state = state.copyWith(query: trimmedQuery, isLoading: true, error: null); final requestToken = ++_activeRequestToken; - _cancelDebounce = ref.read(searchDebounceSchedulerProvider)( + _cancelDebounce = ref.read(debounceSchedulerProvider)( const Duration(milliseconds: 300), () => _searchTypeahead( trimmedQuery, diff --git a/lib/src/features/search/providers/post_search_provider.dart b/lib/src/features/search/providers/post_search_provider.dart index 5954be5f..efce358a 100644 --- a/lib/src/features/search/providers/post_search_provider.dart +++ b/lib/src/features/search/providers/post_search_provider.dart @@ -10,11 +10,11 @@ import 'package:spark/src/core/auth/data/repositories/auth_repository.dart'; import 'package:spark/src/core/network/atproto/atproto.dart'; import 'package:spark/src/core/network/atproto/data/models/feed_models.dart'; import 'package:spark/src/core/network/atproto/data/models/pref_models.dart'; +import 'package:spark/src/core/providers/debounce_scheduler.dart'; import 'package:spark/src/core/utils/label_utils.dart'; import 'package:spark/src/core/utils/logging/log_service.dart'; import 'package:spark/src/core/utils/logging/logger.dart'; import 'package:spark/src/features/search/providers/post_search_state.dart'; -import 'package:spark/src/features/search/providers/search_debounce_scheduler.dart'; import 'package:spark/src/features/settings/providers/preferences_provider.dart'; part 'post_search_provider.g.dart'; @@ -168,7 +168,7 @@ class PostSearch extends _$PostSearch { // Debounce the search _cancelDebounce?.call(); final requestToken = ++_activeSearchToken; - _cancelDebounce = ref.read(searchDebounceSchedulerProvider)( + _cancelDebounce = ref.read(debounceSchedulerProvider)( const Duration(milliseconds: 500), () => _searchPosts(trimmedQuery, requestToken: requestToken), ); diff --git a/lib/src/features/search/providers/search_provider.dart b/lib/src/features/search/providers/search_provider.dart index 9097d90d..a30d9021 100644 --- a/lib/src/features/search/providers/search_provider.dart +++ b/lib/src/features/search/providers/search_provider.dart @@ -5,10 +5,10 @@ import 'package:spark/src/core/auth/data/repositories/auth_repository.dart'; import 'package:sprk_poptart/so/sprk/actor/defs.dart'; import 'package:spark/src/core/network/atproto/data/repositories/actor_repository.dart'; import 'package:spark/src/core/network/atproto/data/repositories/graph_repository.dart'; +import 'package:spark/src/core/providers/debounce_scheduler.dart'; import 'package:spark/src/core/utils/logging/log_service.dart'; import 'package:spark/src/core/utils/logging/logger.dart'; import 'package:spark/src/features/search/providers/search_state.dart'; -import 'package:spark/src/features/search/providers/search_debounce_scheduler.dart'; part 'search_provider.g.dart'; @@ -62,7 +62,7 @@ class Search extends _$Search { // Debounce the search _cancelDebounce?.call(); final requestToken = ++_activeSearchToken; - _cancelDebounce = ref.read(searchDebounceSchedulerProvider)( + _cancelDebounce = ref.read(debounceSchedulerProvider)( const Duration(milliseconds: 500), () => _searchUsers(trimmedQuery, requestToken: requestToken), ); diff --git a/lib/src/features/stories/providers/stories_by_author.dart b/lib/src/features/stories/providers/stories_by_author.dart index ac67d793..0532b3ab 100644 --- a/lib/src/features/stories/providers/stories_by_author.dart +++ b/lib/src/features/stories/providers/stories_by_author.dart @@ -10,8 +10,8 @@ FutureOr< ({Map> storiesByAuthor, String? cursor}) > storiesByAuthor(Ref ref, {int limit = 30, String? cursor}) async { - final dependencies = ref.read(storyProviderDependenciesProvider); - final result = await dependencies.storyRepository.getStoriesTimeline( + final repository = ref.read(storyRepositoryProvider); + final result = await repository.getStoriesTimeline( limit: limit, cursor: cursor, ); diff --git a/lib/src/features/stories/providers/story_auto_delete_provider.dart b/lib/src/features/stories/providers/story_auto_delete_provider.dart index 90b67fea..60872116 100644 --- a/lib/src/features/stories/providers/story_auto_delete_provider.dart +++ b/lib/src/features/stories/providers/story_auto_delete_provider.dart @@ -1,5 +1,6 @@ import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:riverpod_annotation/riverpod_annotation.dart'; +import 'package:spark/src/core/network/atproto/data/repositories/story_repository.dart'; import 'package:spark/src/core/storage/preferences/storage_constants.dart'; import 'package:spark/src/features/stories/providers/story_manager_provider.dart'; import 'package:spark/src/features/stories/providers/story_provider_dependencies.dart'; @@ -39,17 +40,17 @@ Future storyAutoDeleteExecutor(Ref ref) async { final enabledAsync = await ref.watch(storyAutoDeletePrefProvider.future); if (!enabledAsync) return; - final dependencies = ref.read(storyProviderDependenciesProvider); - final logger = dependencies.loggerFor('StoryAutoDeleteExec'); - final did = dependencies.did; - if (!dependencies.atprotoAvailable || did == null) return; + final repository = ref.read(storyRepositoryProvider); + final logger = ref.read(storyLoggerProvider('StoryAutoDeleteExec')); + final did = ref.read(storyCurrentDidProvider); + if (!ref.read(storyAtprotoAvailableProvider) || did == null) return; try { String? cursor; final expiredUris = []; final now = ref.read(storyClockProvider)().toUtc(); do { - final page = await dependencies.loadRecordPage(did: did, cursor: cursor); + final page = await repository.listStoryRecords(did: did, cursor: cursor); for (final rec in page.records) { final createdAt = rec.value['createdAt']; DateTime? ts; @@ -67,7 +68,7 @@ Future storyAutoDeleteExecutor(Ref ref) async { for (final record in expiredUris) { try { - await dependencies.deleteRecord(record.uri); + await repository.deleteStoryRecord(record.uri); } catch (e) { logger.w('Failed deleting expired story ${record.uri}', error: e); } diff --git a/lib/src/features/stories/providers/story_manager_provider.dart b/lib/src/features/stories/providers/story_manager_provider.dart index 1bf4352c..d2428b06 100644 --- a/lib/src/features/stories/providers/story_manager_provider.dart +++ b/lib/src/features/stories/providers/story_manager_provider.dart @@ -1,6 +1,7 @@ import 'package:poptart/poptart.dart'; import 'package:riverpod_annotation/riverpod_annotation.dart'; import 'package:spark/src/core/network/atproto/data/models/feed_models.dart'; +import 'package:spark/src/core/network/atproto/data/repositories/story_repository.dart'; import 'package:spark/src/core/utils/logging/logger.dart'; import 'package:spark/src/features/stories/providers/story_auto_delete_provider.dart'; import 'package:spark/src/features/stories/providers/story_provider_dependencies.dart'; @@ -33,25 +34,25 @@ class StoryManagerState { @riverpod class StoryManager extends _$StoryManager { - late final StoryProviderDependencies _dependencies; + late final StoryRepository _repository; late final SparkLogger _logger; @override Future build() async { - _dependencies = ref.read(storyProviderDependenciesProvider); - _logger = _dependencies.loggerFor('StoryManager'); + _repository = ref.read(storyRepositoryProvider); + _logger = ref.read(storyLoggerProvider('StoryManager')); ref.read(storyAutoDeleteExecutorProvider.future).ignore(); return _loadInitial(); } Future _loadInitial() async { try { - final did = _dependencies.did; + final did = ref.read(storyCurrentDidProvider); if (did == null) { return StoryManagerState(stories: const [], error: 'Not authenticated'); } // Page through all story records directly via atproto to include expired - if (!_dependencies.atprotoAvailable) { + if (!ref.read(storyAtprotoAvailableProvider)) { return StoryManagerState( stories: const [], error: 'AtProto not initialized', @@ -60,7 +61,7 @@ class StoryManager extends _$StoryManager { String? cursor; final uris = []; do { - final result = await _dependencies.loadRecordPage( + final result = await _repository.listStoryRecords( did: did, cursor: cursor, ); @@ -72,9 +73,7 @@ class StoryManager extends _$StoryManager { if (uris.isEmpty) { return StoryManagerState(stories: const []); } - final storyViews = await _dependencies.storyRepository.getStoryViews( - uris, - ); + final storyViews = await _repository.getStoryViews(uris); storyViews.sort((a, b) => b.indexedAt.compareTo(a.indexedAt)); @@ -98,7 +97,7 @@ class StoryManager extends _$StoryManager { final updatedList = List.from(current.stories) ..removeWhere((s) => s.uri == story.uri); state = AsyncData(current.copyWith(stories: updatedList)); - await _dependencies.deleteRecord(story.uri); + await _repository.deleteStoryRecord(story.uri); } catch (e, s) { _logger.e('Error deleting story', error: e, stackTrace: s); // Revert by refreshing fully diff --git a/lib/src/features/stories/providers/story_provider_dependencies.dart b/lib/src/features/stories/providers/story_provider_dependencies.dart index 9a6b1f0c..aebff20b 100644 --- a/lib/src/features/stories/providers/story_provider_dependencies.dart +++ b/lib/src/features/stories/providers/story_provider_dependencies.dart @@ -1,85 +1,25 @@ import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:get_it/get_it.dart'; -import 'package:poptart/poptart.dart'; -import 'package:poptart_lex/com/atproto/repo/list_records.dart' - as repo_list_records; import 'package:spark/src/core/network/atproto/atproto.dart'; import 'package:spark/src/core/storage/preferences/local_storage_interface.dart'; import 'package:spark/src/core/storage/preferences/storage_manager.dart'; import 'package:spark/src/core/utils/logging/log_service.dart'; import 'package:spark/src/core/utils/logging/logger.dart'; -class StoryRecordEntry { - const StoryRecordEntry({required this.uri, required this.value}); - - final AtUri uri; - final Map value; -} - -class StoryRecordPage { - const StoryRecordPage({required this.records, this.cursor}); - - final List records; - final String? cursor; -} - -typedef StoryRecordPageLoader = - Future Function({required String did, String? cursor}); - -class StoryProviderDependencies { - const StoryProviderDependencies({ - required this.readDid, - required this.readAtprotoAvailable, - required this.loadRecordPage, - required this.storyRepository, - required this.deleteRecord, - required this.loggerFor, - }); +final storyRepositoryProvider = Provider((ref) { + return GetIt.instance(); +}); - final String? Function() readDid; - final bool Function() readAtprotoAvailable; - final StoryRecordPageLoader loadRecordPage; - final StoryRepository storyRepository; - final Future Function(AtUri uri) deleteRecord; - final SparkLogger Function(String name) loggerFor; +final storyCurrentDidProvider = Provider((ref) { + return GetIt.instance().authRepository.did; +}); - String? get did => readDid(); - bool get atprotoAvailable => readAtprotoAvailable(); -} +final storyAtprotoAvailableProvider = Provider((ref) { + return GetIt.instance().authRepository.atproto != null; +}); -final storyProviderDependenciesProvider = Provider(( - ref, -) { - final sprk = GetIt.instance(); - return StoryProviderDependencies( - readDid: () => sprk.authRepository.did, - readAtprotoAvailable: () => sprk.authRepository.atproto != null, - loadRecordPage: ({required did, cursor}) async { - final atproto = sprk.authRepository.atproto; - if (atproto == null) { - throw StateError('AtProto not initialized'); - } - final result = await atproto.call( - repo_list_records.comAtprotoRepoListRecords, - parameters: repo_list_records.RepoListRecordsInput( - repo: did, - collection: 'so.sprk.story.post', - cursor: cursor, - limit: 100, - ), - ); - return StoryRecordPage( - records: [ - for (final record in result.data.records) - StoryRecordEntry(uri: record.uri, value: record.value), - ], - cursor: result.data.cursor, - ); - }, - storyRepository: GetIt.instance(), - deleteRecord: (uri) => sprk.repo.deleteRecord(uri: uri), - loggerFor: GetIt.instance().getLogger, - ); +final storyLoggerProvider = Provider.family((ref, name) { + return GetIt.instance().getLogger(name); }); final storyAutoDeletePreferencesProvider = Provider(( diff --git a/test/src/core/network/atproto/data/repositories/feed_repository_impl_test.dart b/test/src/core/network/atproto/data/repositories/feed_repository_impl_test.dart deleted file mode 100644 index 17f5eaa2..00000000 --- a/test/src/core/network/atproto/data/repositories/feed_repository_impl_test.dart +++ /dev/null @@ -1,1039 +0,0 @@ -import 'dart:convert'; -import 'dart:io'; -import 'dart:typed_data'; - -import 'package:flutter_test/flutter_test.dart'; -import 'package:http/http.dart' as http; -import 'package:poptart/poptart.dart'; -import 'package:poptart_lex/com/atproto/repo/strong_ref.dart'; -import 'package:spark/src/core/auth/data/repositories/auth_repository.dart'; -import 'package:spark/src/core/network/atproto/data/models/feed_models.dart'; -import 'package:spark/src/core/network/atproto/data/models/pref_models.dart'; -import 'package:spark/src/core/network/atproto/data/repositories/feed_repository_impl.dart'; -import 'package:spark/src/core/network/atproto/data/repositories/repo_repository.dart'; -import 'package:spark/src/core/network/atproto/data/repositories/sprk_repository.dart'; -import 'package:spark/src/core/utils/logging/logger.dart'; -import 'package:spark/src/core/utils/video_upload_exception.dart'; -import 'package:sprk_poptart/so/sprk/actor/defs/profile_view_basic.dart'; - -void main() { - final fixedNow = DateTime.utc(2026, 7, 22, 12, 34, 56); - - group('FeedRepositoryImpl feed requests', () { - test('getFeed routes timeline requests and maps hydrated posts', () async { - final post = _postView(); - final harness = _Harness( - now: fixedNow, - response: { - 'cursor': 'timeline-next', - 'feed': [ - {'post': post.toJson()}, - ], - }, - ); - final feed = Feed( - type: 'timeline', - config: makeSavedFeed( - id: 'timeline', - type: 'timeline', - value: 'timeline', - pinned: true, - ), - ); - - final result = await harness.repository.getFeed( - feed, - limit: 12, - cursor: 'timeline-cursor', - labelerDids: const ['did:plc:one', 'did:plc:two'], - ); - - final request = harness.transport.singleRequest; - expect(request.uri.path, '/xrpc/so.sprk.feed.getTimeline'); - expect(request.uri.queryParameters['limit'], '12'); - expect(request.uri.queryParameters['cursor'], 'timeline-cursor'); - expect(request.headers['atproto-proxy'], _Harness.sprkDid); - expect( - request.headers['atproto-accept-labelers'], - 'did:plc:one,did:plc:two', - ); - expect(result.cursor, 'timeline-next'); - expect(result.feed.single.post.uri, post.uri); - }); - - test( - 'getFeedView sends Spark feed requests to the Spark service', - () async { - final post = _postView(); - final harness = _Harness( - now: fixedNow, - response: { - 'cursor': 'spark-next', - 'feed': [ - {'post': post.toJson()}, - ], - }, - ); - final feedUri = AtUri( - 'at://did:plc:generator/so.sprk.feed.generator/spark-feed', - ); - - final result = await harness.repository.getFeedView( - feedUri, - limit: 8, - cursor: 'spark-cursor', - labelerDids: const ['did:plc:labeler'], - ); - - final request = harness.transport.singleRequest; - expect(request.uri.path, '/xrpc/so.sprk.feed.getFeed'); - expect(request.uri.queryParameters['feed'], feedUri.toString()); - expect(request.uri.queryParameters['limit'], '8'); - expect(request.uri.queryParameters['cursor'], 'spark-cursor'); - expect(request.headers['atproto-proxy'], _Harness.sprkDid); - expect(request.headers['atproto-accept-labelers'], 'did:plc:labeler'); - expect(result.cursor, 'spark-next'); - expect(result.feed.single.post.uri, post.uri); - }, - ); - - test( - 'getFeedView selects Bluesky and omits Spark labeler headers', - () async { - final harness = _Harness( - now: fixedNow, - response: { - 'cursor': 'bsky-next', - 'feed': [], - }, - ); - final feedUri = AtUri( - 'at://did:plc:generator/app.bsky.feed.generator/bsky-feed', - ); - - final result = await harness.repository.getFeedView( - feedUri, - limit: 6, - cursor: 'bsky-cursor', - labelerDids: const ['did:plc:labeler'], - ); - - final request = harness.transport.singleRequest; - expect(request.uri.path, '/xrpc/app.bsky.feed.getFeed'); - expect(request.uri.queryParameters['feed'], feedUri.toString()); - expect(request.uri.queryParameters['limit'], '6'); - expect(request.uri.queryParameters['cursor'], 'bsky-cursor'); - expect(request.headers['atproto-proxy'], _Harness.bskyDid); - expect(request.headers, isNot(contains('atproto-accept-labelers'))); - expect(result.cursor, 'bsky-next'); - expect(result.feed, isEmpty); - }, - ); - - test('getFeedView rejects unauthenticated requests before transport', () { - final harness = _Harness( - now: fixedNow, - response: {'feed': []}, - authenticated: false, - ); - - expect( - harness.repository.getFeedView( - AtUri('at://did:plc:generator/so.sprk.feed.generator/feed'), - ), - throwsA( - isA().having( - (error) => error.toString(), - 'message', - contains('Not authenticated'), - ), - ), - ); - expect(harness.transport.requests, isEmpty); - }); - - test('getFeedView propagates the transport 500 response', () async { - final harness = _Harness( - now: fixedNow, - statusCode: 500, - response: { - 'error': 'InternalServerError', - 'message': 'feed unavailable', - }, - ); - - await expectLater( - harness.repository.getFeedView( - AtUri('at://did:plc:generator/so.sprk.feed.generator/feed'), - ), - throwsA( - isA().having( - (error) => error.toString(), - 'message', - allOf( - contains('InternalServerError'), - contains('feed unavailable'), - ), - ), - ), - ); - expect(harness.transport.requests, hasLength(1)); - expect( - harness.transport.singleRequest.uri.path, - '/xrpc/so.sprk.feed.getFeed', - ); - }); - }); - - group('FeedRepositoryImpl records', () { - test( - 'likePost selects the record collection and uses the injected clock', - () async { - final harness = _Harness(now: fixedNow); - final sparkPost = AtUri( - 'at://did:plc:author/so.sprk.feed.post/spark-post', - ); - final bskyPost = AtUri( - 'at://did:plc:author/app.bsky.feed.post/bsky-post', - ); - - await harness.repository.likePost('spark-cid', sparkPost); - await harness.repository.likePost('bsky-cid', bskyPost); - - expect(harness.repo.createCalls, hasLength(2)); - final sparkCall = harness.repo.createCalls[0]; - expect(sparkCall.collection, 'so.sprk.feed.like'); - expect(sparkCall.record[r'$type'], 'so.sprk.feed.like'); - final sparkSubject = - sparkCall.record['subject'] as Map; - expect(sparkSubject[r'$type'], 'com.atproto.repo.strongRef'); - expect(sparkSubject['uri'], sparkPost.toString()); - expect(sparkSubject['cid'], 'spark-cid'); - expect(sparkCall.record['createdAt'], fixedNow.toIso8601String()); - - final bskyCall = harness.repo.createCalls[1]; - expect(bskyCall.collection, 'app.bsky.feed.like'); - expect(bskyCall.record[r'$type'], 'app.bsky.feed.like'); - final bskySubject = bskyCall.record['subject'] as Map; - expect(bskySubject[r'$type'], 'com.atproto.repo.strongRef'); - expect(bskySubject['uri'], bskyPost.toString()); - expect(bskySubject['cid'], 'bsky-cid'); - expect(bskyCall.record['createdAt'], fixedNow.toIso8601String()); - }, - ); - - test( - 'unlikePost deletes the interaction without crosspost cleanup', - () async { - final harness = _Harness(now: fixedNow); - final likeUri = AtUri( - 'at://did:plc:viewer/so.sprk.feed.like/interaction', - ); - - await harness.repository.unlikePost(likeUri); - - expect(harness.repo.deleteCalls, hasLength(1)); - expect(harness.repo.deleteCalls.single.uri, likeUri); - expect( - harness.repo.deleteCalls.single.skipBskyCrosspostCleanup, - isTrue, - ); - }, - ); - - test('postComment writes Spark reply roots and parents', () async { - final harness = _Harness(now: fixedNow); - final parentUri = AtUri('at://did:plc:author/so.sprk.feed.post/parent'); - final rootUri = AtUri('at://did:plc:author/so.sprk.feed.post/root'); - - await harness.repository.postComment( - 'A Spark reply', - 'parent-cid', - parentUri, - rootCid: 'root-cid', - rootUri: rootUri, - ); - - final call = harness.repo.createCalls.single; - expect(call.collection, 'so.sprk.feed.reply'); - expect(call.record[r'$type'], 'so.sprk.feed.reply'); - expect(call.record['text'], 'A Spark reply'); - expect(call.record, isNot(contains('facets'))); - expect(call.record['createdAt'], fixedNow.toIso8601String()); - final reply = call.record['reply'] as Map; - expect(reply[r'$type'], 'so.sprk.feed.reply#replyRef'); - _expectStrongRef(reply['root'], uri: rootUri, cid: 'root-cid'); - _expectStrongRef(reply['parent'], uri: parentUri, cid: 'parent-cid'); - }); - - test( - 'postComment writes Bluesky posts and defaults root to parent', - () async { - final harness = _Harness(now: fixedNow); - final parentUri = AtUri( - 'at://did:web:sprk.so/app.bsky.feed.post/parent', - ); - - await harness.repository.postComment( - 'A Bluesky reply', - 'parent-cid', - parentUri, - ); - - final call = harness.repo.createCalls.single; - expect(call.collection, 'app.bsky.feed.post'); - expect(call.record[r'$type'], 'app.bsky.feed.post'); - expect(call.record['text'], 'A Bluesky reply'); - expect(call.record['createdAt'], fixedNow.toIso8601String()); - final reply = call.record['reply'] as Map; - expect(reply[r'$type'], 'app.bsky.feed.post#replyRef'); - _expectStrongRef(reply['root'], uri: parentUri, cid: 'parent-cid'); - _expectStrongRef(reply['parent'], uri: parentUri, cid: 'parent-cid'); - }, - ); - - test('record creation errors propagate through the repository', () { - final harness = _Harness(now: fixedNow) - ..repo.createError = StateError('record write failed'); - - expect( - harness.repository.likePost( - 'cid', - AtUri('at://did:plc:author/so.sprk.feed.post/post'), - ), - throwsA( - isA().having( - (error) => error.message, - 'message', - 'record write failed', - ), - ), - ); - }); - }); - - group('FeedRepositoryImpl threads', () { - test( - 'getThread sends Spark thread parameters and maps not-found posts', - () async { - final threadUri = AtUri( - 'at://did:plc:author/so.sprk.feed.post/missing', - ); - final harness = _Harness( - now: fixedNow, - response: { - 'thread': [ - { - r'$type': 'so.sprk.feed.getPostThread#threadItem', - 'uri': threadUri.toString(), - 'depth': 0, - 'value': { - r'$type': 'so.sprk.feed.defs#notFoundPost', - 'uri': threadUri.toString(), - 'notFound': true, - }, - }, - ], - }, - ); - - final result = await harness.repository.getThread( - threadUri, - depth: 4, - parentHeight: 2, - ); - - final request = harness.transport.singleRequest; - expect(request.uri.path, '/xrpc/so.sprk.feed.getPostThread'); - expect(request.uri.queryParameters['anchor'], threadUri.toString()); - expect(request.uri.queryParameters['depth'], '4'); - expect(request.uri.queryParameters['parentHeight'], '2'); - expect(request.headers['atproto-proxy'], _Harness.sprkDid); - expect(result, isA()); - expect((result as NotFoundPost).uri, threadUri); - expect(result.notFound, isTrue); - }, - ); - - test('getThread rejects an uninitialized AtProto client', () { - final harness = _Harness(now: fixedNow, atprotoInitialized: false); - - expect( - harness.repository.getThread( - AtUri('at://did:plc:author/so.sprk.feed.post/post'), - ), - throwsA( - isA().having( - (error) => error.toString(), - 'message', - contains('AtProto not initialized'), - ), - ), - ); - expect(harness.transport.requests, isEmpty); - }); - }); - - group('FeedRepositoryImpl video upload', () { - test('rejects missing, empty, and oversized files before HTTP', () async { - final videoClient = _VideoClient(); - final files = { - '/missing.mp4': _FakeFile('/missing.mp4', exists: false), - '/empty.mp4': _FakeFile('/empty.mp4'), - '/huge.mp4': _FakeFile( - '/huge.mp4', - lengthOverride: 2 * 1024 * 1024 * 1024 + 1, - ), - }; - final harness = _Harness( - now: fixedNow, - oauth: true, - videoHttpClient: videoClient, - videoFile: (path) => files[path]!, - videoServiceAuthTokenRequest: (_) async => 'service-token', - ); - - await expectLater( - harness.repository.uploadVideo('/missing.mp4'), - throwsA( - isA().having( - (error) => error.toString(), - 'message', - contains('Video file not found'), - ), - ), - ); - await expectLater( - harness.repository.uploadVideo('/empty.mp4'), - throwsA( - isA().having( - (error) => error.toString(), - 'message', - contains('Video file is empty'), - ), - ), - ); - await expectLater( - harness.repository.uploadVideo('/huge.mp4'), - throwsA( - isA() - .having((error) => error.statusCode, 'statusCode', 413) - .having( - (error) => error.uploadSizeBytes, - 'uploadSizeBytes', - 2 * 1024 * 1024 * 1024 + 1, - ) - .having( - (error) => error.limitBytes, - 'limitBytes', - 2 * 1024 * 1024 * 1024, - ), - ), - ); - expect(videoClient.requests, isEmpty); - }); - - test( - 'streams bytes, reports progress, and parses an immediate blob', - () async { - final videoClient = _VideoClient() - ..enqueueJson(200, { - 'blobRef': _blobJson(mimeType: 'video/quicktime', size: 6), - }); - final progress = []; - final harness = _Harness( - now: fixedNow, - oauth: true, - videoHttpClient: videoClient, - videoFile: (_) => _FakeFile( - '/clip.mov', - chunks: const [ - [1, 2], - [3, 4, 5, 6], - ], - ), - videoServiceAuthTokenRequest: (_) async => 'service-token', - ); - - final result = await harness.repository.uploadVideo( - 'file:///clip.mov', - onUploadProgress: progress.add, - ); - - final request = videoClient.requests.single; - expect(request.method, 'POST'); - expect(request.url.path, '/xrpc/so.sprk.video.uploadVideo'); - expect(request.headers['Authorization'], 'Bearer service-token'); - expect(request.headers['Content-Type'], 'video/quicktime'); - expect(request.bodyBytes, [1, 2, 3, 4, 5, 6]); - expect(progress, [0, closeTo(1 / 3, 0.0001), 1, 1]); - expect(result.videoBlob.mimeType, 'video/quicktime'); - expect(result.videoBlob.size, 6); - expect(result.audioBlob, isNull); - expect(result.audioDetails, isNull); - }, - ); - - test('maps upload rejection details and size metadata', () async { - final videoClient = _VideoClient() - ..enqueueJson(413, { - 'error': 'PayloadTooLarge', - 'message': 'limit is 100 MB', - }); - final harness = _Harness( - now: fixedNow, - oauth: true, - videoHttpClient: videoClient, - videoFile: (_) => _FakeFile( - '/clip.mp4', - chunks: const [ - [1, 2, 3], - ], - ), - videoServiceAuthTokenRequest: (_) async => 'service-token', - ); - - await expectLater( - harness.repository.uploadVideo('/clip.mp4'), - throwsA( - isA() - .having((error) => error.statusCode, 'statusCode', 413) - .having((error) => error.uploadSizeBytes, 'uploadSizeBytes', 3) - .having((error) => error.isPayloadTooLarge, 'too large', isTrue) - .having( - (error) => error.message, - 'message', - contains('limit is 100 MB'), - ), - ), - ); - expect(videoClient.requests, hasLength(1)); - }); - - test( - 'polls deterministically and parses video plus audio output', - () async { - final videoClient = _VideoClient() - ..enqueueJson(200, { - 'jobStatus': {'jobId': 'job-1', 'state': 'JOB_STATE_QUEUED'}, - }) - ..enqueueJson(200, { - 'jobStatus': {'jobId': 'job-1', 'state': 'JOB_STATE_PROCESSING'}, - }) - ..enqueueJson(200, { - 'jobStatus': { - 'jobId': 'job-1', - 'state': 'JOB_STATE_COMPLETED', - 'blob': _blobJson(mimeType: 'video/mp4', size: 3), - 'audio': { - 'blob': _blobJson(mimeType: 'audio/aac', size: 2), - 'details': { - r'$type': 'so.sprk.sound.defs#audioDetails', - 'artist': 'Artist', - 'title': 'Title', - }, - }, - }, - }); - final delays = []; - final harness = _Harness( - now: fixedNow, - oauth: true, - videoHttpClient: videoClient, - videoFile: (_) => _FakeFile( - '/clip.mp4', - chunks: const [ - [1, 2, 3], - ], - ), - videoProcessingDelay: (duration) async => delays.add(duration), - videoServiceAuthTokenRequest: (_) async => 'service-token', - ); - - final result = await harness.repository.uploadVideo('/clip.mp4'); - - expect(delays, const [Duration(seconds: 2), Duration(seconds: 2)]); - expect(videoClient.requests, hasLength(3)); - for (final request in videoClient.requests.skip(1)) { - expect(request.method, 'GET'); - expect(request.url.path, '/xrpc/so.sprk.video.getJobStatus'); - expect(request.url.queryParameters['jobId'], 'job-1'); - expect(request.headers['Authorization'], 'Bearer service-token'); - } - expect(result.videoBlob.mimeType, 'video/mp4'); - expect(result.audioBlob?.mimeType, 'audio/aac'); - expect(result.audioDetails?.artist, 'Artist'); - expect(result.audioDetails?.title, 'Title'); - }, - ); - - test( - 'refreshes an expired polling token and retries the request', - () async { - final videoClient = _VideoClient() - ..enqueueJson(200, { - 'jobStatus': {'jobId': 'job-1', 'state': 'JOB_STATE_PROCESSING'}, - }) - ..enqueueJson(401, {'message': 'JWT has expired'}) - ..enqueueJson(200, { - 'jobStatus': { - 'jobId': 'job-1', - 'state': 'JOB_STATE_COMPLETED', - 'blob': _blobJson(mimeType: 'video/mp4', size: 3), - }, - }); - var tokenRequests = 0; - final harness = _Harness( - now: fixedNow, - oauth: true, - videoHttpClient: videoClient, - videoFile: (_) => _FakeFile( - '/clip.mp4', - chunks: const [ - [1, 2, 3], - ], - ), - videoProcessingDelay: (_) async {}, - videoServiceAuthTokenRequest: (_) async { - tokenRequests++; - if (tokenRequests == 2) { - throw StateError('PDS token expired'); - } - return tokenRequests == 1 ? 'old-token' : 'new-token'; - }, - ); - - await harness.repository.uploadVideo('/clip.mp4'); - - expect(tokenRequests, 3); - expect(harness.auth.refreshTokenCalls, 1); - expect(videoClient.requests, hasLength(3)); - expect( - videoClient.requests[1].headers['Authorization'], - 'Bearer old-token', - ); - expect( - videoClient.requests[2].headers['Authorization'], - 'Bearer new-token', - ); - }, - ); - - test('fails after three consecutive polling errors', () async { - final videoClient = _VideoClient() - ..enqueueJson(200, { - 'jobStatus': {'jobId': 'job-1', 'state': 'JOB_STATE_PROCESSING'}, - }) - ..enqueueJson(500, {'message': 'first'}) - ..enqueueJson(502, {'message': 'second'}) - ..enqueueJson(503, {'message': 'third'}); - var delays = 0; - final harness = _Harness( - now: fixedNow, - oauth: true, - videoHttpClient: videoClient, - videoFile: (_) => _FakeFile( - '/clip.mp4', - chunks: const [ - [1], - ], - ), - videoProcessingDelay: (_) async => delays++, - videoServiceAuthTokenRequest: (_) async => 'service-token', - ); - - await expectLater( - harness.repository.uploadVideo('/clip.mp4'), - throwsA( - isA().having( - (error) => error.message, - 'message', - contains('third'), - ), - ), - ); - expect(delays, 3); - expect(videoClient.requests, hasLength(4)); - }); - - test( - 'fails failed jobs and times out processing without wall-clock waits', - () async { - final failedClient = _VideoClient() - ..enqueueJson(200, { - 'jobStatus': { - 'jobId': 'failed-job', - 'state': 'JOB_STATE_FAILED', - 'error': 'TranscodeFailed', - 'message': 'unsupported codec', - }, - }); - final failedHarness = _Harness( - now: fixedNow, - oauth: true, - videoHttpClient: failedClient, - videoFile: (_) => _FakeFile( - '/clip.mp4', - chunks: const [ - [1], - ], - ), - videoServiceAuthTokenRequest: (_) async => 'service-token', - ); - await expectLater( - failedHarness.repository.uploadVideo('/clip.mp4'), - throwsA( - isA().having( - (error) => error.message, - 'message', - contains('unsupported codec'), - ), - ), - ); - - final timeoutClient = _VideoClient() - ..enqueueJson(200, { - 'jobStatus': {'jobId': 'slow-job', 'state': 'JOB_STATE_PROCESSING'}, - }); - for (var i = 0; i < 120; i++) { - timeoutClient.enqueueJson(200, { - 'jobStatus': {'jobId': 'slow-job', 'state': 'JOB_STATE_PROCESSING'}, - }); - } - var delays = 0; - final timeoutHarness = _Harness( - now: fixedNow, - oauth: true, - videoHttpClient: timeoutClient, - videoFile: (_) => _FakeFile( - '/clip.mp4', - chunks: const [ - [1], - ], - ), - videoProcessingDelay: (_) async => delays++, - videoServiceAuthTokenRequest: (_) async => 'service-token', - ); - await expectLater( - timeoutHarness.repository.uploadVideo('/clip.mp4'), - throwsA( - isA().having( - (error) => error.message, - 'message', - contains('timed out'), - ), - ), - ); - expect(delays, 121); - expect(timeoutClient.requests, hasLength(121)); - }, - ); - }); -} - -Map _blobJson({required String mimeType, required int size}) => - { - r'$type': 'blob', - 'mimeType': mimeType, - 'size': size, - 'ref': {r'$link': 'bafkreigh2akiscaildc2'}, - }; - -PostView _postView() { - return PostView( - uri: AtUri('at://did:plc:author/so.sprk.feed.post/post'), - cid: 'post-cid', - author: const ProfileViewBasic( - did: 'did:plc:author', - handle: 'author.test', - ), - record: const { - r'$type': 'so.sprk.feed.post', - 'caption': {'text': 'A post', 'facets': []}, - 'createdAt': '2026-07-22T12:00:00.000Z', - }, - indexedAt: DateTime.utc(2026, 7, 22, 12), - ); -} - -void _expectStrongRef( - Object? value, { - required AtUri uri, - required String cid, -}) { - final ref = value as Map; - expect(ref[r'$type'], 'com.atproto.repo.strongRef'); - expect(ref['uri'], uri.toString()); - expect(ref['cid'], cid); -} - -class _Harness { - _Harness({ - required DateTime now, - Map response = const { - 'feed': [], - }, - int statusCode = 200, - bool authenticated = true, - bool atprotoInitialized = true, - bool oauth = false, - http.Client? videoHttpClient, - Future Function(Duration)? videoProcessingDelay, - File Function(String)? videoFile, - Future Function(PoptartClient)? videoServiceAuthTokenRequest, - }) : transport = _Transport(response: response, statusCode: statusCode), - repo = _FakeRepoRepository() { - final client = oauth - ? PoptartClient.fromOAuthSession( - restoreOAuthSession( - accessToken: 'opaque-access-token', - refreshToken: 'opaque-refresh-token', - scope: 'atproto', - expiresAt: DateTime.utc(2030), - sub: 'did:plc:viewer', - clientId: 'https://spark.test/client-metadata.json', - pdsEndpoint: 'pds.test', - publicKey: 'unused-public-key', - privateKey: 'unused-private-key', - ), - service: 'pds.test', - getClient: transport.get, - ) - : PoptartClient.anonymous( - service: 'pds.test', - getClient: transport.get, - ); - auth = _FakeAuthRepository( - authenticated: authenticated, - atproto: atprotoInitialized ? client : null, - ); - sprk = _FakeSprkRepository(auth: auth, repo: repo); - repository = FeedRepositoryImpl( - sprk, - logger: SparkLogger(), - now: () => now, - videoHttpClient: videoHttpClient, - videoProcessingDelay: videoProcessingDelay, - videoFile: videoFile, - videoServiceAuthTokenRequest: videoServiceAuthTokenRequest, - ); - } - - static const sprkDid = 'did:web:sprk.test'; - static const bskyDid = 'did:web:bsky.test'; - - final _Transport transport; - final _FakeRepoRepository repo; - late final _FakeAuthRepository auth; - late final _FakeSprkRepository sprk; - late final FeedRepositoryImpl repository; -} - -class _Transport { - _Transport({required this.response, required this.statusCode}); - - final Map response; - final int statusCode; - final List<_Request> requests = []; - - _Request get singleRequest => requests.single; - - Future get(Uri uri, {Map? headers}) async { - requests.add(_Request(uri: uri, headers: headers ?? const {})); - return http.Response( - jsonEncode(response), - statusCode, - headers: const {'content-type': 'application/json'}, - request: http.Request('GET', uri), - ); - } -} - -class _Request { - const _Request({required this.uri, required this.headers}); - - final Uri uri; - final Map headers; -} - -class _VideoClient extends http.BaseClient { - final List<_QueuedResponse> _responses = []; - final List<_VideoRequest> requests = []; - - void enqueueJson(int statusCode, Map body) { - _responses.add(_QueuedResponse(statusCode, jsonEncode(body))); - } - - @override - Future send(http.BaseRequest request) async { - final bodyBytes = await request.finalize().toBytes(); - requests.add( - _VideoRequest( - method: request.method, - url: request.url, - headers: Map.from(request.headers), - bodyBytes: bodyBytes, - ), - ); - if (_responses.isEmpty) { - throw StateError('No queued video response for ${request.method}'); - } - final response = _responses.removeAt(0); - return http.StreamedResponse( - Stream>.value(utf8.encode(response.body)), - response.statusCode, - headers: const {'content-type': 'application/json'}, - request: request, - ); - } -} - -class _QueuedResponse { - const _QueuedResponse(this.statusCode, this.body); - - final int statusCode; - final String body; -} - -class _VideoRequest { - const _VideoRequest({ - required this.method, - required this.url, - required this.headers, - required this.bodyBytes, - }); - - final String method; - final Uri url; - final Map headers; - final Uint8List bodyBytes; -} - -class _FakeFile implements File { - _FakeFile( - this.path, { - bool exists = true, - this.chunks = const >[], - this.lengthOverride, - }) : _shouldExist = exists; - - @override - final String path; - final bool _shouldExist; - final List> chunks; - final int? lengthOverride; - - @override - bool existsSync() => _shouldExist; - - @override - Future length() async => - lengthOverride ?? - chunks.fold(0, (total, chunk) => total + chunk.length); - - @override - Stream> openRead([int? start, int? end]) => - Stream>.fromIterable(chunks); - - @override - dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); -} - -class _FakeAuthRepository implements AuthRepository { - _FakeAuthRepository({required this.authenticated, required this.atproto}); - - final bool authenticated; - int refreshTokenCalls = 0; - - @override - final PoptartClient? atproto; - - @override - bool get isAuthenticated => authenticated; - - @override - Future refreshToken() async { - refreshTokenCalls++; - return true; - } - - @override - dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); -} - -class _FakeSprkRepository implements SprkRepository { - _FakeSprkRepository({required this.auth, required this.repo}); - - final _FakeAuthRepository auth; - - @override - final RepoRepository repo; - - @override - AuthRepository get authRepository => auth; - - @override - String get sprkDid => _Harness.sprkDid; - - @override - String get bskyDid => _Harness.bskyDid; - - @override - Future executeWithRetry(Future Function() apiCall) => apiCall(); - - @override - dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); -} - -class _FakeRepoRepository implements RepoRepository { - final List<_CreateCall> createCalls = []; - final List<_DeleteCall> deleteCalls = []; - Object? createError; - - @override - Future createRecord({ - required String collection, - required Map record, - String? rkey, - String? repo, - }) async { - createCalls.add(_CreateCall(collection: collection, record: record)); - if (createError case final Object error) { - throw error; - } - return RepoStrongRef( - uri: AtUri('at://did:plc:viewer/$collection/result'), - cid: 'result-cid', - ); - } - - @override - Future deleteRecord({ - required AtUri uri, - bool skipBskyCrosspostCleanup = false, - }) async { - deleteCalls.add( - _DeleteCall(uri: uri, skipBskyCrosspostCleanup: skipBskyCrosspostCleanup), - ); - } - - @override - dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); -} - -class _CreateCall { - const _CreateCall({required this.collection, required this.record}); - - final String collection; - final Map record; -} - -class _DeleteCall { - const _DeleteCall({ - required this.uri, - required this.skipBskyCrosspostCleanup, - }); - - final AtUri uri; - final bool skipBskyCrosspostCleanup; -} diff --git a/test/src/core/network/atproto/data/repositories/feed_repository_records_test.dart b/test/src/core/network/atproto/data/repositories/feed_repository_records_test.dart new file mode 100644 index 00000000..168c723a --- /dev/null +++ b/test/src/core/network/atproto/data/repositories/feed_repository_records_test.dart @@ -0,0 +1,218 @@ +import 'package:flutter_test/flutter_test.dart'; +import 'package:poptart/poptart.dart'; +import 'package:spark/src/core/network/atproto/data/models/feed_models.dart'; +import 'package:spark/src/core/network/atproto/data/repositories/feed_repository_impl.dart'; +import 'package:spark/src/core/utils/logging/logger.dart'; + +import 'repository_test_support.dart'; + +void main() { + final fixedNow = DateTime.utc(2026, 7, 22, 12, 34, 56); + + FeedRepositoryImpl repository(RepositoryHarness harness) { + return FeedRepositoryImpl( + harness.sprk, + logger: SparkLogger(), + now: () => fixedNow, + ); + } + + group('FeedRepositoryImpl records', () { + test( + 'likePost selects the record collection and uses the injected clock', + () async { + final harness = RepositoryHarness(); + final sparkPost = AtUri( + 'at://did:plc:author/so.sprk.feed.post/spark-post', + ); + final bskyPost = AtUri( + 'at://did:plc:author/app.bsky.feed.post/bsky-post', + ); + final feedRepository = repository(harness); + + await feedRepository.likePost('spark-cid', sparkPost); + await feedRepository.likePost('bsky-cid', bskyPost); + + expect(harness.repo.createCalls, hasLength(2)); + final sparkCall = harness.repo.createCalls[0]; + expect(sparkCall.collection, 'so.sprk.feed.like'); + expect(sparkCall.record[r'$type'], 'so.sprk.feed.like'); + _expectStrongRef( + sparkCall.record['subject'], + uri: sparkPost, + cid: 'spark-cid', + ); + expect(sparkCall.record['createdAt'], fixedNow.toIso8601String()); + + final bskyCall = harness.repo.createCalls[1]; + expect(bskyCall.collection, 'app.bsky.feed.like'); + expect(bskyCall.record[r'$type'], 'app.bsky.feed.like'); + _expectStrongRef( + bskyCall.record['subject'], + uri: bskyPost, + cid: 'bsky-cid', + ); + expect(bskyCall.record['createdAt'], fixedNow.toIso8601String()); + }, + ); + + test( + 'unlikePost deletes the interaction without crosspost cleanup', + () async { + final harness = RepositoryHarness(); + final likeUri = AtUri( + 'at://did:plc:viewer/so.sprk.feed.like/interaction', + ); + + await repository(harness).unlikePost(likeUri); + + expect(harness.repo.deleteCalls, hasLength(1)); + expect(harness.repo.deleteCalls.single.uri, likeUri); + expect( + harness.repo.deleteCalls.single.skipBskyCrosspostCleanup, + isTrue, + ); + }, + ); + + test('postComment writes Spark reply roots and parents', () async { + final harness = RepositoryHarness(); + final parentUri = AtUri('at://did:plc:author/so.sprk.feed.post/parent'); + final rootUri = AtUri('at://did:plc:author/so.sprk.feed.post/root'); + + await repository(harness).postComment( + 'A Spark reply', + 'parent-cid', + parentUri, + rootCid: 'root-cid', + rootUri: rootUri, + ); + + final call = harness.repo.createCalls.single; + expect(call.collection, 'so.sprk.feed.reply'); + expect(call.record[r'$type'], 'so.sprk.feed.reply'); + expect(call.record['text'], 'A Spark reply'); + expect(call.record, isNot(contains('facets'))); + expect(call.record['createdAt'], fixedNow.toIso8601String()); + final reply = call.record['reply'] as Map; + expect(reply[r'$type'], 'so.sprk.feed.reply#replyRef'); + _expectStrongRef(reply['root'], uri: rootUri, cid: 'root-cid'); + _expectStrongRef(reply['parent'], uri: parentUri, cid: 'parent-cid'); + }); + + test( + 'postComment writes Bluesky posts and defaults root to parent', + () async { + final harness = RepositoryHarness(); + final parentUri = AtUri( + 'at://did:web:sprk.so/app.bsky.feed.post/parent', + ); + + await repository( + harness, + ).postComment('A Bluesky reply', 'parent-cid', parentUri); + + final call = harness.repo.createCalls.single; + expect(call.collection, 'app.bsky.feed.post'); + expect(call.record[r'$type'], 'app.bsky.feed.post'); + expect(call.record['text'], 'A Bluesky reply'); + expect(call.record['createdAt'], fixedNow.toIso8601String()); + final reply = call.record['reply'] as Map; + expect(reply[r'$type'], 'app.bsky.feed.post#replyRef'); + _expectStrongRef(reply['root'], uri: parentUri, cid: 'parent-cid'); + _expectStrongRef(reply['parent'], uri: parentUri, cid: 'parent-cid'); + }, + ); + + test('record creation errors propagate through the repository', () { + final harness = RepositoryHarness() + ..repo.createError = StateError('record write failed'); + + expect( + repository( + harness, + ).likePost('cid', AtUri('at://did:plc:author/so.sprk.feed.post/post')), + throwsA( + isA().having( + (error) => error.message, + 'message', + 'record write failed', + ), + ), + ); + }); + }); + + group('FeedRepositoryImpl threads', () { + test( + 'getThread sends Spark thread parameters and maps not-found posts', + () async { + final threadUri = AtUri( + 'at://did:plc:author/so.sprk.feed.post/missing', + ); + final harness = RepositoryHarness( + getResponse: { + 'thread': [ + { + r'$type': 'so.sprk.feed.getPostThread#threadItem', + 'uri': threadUri.toString(), + 'depth': 0, + 'value': { + r'$type': 'so.sprk.feed.defs#notFoundPost', + 'uri': threadUri.toString(), + 'notFound': true, + }, + }, + ], + }, + ); + + final result = await repository( + harness, + ).getThread(threadUri, depth: 4, parentHeight: 2); + + final request = harness.transport.singleRequest; + expect(request.uri.path, '/xrpc/so.sprk.feed.getPostThread'); + expect(request.uri.queryParameters['anchor'], threadUri.toString()); + expect(request.uri.queryParameters['depth'], '4'); + expect(request.uri.queryParameters['parentHeight'], '2'); + expect( + request.headers['atproto-proxy'], + FakeSprkRepository.testSprkDid, + ); + expect(result, isA()); + expect((result as NotFoundPost).uri, threadUri); + expect(result.notFound, isTrue); + }, + ); + + test('getThread rejects an uninitialized AtProto client', () { + final harness = RepositoryHarness(atprotoInitialized: false); + + expect( + repository( + harness, + ).getThread(AtUri('at://did:plc:author/so.sprk.feed.post/post')), + throwsA( + isA().having( + (error) => error.toString(), + 'message', + contains('AtProto not initialized'), + ), + ), + ); + expect(harness.transport.requests, isEmpty); + }); + }); +} + +void _expectStrongRef( + Object? value, { + required AtUri uri, + required String cid, +}) { + final ref = value as Map; + expect(ref[r'$type'], 'com.atproto.repo.strongRef'); + expect(ref['uri'], uri.toString()); + expect(ref['cid'], cid); +} diff --git a/test/src/core/network/atproto/data/repositories/feed_repository_requests_test.dart b/test/src/core/network/atproto/data/repositories/feed_repository_requests_test.dart new file mode 100644 index 00000000..3ac114c3 --- /dev/null +++ b/test/src/core/network/atproto/data/repositories/feed_repository_requests_test.dart @@ -0,0 +1,202 @@ +import 'package:flutter_test/flutter_test.dart'; +import 'package:poptart/poptart.dart'; +import 'package:spark/src/core/network/atproto/data/models/feed_models.dart'; +import 'package:spark/src/core/network/atproto/data/models/pref_models.dart'; +import 'package:spark/src/core/network/atproto/data/repositories/feed_repository_impl.dart'; +import 'package:spark/src/core/utils/logging/logger.dart'; +import 'package:sprk_poptart/so/sprk/actor/defs/profile_view_basic.dart'; + +import 'repository_test_support.dart'; + +void main() { + final fixedNow = DateTime.utc(2026, 7, 22, 12, 34, 56); + + FeedRepositoryImpl repository(RepositoryHarness harness) { + return FeedRepositoryImpl( + harness.sprk, + logger: SparkLogger(), + now: () => fixedNow, + ); + } + + group('FeedRepositoryImpl feed requests', () { + test('getFeed routes timeline requests and maps hydrated posts', () async { + final post = _postView(); + final harness = RepositoryHarness( + getResponse: { + 'cursor': 'timeline-next', + 'feed': [ + {'post': post.toJson()}, + ], + }, + ); + final feed = Feed( + type: 'timeline', + config: makeSavedFeed( + id: 'timeline', + type: 'timeline', + value: 'timeline', + pinned: true, + ), + ); + + final result = await repository(harness).getFeed( + feed, + limit: 12, + cursor: 'timeline-cursor', + labelerDids: const ['did:plc:one', 'did:plc:two'], + ); + + final request = harness.transport.singleRequest; + expect(request.uri.path, '/xrpc/so.sprk.feed.getTimeline'); + expect(request.uri.queryParameters['limit'], '12'); + expect(request.uri.queryParameters['cursor'], 'timeline-cursor'); + expect(request.headers['atproto-proxy'], FakeSprkRepository.testSprkDid); + expect( + request.headers['atproto-accept-labelers'], + 'did:plc:one,did:plc:two', + ); + expect(result.cursor, 'timeline-next'); + expect(result.feed.single.post.uri, post.uri); + }); + + test( + 'getFeedView sends Spark feed requests to the Spark service', + () async { + final post = _postView(); + final harness = RepositoryHarness( + getResponse: { + 'cursor': 'spark-next', + 'feed': [ + {'post': post.toJson()}, + ], + }, + ); + final feedUri = AtUri( + 'at://did:plc:generator/so.sprk.feed.generator/spark-feed', + ); + + final result = await repository(harness).getFeedView( + feedUri, + limit: 8, + cursor: 'spark-cursor', + labelerDids: const ['did:plc:labeler'], + ); + + final request = harness.transport.singleRequest; + expect(request.uri.path, '/xrpc/so.sprk.feed.getFeed'); + expect(request.uri.queryParameters['feed'], feedUri.toString()); + expect(request.uri.queryParameters['limit'], '8'); + expect(request.uri.queryParameters['cursor'], 'spark-cursor'); + expect( + request.headers['atproto-proxy'], + FakeSprkRepository.testSprkDid, + ); + expect(request.headers['atproto-accept-labelers'], 'did:plc:labeler'); + expect(result.cursor, 'spark-next'); + expect(result.feed.single.post.uri, post.uri); + }, + ); + + test( + 'getFeedView selects Bluesky and omits Spark labeler headers', + () async { + final harness = RepositoryHarness( + getResponse: const { + 'cursor': 'bsky-next', + 'feed': [], + }, + ); + final feedUri = AtUri( + 'at://did:plc:generator/app.bsky.feed.generator/bsky-feed', + ); + + final result = await repository(harness).getFeedView( + feedUri, + limit: 6, + cursor: 'bsky-cursor', + labelerDids: const ['did:plc:labeler'], + ); + + final request = harness.transport.singleRequest; + expect(request.uri.path, '/xrpc/app.bsky.feed.getFeed'); + expect(request.uri.queryParameters['feed'], feedUri.toString()); + expect(request.uri.queryParameters['limit'], '6'); + expect(request.uri.queryParameters['cursor'], 'bsky-cursor'); + expect( + request.headers['atproto-proxy'], + FakeSprkRepository.testBskyDid, + ); + expect(request.headers, isNot(contains('atproto-accept-labelers'))); + expect(result.cursor, 'bsky-next'); + expect(result.feed, isEmpty); + }, + ); + + test('getFeedView rejects unauthenticated requests before transport', () { + final harness = RepositoryHarness(authenticated: false); + + expect( + repository(harness).getFeedView( + AtUri('at://did:plc:generator/so.sprk.feed.generator/feed'), + ), + throwsA( + isA().having( + (error) => error.toString(), + 'message', + contains('Not authenticated'), + ), + ), + ); + expect(harness.transport.requests, isEmpty); + }); + + test('getFeedView propagates the transport 500 response', () async { + final harness = RepositoryHarness( + getStatusCode: 500, + getResponse: const { + 'error': 'InternalServerError', + 'message': 'feed unavailable', + }, + ); + + await expectLater( + repository(harness).getFeedView( + AtUri('at://did:plc:generator/so.sprk.feed.generator/feed'), + ), + throwsA( + isA().having( + (error) => error.toString(), + 'message', + allOf( + contains('InternalServerError'), + contains('feed unavailable'), + ), + ), + ), + ); + expect(harness.transport.requests, hasLength(1)); + expect( + harness.transport.singleRequest.uri.path, + '/xrpc/so.sprk.feed.getFeed', + ); + }); + }); +} + +PostView _postView() { + return PostView( + uri: AtUri('at://did:plc:author/so.sprk.feed.post/post'), + cid: 'post-cid', + author: const ProfileViewBasic( + did: 'did:plc:author', + handle: 'author.test', + ), + record: const { + r'$type': 'so.sprk.feed.post', + 'caption': {'text': 'A post', 'facets': []}, + 'createdAt': '2026-07-22T12:00:00.000Z', + }, + indexedAt: DateTime.utc(2026, 7, 22, 12), + ); +} diff --git a/test/src/core/network/atproto/data/repositories/repository_test_support.dart b/test/src/core/network/atproto/data/repositories/repository_test_support.dart index cd8a8c47..b3541e3b 100644 --- a/test/src/core/network/atproto/data/repositories/repository_test_support.dart +++ b/test/src/core/network/atproto/data/repositories/repository_test_support.dart @@ -11,24 +11,42 @@ class RepositoryHarness { RepositoryHarness({ bool authenticated = true, bool atprotoInitialized = true, + bool oauth = false, String? did = 'did:plc:viewer', Map? getResponse, + int getStatusCode = 200, }) : transport = TestTransport() { + final atproto = oauth + ? PoptartClient.fromOAuthSession( + restoreOAuthSession( + accessToken: 'opaque-access-token', + refreshToken: 'opaque-refresh-token', + scope: 'atproto', + expiresAt: DateTime.utc(2030), + sub: 'did:plc:viewer', + clientId: 'https://spark.test/client-metadata.json', + pdsEndpoint: 'pds.test', + publicKey: 'unused-public-key', + privateKey: 'unused-private-key', + ), + service: 'pds.test', + getClient: transport.get, + postClient: transport.post, + ) + : PoptartClient.anonymous( + service: 'pds.test', + getClient: transport.get, + postClient: transport.post, + ); auth = FakeAuthRepository( authenticated: authenticated, did: did, - atproto: atprotoInitialized - ? PoptartClient.anonymous( - service: 'pds.test', - getClient: transport.get, - postClient: transport.post, - ) - : null, + atproto: atprotoInitialized ? atproto : null, ); repo = FakeRepoRepository(); sprk = FakeSprkRepository(auth: auth, repo: repo); if (atprotoInitialized && getResponse != null) { - transport.enqueueGet(getResponse); + transport.enqueueGet(getResponse, statusCode: getStatusCode); } } @@ -139,6 +157,8 @@ class FakeAuthRepository implements AuthRepository { }); final bool authenticated; + int refreshTokenCalls = 0; + bool refreshTokenResult = true; @override final String? did; @@ -153,7 +173,10 @@ class FakeAuthRepository implements AuthRepository { Future get initializationComplete => Future.value(); @override - Future refreshToken() async => true; + Future refreshToken() async { + refreshTokenCalls++; + return refreshTokenResult; + } @override dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); @@ -162,6 +185,9 @@ class FakeAuthRepository implements AuthRepository { class FakeSprkRepository implements SprkRepository { FakeSprkRepository({required this.auth, required this.repo}); + static const testSprkDid = 'did:web:sprk.test#sprk_appview'; + static const testBskyDid = 'did:web:bsky.test#bsky_appview'; + final FakeAuthRepository auth; @override @@ -171,10 +197,10 @@ class FakeSprkRepository implements SprkRepository { AuthRepository get authRepository => auth; @override - String get sprkDid => 'did:web:sprk.test#sprk_appview'; + String get sprkDid => testSprkDid; @override - String get bskyDid => 'did:web:bsky.test#bsky_appview'; + String get bskyDid => testBskyDid; @override String get modDid => 'did:web:mod.sprk.test'; @@ -192,6 +218,7 @@ class FakeSprkRepository implements SprkRepository { class FakeRepoRepository implements RepoRepository { final List createCalls = []; final List deleteCalls = []; + Object? createError; @override Future createRecord({ @@ -208,6 +235,9 @@ class FakeRepoRepository implements RepoRepository { repo: repo, ), ); + if (createError case final Object error) { + throw error; + } return RepoStrongRef( uri: AtUri('at://did:plc:viewer/$collection/result'), cid: 'result-cid', diff --git a/test/src/core/network/atproto/data/repositories/story_repository_impl_test.dart b/test/src/core/network/atproto/data/repositories/story_repository_impl_test.dart index 6760cd13..63c7553d 100644 --- a/test/src/core/network/atproto/data/repositories/story_repository_impl_test.dart +++ b/test/src/core/network/atproto/data/repositories/story_repository_impl_test.dart @@ -43,6 +43,37 @@ void main() { ]); }); + test('listStoryRecords owns record paging parameters', () async { + final harness = RepositoryHarness( + getResponse: const {'records': [], 'cursor': 'next-page'}, + ); + final repository = StoryRepositoryImpl(harness.sprk); + + final result = await repository.listStoryRecords( + did: 'did:plc:viewer', + cursor: 'current-page', + ); + + expect(result.records, isEmpty); + expect(result.cursor, 'next-page'); + final request = harness.transport.singleRequest; + expect(request.uri.path, '/xrpc/com.atproto.repo.listRecords'); + expect(request.uri.queryParameters['repo'], 'did:plc:viewer'); + expect(request.uri.queryParameters['collection'], 'so.sprk.story.post'); + expect(request.uri.queryParameters['cursor'], 'current-page'); + expect(request.uri.queryParameters['limit'], '100'); + }); + + test('deleteStoryRecord delegates to the record repository', () async { + final harness = RepositoryHarness(); + final repository = StoryRepositoryImpl(harness.sprk); + final uri = AtUri('at://did:plc:viewer/so.sprk.story.post/story'); + + await repository.deleteStoryRecord(uri); + + expect(harness.repo.deleteCalls.single.uri, uri); + }); + test( 'postStory normalizes empty optionals and uses injected time', () async { diff --git a/test/src/core/network/atproto/data/services/video_upload_service_test.dart b/test/src/core/network/atproto/data/services/video_upload_service_test.dart new file mode 100644 index 00000000..145d4d9a --- /dev/null +++ b/test/src/core/network/atproto/data/services/video_upload_service_test.dart @@ -0,0 +1,502 @@ +import 'dart:convert'; +import 'dart:io'; +import 'dart:typed_data'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:http/http.dart' as http; +import 'package:poptart/poptart.dart'; +import 'package:spark/src/core/network/atproto/data/repositories/feed_repository_impl.dart'; +import 'package:spark/src/core/network/atproto/data/services/video_upload_service.dart'; +import 'package:spark/src/core/utils/logging/logger.dart'; +import 'package:spark/src/core/utils/video_upload_exception.dart'; + +import '../repositories/repository_test_support.dart'; + +void main() { + final fixedNow = DateTime.utc(2026, 7, 22, 12, 34, 56); + + FeedRepositoryImpl repository({ + required RepositoryHarness harness, + required _VideoClient videoClient, + required File Function(String) file, + Future Function(Duration)? processingDelay, + Future Function(PoptartClient)? serviceAuthTokenRequest, + }) { + return FeedRepositoryImpl( + harness.sprk, + logger: SparkLogger(), + now: () => fixedNow, + videoUploadService: VideoUploadClient( + harness.auth, + logger: SparkLogger(), + now: () => fixedNow, + httpClientFactory: () => videoClient, + file: file, + processingDelay: processingDelay, + serviceAuthTokenRequest: serviceAuthTokenRequest, + ), + ); + } + + group('VideoUploadClient', () { + test('rejects missing, empty, and oversized files before HTTP', () async { + final videoClient = _VideoClient(); + final files = { + '/missing.mp4': _FakeFile('/missing.mp4', exists: false), + '/empty.mp4': _FakeFile('/empty.mp4'), + '/huge.mp4': _FakeFile( + '/huge.mp4', + lengthOverride: 2 * 1024 * 1024 * 1024 + 1, + ), + }; + final harness = RepositoryHarness(oauth: true); + final feedRepository = repository( + harness: harness, + videoClient: videoClient, + file: (path) => files[path]!, + serviceAuthTokenRequest: (_) async => 'service-token', + ); + + await expectLater( + feedRepository.uploadVideo('/missing.mp4'), + throwsA( + isA().having( + (error) => error.toString(), + 'message', + contains('Video file not found'), + ), + ), + ); + await expectLater( + feedRepository.uploadVideo('/empty.mp4'), + throwsA( + isA().having( + (error) => error.toString(), + 'message', + contains('Video file is empty'), + ), + ), + ); + await expectLater( + feedRepository.uploadVideo('/huge.mp4'), + throwsA( + isA() + .having((error) => error.statusCode, 'statusCode', 413) + .having( + (error) => error.uploadSizeBytes, + 'uploadSizeBytes', + 2 * 1024 * 1024 * 1024 + 1, + ) + .having( + (error) => error.limitBytes, + 'limitBytes', + 2 * 1024 * 1024 * 1024, + ), + ), + ); + expect(videoClient.requests, isEmpty); + expect(videoClient.closed, isFalse); + }); + + test( + 'streams bytes, reports progress, parses a blob, and closes HTTP', + () async { + final videoClient = _VideoClient() + ..enqueueJson(200, { + 'blobRef': _blobJson(mimeType: 'video/quicktime', size: 6), + }); + final progress = []; + final harness = RepositoryHarness(oauth: true); + final feedRepository = repository( + harness: harness, + videoClient: videoClient, + file: (_) => _FakeFile( + '/clip.mov', + chunks: const [ + [1, 2], + [3, 4, 5, 6], + ], + ), + serviceAuthTokenRequest: (_) async => 'service-token', + ); + + final result = await feedRepository.uploadVideo( + 'file:///clip.mov', + onUploadProgress: progress.add, + ); + + final request = videoClient.requests.single; + expect(request.method, 'POST'); + expect(request.url.path, '/xrpc/so.sprk.video.uploadVideo'); + expect(request.headers['Authorization'], 'Bearer service-token'); + expect(request.headers['Content-Type'], 'video/quicktime'); + expect(request.bodyBytes, [1, 2, 3, 4, 5, 6]); + expect(progress, [0, closeTo(1 / 3, 0.0001), 1, 1]); + expect(result.videoBlob.mimeType, 'video/quicktime'); + expect(result.videoBlob.size, 6); + expect(result.audioBlob, isNull); + expect(result.audioDetails, isNull); + expect(videoClient.closed, isTrue); + }, + ); + + test('maps upload rejection details and closes HTTP', () async { + final videoClient = _VideoClient() + ..enqueueJson(413, { + 'error': 'PayloadTooLarge', + 'message': 'limit is 100 MB', + }); + final harness = RepositoryHarness(oauth: true); + final feedRepository = repository( + harness: harness, + videoClient: videoClient, + file: (_) => _FakeFile( + '/clip.mp4', + chunks: const [ + [1, 2, 3], + ], + ), + serviceAuthTokenRequest: (_) async => 'service-token', + ); + + await expectLater( + feedRepository.uploadVideo('/clip.mp4'), + throwsA( + isA() + .having((error) => error.statusCode, 'statusCode', 413) + .having((error) => error.uploadSizeBytes, 'uploadSizeBytes', 3) + .having((error) => error.isPayloadTooLarge, 'too large', isTrue) + .having( + (error) => error.message, + 'message', + contains('limit is 100 MB'), + ), + ), + ); + expect(videoClient.requests, hasLength(1)); + expect(videoClient.closed, isTrue); + }); + + test( + 'polls deterministically and parses video plus audio output', + () async { + final videoClient = _VideoClient() + ..enqueueJson(200, { + 'jobStatus': {'jobId': 'job-1', 'state': 'JOB_STATE_QUEUED'}, + }) + ..enqueueJson(200, { + 'jobStatus': {'jobId': 'job-1', 'state': 'JOB_STATE_PROCESSING'}, + }) + ..enqueueJson(200, { + 'jobStatus': { + 'jobId': 'job-1', + 'state': 'JOB_STATE_COMPLETED', + 'blob': _blobJson(mimeType: 'video/mp4', size: 3), + 'audio': { + 'blob': _blobJson(mimeType: 'audio/aac', size: 2), + 'details': { + r'$type': 'so.sprk.sound.defs#audioDetails', + 'artist': 'Artist', + 'title': 'Title', + }, + }, + }, + }); + final delays = []; + final harness = RepositoryHarness(oauth: true); + final feedRepository = repository( + harness: harness, + videoClient: videoClient, + file: (_) => _FakeFile( + '/clip.mp4', + chunks: const [ + [1, 2, 3], + ], + ), + processingDelay: (duration) async => delays.add(duration), + serviceAuthTokenRequest: (_) async => 'service-token', + ); + + final result = await feedRepository.uploadVideo('/clip.mp4'); + + expect(delays, const [Duration(seconds: 2), Duration(seconds: 2)]); + expect(videoClient.requests, hasLength(3)); + for (final request in videoClient.requests.skip(1)) { + expect(request.method, 'GET'); + expect(request.url.path, '/xrpc/so.sprk.video.getJobStatus'); + expect(request.url.queryParameters['jobId'], 'job-1'); + expect(request.headers['Authorization'], 'Bearer service-token'); + } + expect(result.videoBlob.mimeType, 'video/mp4'); + expect(result.audioBlob?.mimeType, 'audio/aac'); + expect(result.audioDetails?.artist, 'Artist'); + expect(result.audioDetails?.title, 'Title'); + expect(videoClient.closed, isTrue); + }, + ); + + test( + 'refreshes an expired polling token and retries the request', + () async { + final videoClient = _VideoClient() + ..enqueueJson(200, { + 'jobStatus': {'jobId': 'job-1', 'state': 'JOB_STATE_PROCESSING'}, + }) + ..enqueueJson(401, {'message': 'JWT has expired'}) + ..enqueueJson(200, { + 'jobStatus': { + 'jobId': 'job-1', + 'state': 'JOB_STATE_COMPLETED', + 'blob': _blobJson(mimeType: 'video/mp4', size: 3), + }, + }); + var tokenRequests = 0; + final harness = RepositoryHarness(oauth: true); + final feedRepository = repository( + harness: harness, + videoClient: videoClient, + file: (_) => _FakeFile( + '/clip.mp4', + chunks: const [ + [1, 2, 3], + ], + ), + processingDelay: (_) async {}, + serviceAuthTokenRequest: (_) async { + tokenRequests++; + if (tokenRequests == 2) { + throw StateError('PDS token expired'); + } + return tokenRequests == 1 ? 'old-token' : 'new-token'; + }, + ); + + await feedRepository.uploadVideo('/clip.mp4'); + + expect(tokenRequests, 3); + expect(harness.auth.refreshTokenCalls, 1); + expect(videoClient.requests, hasLength(3)); + expect( + videoClient.requests[1].headers['Authorization'], + 'Bearer old-token', + ); + expect( + videoClient.requests[2].headers['Authorization'], + 'Bearer new-token', + ); + expect(videoClient.closed, isTrue); + }, + ); + + test('fails after three consecutive polling errors', () async { + final videoClient = _VideoClient() + ..enqueueJson(200, { + 'jobStatus': {'jobId': 'job-1', 'state': 'JOB_STATE_PROCESSING'}, + }) + ..enqueueJson(500, {'message': 'first'}) + ..enqueueJson(502, {'message': 'second'}) + ..enqueueJson(503, {'message': 'third'}); + var delays = 0; + final harness = RepositoryHarness(oauth: true); + final feedRepository = repository( + harness: harness, + videoClient: videoClient, + file: (_) => _FakeFile( + '/clip.mp4', + chunks: const [ + [1], + ], + ), + processingDelay: (_) async => delays++, + serviceAuthTokenRequest: (_) async => 'service-token', + ); + + await expectLater( + feedRepository.uploadVideo('/clip.mp4'), + throwsA( + isA().having( + (error) => error.message, + 'message', + contains('third'), + ), + ), + ); + expect(delays, 3); + expect(videoClient.requests, hasLength(4)); + expect(videoClient.closed, isTrue); + }); + + test( + 'fails failed jobs and times out processing without wall-clock waits', + () async { + final failedClient = _VideoClient() + ..enqueueJson(200, { + 'jobStatus': { + 'jobId': 'failed-job', + 'state': 'JOB_STATE_FAILED', + 'error': 'TranscodeFailed', + 'message': 'unsupported codec', + }, + }); + final failedHarness = RepositoryHarness(oauth: true); + final failedRepository = repository( + harness: failedHarness, + videoClient: failedClient, + file: (_) => _FakeFile( + '/clip.mp4', + chunks: const [ + [1], + ], + ), + serviceAuthTokenRequest: (_) async => 'service-token', + ); + await expectLater( + failedRepository.uploadVideo('/clip.mp4'), + throwsA( + isA().having( + (error) => error.message, + 'message', + contains('unsupported codec'), + ), + ), + ); + expect(failedClient.closed, isTrue); + + final timeoutClient = _VideoClient() + ..enqueueJson(200, { + 'jobStatus': {'jobId': 'slow-job', 'state': 'JOB_STATE_PROCESSING'}, + }); + for (var i = 0; i < 120; i++) { + timeoutClient.enqueueJson(200, { + 'jobStatus': {'jobId': 'slow-job', 'state': 'JOB_STATE_PROCESSING'}, + }); + } + var delays = 0; + final timeoutHarness = RepositoryHarness(oauth: true); + final timeoutRepository = repository( + harness: timeoutHarness, + videoClient: timeoutClient, + file: (_) => _FakeFile( + '/clip.mp4', + chunks: const [ + [1], + ], + ), + processingDelay: (_) async => delays++, + serviceAuthTokenRequest: (_) async => 'service-token', + ); + await expectLater( + timeoutRepository.uploadVideo('/clip.mp4'), + throwsA( + isA().having( + (error) => error.message, + 'message', + contains('timed out'), + ), + ), + ); + expect(delays, 121); + expect(timeoutClient.requests, hasLength(121)); + expect(timeoutClient.closed, isTrue); + }, + ); + }); +} + +Map _blobJson({required String mimeType, required int size}) => + { + r'$type': 'blob', + 'mimeType': mimeType, + 'size': size, + 'ref': {r'$link': 'bafkreigh2akiscaildc2'}, + }; + +class _VideoClient extends http.BaseClient { + final List<_QueuedResponse> _responses = []; + final List<_VideoRequest> requests = []; + bool closed = false; + + void enqueueJson(int statusCode, Map body) { + _responses.add(_QueuedResponse(statusCode, jsonEncode(body))); + } + + @override + Future send(http.BaseRequest request) async { + final bodyBytes = await request.finalize().toBytes(); + requests.add( + _VideoRequest( + method: request.method, + url: request.url, + headers: Map.from(request.headers), + bodyBytes: bodyBytes, + ), + ); + if (_responses.isEmpty) { + throw StateError('No queued video response for ${request.method}'); + } + final response = _responses.removeAt(0); + return http.StreamedResponse( + Stream>.value(utf8.encode(response.body)), + response.statusCode, + headers: const {'content-type': 'application/json'}, + request: request, + ); + } + + @override + void close() { + closed = true; + super.close(); + } +} + +class _QueuedResponse { + const _QueuedResponse(this.statusCode, this.body); + + final int statusCode; + final String body; +} + +class _VideoRequest { + const _VideoRequest({ + required this.method, + required this.url, + required this.headers, + required this.bodyBytes, + }); + + final String method; + final Uri url; + final Map headers; + final Uint8List bodyBytes; +} + +class _FakeFile implements File { + _FakeFile( + this.path, { + bool exists = true, + this.chunks = const >[], + this.lengthOverride, + }) : _shouldExist = exists; + + @override + final String path; + final bool _shouldExist; + final List> chunks; + final int? lengthOverride; + + @override + bool existsSync() => _shouldExist; + + @override + Future length() async => + lengthOverride ?? + chunks.fold(0, (total, chunk) => total + chunk.length); + + @override + Stream> openRead([int? start, int? end]) => + Stream>.fromIterable(chunks); + + @override + dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); +} diff --git a/test/src/core/pro_video_editor/providers/sound_picker_search_provider_test.dart b/test/src/core/pro_video_editor/providers/sound_picker_search_provider_test.dart index fe8fccc7..29e85d39 100644 --- a/test/src/core/pro_video_editor/providers/sound_picker_search_provider_test.dart +++ b/test/src/core/pro_video_editor/providers/sound_picker_search_provider_test.dart @@ -9,6 +9,7 @@ import 'package:poptart/poptart.dart'; import 'package:poptart_lex/com/atproto/repo/strong_ref.dart'; import 'package:spark/src/core/network/atproto/data/models/models.dart'; import 'package:spark/src/core/network/atproto/data/repositories/sound_repository.dart'; +import 'package:spark/src/core/providers/debounce_scheduler.dart'; import 'package:spark/src/core/pro_video_editor/providers/sound_picker_search_provider.dart'; import 'package:spark/src/core/pro_video_editor/providers/sound_picker_search_state.dart'; import 'package:spark/src/core/utils/logging/log_service.dart'; @@ -44,9 +45,7 @@ void main() { ProviderContainer createContainer() { final container = ProviderContainer( overrides: [ - soundPickerSearchDebounceSchedulerProvider.overrideWithValue( - debounceScheduler.schedule, - ), + debounceSchedulerProvider.overrideWithValue(debounceScheduler.schedule), ], ); addTearDown(container.dispose); diff --git a/test/src/features/posting/providers/recording_provider_test.dart b/test/src/features/posting/providers/recording_provider_test.dart index a131b7a2..44ed2a09 100644 --- a/test/src/features/posting/providers/recording_provider_test.dart +++ b/test/src/features/posting/providers/recording_provider_test.dart @@ -25,19 +25,31 @@ void main() { recordingTickSchedulerProvider.overrideWithValue(scheduler.schedule), ], ); - addTearDown(container.dispose); final subscription = container.listen( recordingProvider, (previous, next) {}, ); - addTearDown(subscription.close); + final tempDir = await Directory.systemTemp.createTemp( + 'recording-provider-timing-test', + ); final notifier = container.read(recordingProvider.notifier); + final cleanupComplete = notifier.disposalCleanupComplete; + addTearDown(() async { + subscription.close(); + container.dispose(); + await cleanupComplete; + if (await tempDir.exists()) { + await tempDir.delete(recursive: true); + } + }); + final firstSegmentPath = '${tempDir.path}/segment-1.mp4'; + final secondSegmentPath = '${tempDir.path}/segment-2.mp4'; notifier.startRecording(); scheduler.tick(3); notifier.stopRecording(); - notifier.addSegment(XFile('/tmp/segment-1.mp4')); + notifier.addSegment(XFile(firstSegmentPath)); final pausedState = container.read(recordingProvider); expect(pausedState.isRecording, isFalse); @@ -47,15 +59,12 @@ void main() { notifier.startRecording(); scheduler.tick(2); notifier.stopRecording(); - notifier.addSegment(XFile('/tmp/segment-2.mp4')); + notifier.addSegment(XFile(secondSegmentPath)); final resumedState = container.read(recordingProvider); expect(resumedState.isRecording, isFalse); expect(resumedState.elapsedDuration, const Duration(milliseconds: 500)); - expect(resumedState.segmentPaths, [ - '/tmp/segment-1.mp4', - '/tmp/segment-2.mp4', - ]); + expect(resumedState.segmentPaths, [firstSegmentPath, secondSegmentPath]); expect(resumedState.canFinalize, isTrue); }); diff --git a/test/src/features/search/providers/search_providers_test.dart b/test/src/features/search/providers/search_providers_test.dart index 70a332cf..468bccf2 100644 --- a/test/src/features/search/providers/search_providers_test.dart +++ b/test/src/features/search/providers/search_providers_test.dart @@ -15,13 +15,13 @@ import 'package:spark/src/core/network/atproto/data/repositories/actor_repositor import 'package:spark/src/core/network/atproto/data/repositories/feed_repository.dart'; import 'package:spark/src/core/network/atproto/data/repositories/graph_repository.dart'; import 'package:spark/src/core/network/atproto/data/repositories/sprk_repository.dart'; +import 'package:spark/src/core/providers/debounce_scheduler.dart'; import 'package:spark/src/core/utils/logging/log_level.dart'; import 'package:spark/src/core/utils/logging/log_output.dart'; import 'package:spark/src/core/utils/logging/log_service.dart'; import 'package:spark/src/core/utils/logging/logger.dart'; import 'package:spark/src/features/search/providers/actor_typeahead_provider.dart'; import 'package:spark/src/features/search/providers/post_search_provider.dart'; -import 'package:spark/src/features/search/providers/search_debounce_scheduler.dart'; import 'package:spark/src/features/search/providers/search_provider.dart'; import 'package:spark/src/features/search/providers/suggested_feeds_provider.dart'; import 'package:sprk_poptart/so/sprk/actor/defs.dart'; @@ -58,7 +58,7 @@ void main() { final result = ProviderContainer.test( retry: (retryCount, error) => null, overrides: [ - searchDebounceSchedulerProvider.overrideWithValue(scheduler.schedule), + debounceSchedulerProvider.overrideWithValue(scheduler.schedule), ...overrides, ], ); diff --git a/test/src/features/stories/providers/story_providers_test.dart b/test/src/features/stories/providers/story_providers_test.dart index e7c8581e..31eaf130 100644 --- a/test/src/features/stories/providers/story_providers_test.dart +++ b/test/src/features/stories/providers/story_providers_test.dart @@ -1,6 +1,7 @@ import 'dart:async'; import 'package:flutter_riverpod/flutter_riverpod.dart'; +import 'package:flutter_riverpod/misc.dart' show Override; import 'package:flutter_test/flutter_test.dart'; import 'package:poptart/poptart.dart'; import 'package:spark/src/core/network/atproto/atproto.dart'; @@ -396,13 +397,11 @@ void main() { ); final container = ProviderContainer( overrides: [ - storyProviderDependenciesProvider.overrideWithValue( - _dependencies( - did: 'did:plc:me', - atprotoAvailable: true, - repository: repository, - ), - ), + ..._dependencies( + did: 'did:plc:me', + atprotoAvailable: true, + repository: repository, + ).values, ], ); addTearDown(container.dispose); @@ -419,10 +418,10 @@ void main() { }); } -ProviderContainer _managerContainer(StoryProviderDependencies dependencies) { +ProviderContainer _managerContainer(_StoryOverrides dependencies) { final container = ProviderContainer( overrides: [ - storyProviderDependenciesProvider.overrideWithValue(dependencies), + ...dependencies.values, storyAutoDeleteExecutorProvider.overrideWith((ref) async {}), ], ); @@ -432,14 +431,14 @@ ProviderContainer _managerContainer(StoryProviderDependencies dependencies) { ProviderContainer _autoDeleteContainer({ required LocalStorageInterface storage, - required StoryProviderDependencies dependencies, + required _StoryOverrides dependencies, DateTime? now, Future Function()? refresh, }) { final container = ProviderContainer( overrides: [ storyAutoDeletePreferencesProvider.overrideWithValue(storage), - storyProviderDependenciesProvider.overrideWithValue(dependencies), + ...dependencies.values, storyClockProvider.overrideWithValue( () => now ?? DateTime.utc(2026, 7, 22, 12), ), @@ -450,25 +449,44 @@ ProviderContainer _autoDeleteContainer({ return container; } -StoryProviderDependencies _dependencies({ +_StoryOverrides _dependencies({ required String? did, required bool atprotoAvailable, required _FakeStoryRepository repository, - StoryRecordPageLoader? loadRecordPage, + Future Function({required String did, String? cursor})? + loadRecordPage, Future Function(AtUri uri)? deleteRecord, }) { - return StoryProviderDependencies( - readDid: () => did, - readAtprotoAvailable: () => atprotoAvailable, - loadRecordPage: - loadRecordPage ?? - ({required did, cursor}) async => const StoryRecordPage(records: []), - storyRepository: repository, - deleteRecord: deleteRecord ?? (_) async {}, - loggerFor: (name) => SparkLogger(name: name), + repository.recordPageLoader = + loadRecordPage ?? + ({required did, cursor}) async => const StoryRecordPage(records: []); + repository.deleteRecord = deleteRecord ?? (_) async {}; + return _StoryOverrides( + did: did, + atprotoAvailable: atprotoAvailable, + repository: repository, ); } +class _StoryOverrides { + const _StoryOverrides({ + required this.did, + required this.atprotoAvailable, + required this.repository, + }); + + final String? did; + final bool atprotoAvailable; + final StoryRepository repository; + + List get values => [ + storyCurrentDidProvider.overrideWithValue(did), + storyAtprotoAvailableProvider.overrideWithValue(atprotoAvailable), + storyRepositoryProvider.overrideWithValue(repository), + storyLoggerProvider.overrideWith((ref, name) => SparkLogger(name: name)), + ]; +} + StoryRecordEntry _record(String id, {required DateTime createdAt}) { return StoryRecordEntry( uri: AtUri('at://did:plc:me/so.sprk.story.post/$id'), @@ -487,6 +505,9 @@ StoryView _story(String id, {required int hour, ProfileViewBasic? author}) { } class _FakeStoryRepository implements StoryRepository { + Future Function({required String did, String? cursor})? + recordPageLoader; + Future Function(AtUri uri)? deleteRecord; List storyViews = []; final List> storyViewCalls = []; ({Map> storiesByAuthor, String? cursor}) @@ -496,6 +517,19 @@ class _FakeStoryRepository implements StoryRepository { ); final List<({int limit, String? cursor})> timelineCalls = []; + @override + Future listStoryRecords({ + required String did, + String? cursor, + }) { + return recordPageLoader!(did: did, cursor: cursor); + } + + @override + Future deleteStoryRecord(AtUri uri) { + return deleteRecord!(uri); + } + @override Future> getStoryViews(List storyUris) async { storyViewCalls.add(List.of(storyUris));