199 lines
5.9 KiB
Dart
199 lines
5.9 KiB
Dart
import 'dart:async';
|
|
|
|
import '../../core/errors/failure.dart';
|
|
import '../../core/utils/result.dart';
|
|
import '../models/content_provider_config.dart';
|
|
import '../models/post.dart';
|
|
import '../models/post_comment.dart';
|
|
import '../models/provider_health.dart';
|
|
import '../models/top_period_filter.dart';
|
|
import '../repositories/provider_repository.dart';
|
|
import 'content_provider.dart';
|
|
import 'provider_factory.dart';
|
|
|
|
class ProviderManager {
|
|
ProviderManager(this._repository, this._factory);
|
|
|
|
final ProviderRepository _repository;
|
|
final ProviderFactory _factory;
|
|
|
|
Future<Result<List<ContentProviderConfig>>> loadConfigs({
|
|
bool enabledOnly = true,
|
|
}) async {
|
|
await _repository.ensureSeedProviders();
|
|
return _repository.getProviders(enabledOnly: enabledOnly);
|
|
}
|
|
|
|
Future<Result<List<ContentProvider>>> activeProviders() async {
|
|
final result = await loadConfigs();
|
|
if (result is Error<List<ContentProviderConfig>>) {
|
|
return Error(result.failure);
|
|
}
|
|
final configs = (result as Success<List<ContentProviderConfig>>).data
|
|
..sort((a, b) => a.priority.compareTo(b.priority));
|
|
return Success(configs.map(_factory.create).toList());
|
|
}
|
|
|
|
Future<Result<void>> enableProvider(String id, bool enabled) async {
|
|
final result = await _repository.getProvider(id);
|
|
return result.fold(
|
|
onSuccess: (config) {
|
|
if (config == null) {
|
|
return const Error<void>(
|
|
Failure(code: 'not_found', message: 'Provider not found'),
|
|
);
|
|
}
|
|
return _repository.saveProvider(
|
|
config.copyWith(enabled: enabled, updatedAt: DateTime.now()),
|
|
);
|
|
},
|
|
onError: Error<void>.new,
|
|
);
|
|
}
|
|
|
|
Future<Result<void>> addCustomProvider(ContentProviderConfig config) {
|
|
return _repository.saveProvider(config.copyWith(updatedAt: DateTime.now()));
|
|
}
|
|
|
|
Future<Result<void>> deleteProvider(String id) =>
|
|
_repository.deleteProvider(id);
|
|
|
|
Future<Result<List<ProviderHealth>>> checkAll({int concurrency = 3}) async {
|
|
final providersResult = await activeProviders();
|
|
if (providersResult is Error<List<ContentProvider>>) {
|
|
return Error(providersResult.failure);
|
|
}
|
|
final providers = (providersResult as Success<List<ContentProvider>>).data;
|
|
if (providers.isEmpty) return const Success([]);
|
|
final results = await _runLimited(
|
|
providers,
|
|
concurrency,
|
|
(provider) async {
|
|
final health = await provider.checkHealth();
|
|
await _repository.saveHealth(health);
|
|
return health;
|
|
},
|
|
);
|
|
return Success(results);
|
|
}
|
|
|
|
Future<Result<List<Post>>> searchAcrossProviders({
|
|
required List<String> tags,
|
|
required int page,
|
|
int limit = 50,
|
|
String? rating,
|
|
String? providerId,
|
|
TopPeriodFilter topPeriod = TopPeriodFilter.none,
|
|
}) async {
|
|
final providersResult = await activeProviders();
|
|
if (providersResult is Error<List<ContentProvider>>) {
|
|
return Error(providersResult.failure);
|
|
}
|
|
var providers = (providersResult as Success<List<ContentProvider>>).data;
|
|
if (providerId != null) {
|
|
providers =
|
|
providers.where((provider) => provider.id == providerId).toList();
|
|
}
|
|
|
|
final posts = <Post>[];
|
|
for (final provider in providers) {
|
|
try {
|
|
final providerPosts = await provider.searchPosts(
|
|
tags: tags,
|
|
page: page,
|
|
limit: limit,
|
|
rating: rating,
|
|
topPeriod: topPeriod,
|
|
);
|
|
posts.addAll(providerPosts);
|
|
} catch (_) {
|
|
await _repository.saveHealth(
|
|
ProviderHealth(
|
|
providerId: provider.id,
|
|
status: ProviderStatus.offline,
|
|
pingMs: 0,
|
|
lastCheckedAt: DateTime.now(),
|
|
errorMessage: 'Search failed',
|
|
),
|
|
);
|
|
}
|
|
}
|
|
return Success(posts);
|
|
}
|
|
|
|
Future<Result<List<PostComment>>> getComments(
|
|
String providerId,
|
|
String postId,
|
|
) async {
|
|
final providersResult = await activeProviders();
|
|
if (providersResult is Error<List<ContentProvider>>) {
|
|
return Error(providersResult.failure);
|
|
}
|
|
final providers = (providersResult as Success<List<ContentProvider>>).data;
|
|
final matches = providers.where((provider) => provider.id == providerId);
|
|
if (matches.isEmpty) {
|
|
return const Error(
|
|
Failure(code: 'not_found', message: 'Provider not found'),
|
|
);
|
|
}
|
|
final provider = matches.first;
|
|
if (provider is! CommentProvider) return const Success([]);
|
|
try {
|
|
return Success(await (provider as CommentProvider).getComments(postId));
|
|
} catch (error) {
|
|
return Error(
|
|
Failure(
|
|
code: 'comments_unavailable',
|
|
message: 'Comments unavailable',
|
|
details: error,
|
|
),
|
|
);
|
|
}
|
|
}
|
|
|
|
Future<Result<Post?>> getPost(String providerId, String postId) async {
|
|
final providersResult = await activeProviders();
|
|
if (providersResult is Error<List<ContentProvider>>) {
|
|
return Error(providersResult.failure);
|
|
}
|
|
final providers = (providersResult as Success<List<ContentProvider>>).data;
|
|
final matches = providers.where((provider) => provider.id == providerId);
|
|
if (matches.isEmpty) {
|
|
return const Error(
|
|
Failure(code: 'not_found', message: 'Provider not found'));
|
|
}
|
|
try {
|
|
return Success(await matches.first.getPost(postId));
|
|
} catch (error) {
|
|
return Error(
|
|
Failure(
|
|
code: 'provider_get_post',
|
|
message: 'Could not load post',
|
|
details: error,
|
|
),
|
|
);
|
|
}
|
|
}
|
|
|
|
Future<List<R>> _runLimited<T, R>(
|
|
List<T> items,
|
|
int concurrency,
|
|
Future<R> Function(T item) action,
|
|
) async {
|
|
final results = <R>[];
|
|
var index = 0;
|
|
|
|
Future<void> worker() async {
|
|
while (index < items.length) {
|
|
final current = items[index++];
|
|
results.add(await action(current));
|
|
}
|
|
}
|
|
|
|
await Future.wait(
|
|
List.generate(concurrency.clamp(1, items.length), (_) => worker()),
|
|
);
|
|
return results;
|
|
}
|
|
}
|