import 'package:flutter/foundation.dart'; import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:uuid/uuid.dart'; import '../brick/cache_helpers.dart'; import '../models/programmer_task.model.dart'; import '../models/programmer_task_activity_log.model.dart'; import '../utils/snackbar.dart' show isOfflineSaveError; import 'connectivity_provider.dart'; import 'profile_provider.dart'; import 'realtime_controller.dart'; import 'stream_recovery.dart'; import 'supabase_provider.dart'; /// Lifecycle action types that determine whether a task is currently paused. const _pauseStateActions = ['started', 'paused', 'resumed', 'completed', 'cancelled']; int _statusRank(String status) { switch (status) { case ProgrammerTaskStatus.inProgress: return 0; case ProgrammerTaskStatus.queued: return 1; case ProgrammerTaskStatus.completed: return 2; case ProgrammerTaskStatus.cancelled: return 3; default: return 4; } } void _sortTasks(List tasks) { tasks.sort((a, b) { final r = _statusRank(a.status).compareTo(_statusRank(b.status)); if (r != 0) return r; return b.createdAt.compareTo(a.createdAt); }); } // --------------------------------------------------------------------------- // Offline-pending state (merged into the live stream; replayed on reconnect) // --------------------------------------------------------------------------- final offlinePendingProgrammerTasksProvider = StateProvider>((ref) => const []); final offlinePendingProgrammerTaskRawProvider = StateProvider>>((ref) => const []); final offlinePendingProgrammerTaskUpdatesProvider = StateProvider>>((ref) => const {}); final offlinePendingProgrammerTaskLogsProvider = StateProvider>>((ref) => const []); // --------------------------------------------------------------------------- // Stream providers // --------------------------------------------------------------------------- final programmerTasksProvider = StreamProvider>((ref) { final userId = ref.watch(currentUserIdProvider); if (userId == null) return const Stream.empty(); final client = ref.watch(supabaseClientProvider); final pendingNew = ref.watch(offlinePendingProgrammerTasksProvider); final pendingUpdates = ref.watch(offlinePendingProgrammerTaskUpdatesProvider); List applyPending(List rows) { final byId = {for (final r in rows) r.id: r}; pendingUpdates.forEach((id, fields) { final existing = byId[id]; if (existing == null) return; byId[id] = _applyUpdateFields(existing, fields); }); for (final p in pendingNew) { byId.putIfAbsent(p.id, () => p); } final result = byId.values.toList(); _sortTasks(result); return result; } final wrapper = StreamRecoveryWrapper( stream: client .from('programmer_tasks') .stream(primaryKey: ['id']) .order('created_at', ascending: false), onPollData: () async { final data = await client .from('programmer_tasks') .select() .order('created_at', ascending: false) .range(0, 199); return data.map(ProgrammerTask.fromMap).toList(); }, fromMap: ProgrammerTask.fromMap, channelName: 'programmer_tasks', onStatusChanged: ref.read(realtimeControllerProvider).handleChannelStatus, onOfflineData: () async { final all = await cachedListFromBrick(); return applyPending(all); }, onCacheMirror: (rows) => mirrorBatchToBrick(rows, tag: 'programmer_tasks'), ); ref.onDispose(wrapper.dispose); // Reconnect-replay: POST queued items directly, bypassing Brick's queue. ref.listen(isOnlineProvider, (wasOnline, isNowOnline) async { if (wasOnline == true || !isNowOnline) return; // 1. Replay offline-created tasks. final pendingRaw = List>.from( ref.read(offlinePendingProgrammerTaskRawProvider), ); if (pendingRaw.isNotEmpty) { final syncedIds = []; for (final payload in pendingRaw) { final id = payload['id'] as String; try { await client.from('programmer_tasks').insert(payload); try { await client.from('programmer_task_activity_logs').insert({ 'task_id': id, 'actor_id': payload['creator_id'], 'action_type': 'created', }); } catch (_) {} syncedIds.add(id); } catch (e) { final msg = e.toString(); if (msg.contains('23505') || msg.contains('duplicate key')) { syncedIds.add(id); } else { debugPrint('[programmerTasksProvider] failed to sync id=$id: $e'); } } } if (syncedIds.isNotEmpty) { ref.read(offlinePendingProgrammerTaskRawProvider.notifier).state = List.unmodifiable( ref .read(offlinePendingProgrammerTaskRawProvider) .where((p) => !syncedIds.contains(p['id'] as String)), ); ref.read(offlinePendingProgrammerTasksProvider.notifier).state = List.unmodifiable( ref .read(offlinePendingProgrammerTasksProvider) .where((t) => !syncedIds.contains(t.id)), ); } } // 2. Replay offline field/status updates. final pendingUpdatesMap = Map>.from( ref.read(offlinePendingProgrammerTaskUpdatesProvider), ); if (pendingUpdatesMap.isNotEmpty) { final syncedIds = []; for (final entry in pendingUpdatesMap.entries) { try { await client .from('programmer_tasks') .update(entry.value) .eq('id', entry.key); syncedIds.add(entry.key); } catch (e) { debugPrint( '[programmerTasksProvider] failed to sync update id=${entry.key}: $e', ); } } if (syncedIds.isNotEmpty) { final remaining = Map>.from( ref.read(offlinePendingProgrammerTaskUpdatesProvider), ); for (final id in syncedIds) { remaining.remove(id); } ref.read(offlinePendingProgrammerTaskUpdatesProvider.notifier).state = remaining; } } // 3. Replay queued activity logs (started/paused/resumed/etc.). final pendingLogs = List>.from( ref.read(offlinePendingProgrammerTaskLogsProvider), ); if (pendingLogs.isNotEmpty) { final synced = >[]; for (final log in pendingLogs) { try { await client.from('programmer_task_activity_logs').insert(log); synced.add(log); } catch (e) { debugPrint('[programmerTasksProvider] failed to sync log: $e'); } } if (synced.isNotEmpty) { ref.read(offlinePendingProgrammerTaskLogsProvider.notifier).state = List.unmodifiable( ref .read(offlinePendingProgrammerTaskLogsProvider) .where((l) => !synced.contains(l)), ); } } }); return wrapper.stream.map((result) => applyPending(result.data)); }); final programmerTaskByIdProvider = Provider.family(( ref, id, ) { final tasks = ref.watch(programmerTasksProvider).valueOrNull; if (tasks == null) return null; try { return tasks.firstWhere((t) => t.id == id); } catch (_) { return null; } }); /// The current user's actively-running task id (status in_progress AND latest /// activity log is not a pause), or null. Lets the list distinguish a *running* /// task from a *paused* one (both share the `in_progress` status). Re-queries /// whenever the task list changes. Given the single-active invariant, at most /// one task is running per user. final myRunningProgrammerTaskIdProvider = FutureProvider((ref) async { // Re-run whenever tasks change (a start/pause/resume mutates a row). ref.watch(programmerTasksProvider); final running = await ref.watch(programmerTasksControllerProvider).findRunningTaskForCurrentUser(); return running?.id; }); final programmerTaskActivityLogsProvider = StreamProvider.family, String>((ref, taskId) { final client = ref.watch(supabaseClientProvider); final wrapper = StreamRecoveryWrapper( stream: client .from('programmer_task_activity_logs') .stream(primaryKey: ['id']) .eq('task_id', taskId) .order('created_at', ascending: false), onPollData: () async { final data = await client .from('programmer_task_activity_logs') .select() .eq('task_id', taskId) .order('created_at', ascending: false); return data.map(ProgrammerTaskActivityLog.fromMap).toList(); }, fromMap: ProgrammerTaskActivityLog.fromMap, channelName: 'programmer_task_activity_logs_$taskId', onStatusChanged: ref .read(realtimeControllerProvider) .handleChannelStatus, onOfflineData: () async { final all = await cachedListFromBrick(); return all.where((l) => l.taskId == taskId).toList() ..sort((a, b) => b.createdAt.compareTo(a.createdAt)); }, onCacheMirror: (rows) => mirrorBatchToBrick( rows, tag: 'programmer_task_activity_logs', ), ); ref.onDispose(wrapper.dispose); return wrapper.stream.map((result) => result.data); }); // --------------------------------------------------------------------------- // Controller // --------------------------------------------------------------------------- final programmerTasksControllerProvider = Provider(( ref, ) { final client = ref.watch(supabaseClientProvider); return ProgrammerTasksController(client, ref); }); class ProgrammerTasksController { ProgrammerTasksController(this._client, [this._ref]); final dynamic _client; final Ref? _ref; bool get _isOnline => _ref?.read(isOnlineProvider) ?? true; /// Whether [taskId]'s most recent lifecycle event is a pause (i.e. the task /// is in_progress but not actively running). Future _isCurrentlyPaused(String taskId) async { try { final rows = await _client .from('programmer_task_activity_logs') .select('action_type, created_at') .eq('task_id', taskId) .inFilter('action_type', _pauseStateActions) .order('created_at', ascending: false) .limit(1); if (rows is List && rows.isNotEmpty) { return (rows.first['action_type']?.toString() ?? '') == 'paused'; } } catch (_) {} return false; } Future _insertLog( String taskId, String actionType, { Map? meta, }) async { final row = { 'task_id': taskId, 'actor_id': _client.auth.currentUser?.id, 'action_type': actionType, 'meta': ?meta, }; if (!_isOnline) { _queueLog(row); return; } try { await _client.from('programmer_task_activity_logs').insert(row); } catch (e) { if (isOfflineSaveError(e)) { _queueLog(row); } else { rethrow; } } } /// Creates a task. [assigneeId] defaults to the current user (self-logged); /// pass another programmer's id to assign it to them. Returns the task id. Future createTask({ required String title, required String category, String? description, String? assigneeId, String? projectId, int priority = 1, }) async { final userId = _client.auth.currentUser?.id; if (userId == null) throw Exception('Not authenticated'); final id = const Uuid().v4(); final assignee = assigneeId ?? userId; final payload = { 'id': id, 'title': title, 'description': description, 'category': category, 'status': ProgrammerTaskStatus.queued, 'priority': priority, 'assignee_id': assignee, 'creator_id': userId, 'project_id': ?projectId, }; try { await _client.from('programmer_tasks').insert(payload); await _insertLog(id, 'created'); if (assignee != userId) { await _insertLog(id, 'assigned', meta: {'assignee_id': assignee}); } return id; } catch (e) { if (!isOfflineSaveError(e)) rethrow; final now = DateTime.now().toUtc(); final local = ProgrammerTask( id: id, title: title, description: description, category: category, status: ProgrammerTaskStatus.queued, priority: priority, assigneeId: assignee, creatorId: userId, projectId: projectId, createdAt: now, updatedAt: now, ); _queueNewTask(local, payload); return id; } } /// The current user's actively-running task (status in_progress, latest event /// not a pause), or null. Used to drive the pause-on-switch prompt. Online /// only — returns null offline (the prompt is a best-effort UX aid). Future findRunningTaskForCurrentUser() async { final userId = _client.auth.currentUser?.id; if (userId == null || !_isOnline) return null; try { final rows = await _client .from('programmer_tasks') .select() .eq('assignee_id', userId) .eq('status', ProgrammerTaskStatus.inProgress); if (rows is! List) return null; for (final row in rows) { final task = ProgrammerTask.fromMap(row as Map); if (!await _isCurrentlyPaused(task.id)) return task; } } catch (_) {} return null; } /// Moves a queued task to in_progress (sets started_at on first start) and /// logs a `started` event. Future startTask({required String taskId}) async { final updates = { 'status': ProgrammerTaskStatus.inProgress, 'started_at': DateTime.now().toUtc().toIso8601String(), }; await _updateTaskRow(taskId, updates, onlyIfStartNull: true); await _insertLog(taskId, 'started'); } /// Logs a `paused` event for a running task (idempotent). Future pauseTask({required String taskId}) async { if (_isOnline && await _isCurrentlyPaused(taskId)) return; await _insertLog(taskId, 'paused'); } /// Logs a `resumed` event for a paused task (idempotent). Future resumeTask({required String taskId}) async { if (_isOnline && !await _isCurrentlyPaused(taskId)) return; await _insertLog(taskId, 'resumed'); } Future completeTask({required String taskId}) async { await _updateTaskRow(taskId, { 'status': ProgrammerTaskStatus.completed, 'completed_at': DateTime.now().toUtc().toIso8601String(), }); await _insertLog(taskId, 'completed'); } Future cancelTask({ required String taskId, required String reason, }) async { await _updateTaskRow(taskId, { 'status': ProgrammerTaskStatus.cancelled, 'cancelled_at': DateTime.now().toUtc().toIso8601String(), 'cancellation_reason': reason, }); await _insertLog(taskId, 'cancelled', meta: {'reason': reason}); } Future reassign({ required String taskId, required String newAssigneeId, }) async { await _updateTaskRow(taskId, {'assignee_id': newAssigneeId}); await _insertLog(taskId, 'reassigned', meta: {'assignee_id': newAssigneeId}); } Future updateTask({ required String taskId, String? title, String? description, String? category, int? priority, }) async { final updates = {}; if (title != null) updates['title'] = title; if (description != null) updates['description'] = description; if (category != null) updates['category'] = category; if (priority != null) updates['priority'] = priority; if (updates.isEmpty) return; await _updateTaskRow(taskId, updates); await _insertLog(taskId, 'updated', meta: {'fields': updates.keys.toList()}); } /// Assigns the task to a project (pass null to clear it). Future setProject({ required String taskId, required String? projectId, }) async { await _updateTaskRow(taskId, {'project_id': projectId}); await _insertLog(taskId, 'updated', meta: { 'fields': ['project_id'], }); } /// Records a time adjustment (negative [seconds]) against a task's worked /// duration — e.g. when the current user helped on another task and chose to /// deduct that time from their own running task. Future addAdjustment({ required String taskId, required int seconds, String? reason, String? sourceTaskId, }) async { await _insertLog(taskId, 'adjustment', meta: { 'seconds': seconds, 'reason': ?reason, 'source_task_id': ?sourceTaskId, }); } /// Applies a field patch to a task row. When [onlyIfStartNull] is set, /// `started_at` is dropped if the task already has one (preserves the first /// execution start across pause/resume cycles). Future _updateTaskRow( String taskId, Map updates, { bool onlyIfStartNull = false, }) async { var payload = updates; if (!_isOnline) { _queueUpdate(taskId, payload); return; } try { if (onlyIfStartNull && payload.containsKey('started_at')) { final existing = await _client .from('programmer_tasks') .select('started_at') .eq('id', taskId) .maybeSingle(); if (existing is Map && existing['started_at'] != null) { payload = Map.from(payload)..remove('started_at'); } } await _client.from('programmer_tasks').update(payload).eq('id', taskId); } catch (e) { if (isOfflineSaveError(e)) { _queueUpdate(taskId, payload); } else { rethrow; } } } // ------------------------------------------------------------------ queues void _queueNewTask(ProgrammerTask model, Map payload) { final ref = _ref; if (ref == null) return; ref.read(offlinePendingProgrammerTasksProvider.notifier).state = List.unmodifiable([ ...ref.read(offlinePendingProgrammerTasksProvider), model, ]); ref.read(offlinePendingProgrammerTaskRawProvider.notifier).state = List.unmodifiable([ ...ref.read(offlinePendingProgrammerTaskRawProvider), payload, ]); mirrorBatchToBrick([model], tag: 'offline_programmer_task'); } void _queueUpdate(String taskId, Map fields) { final ref = _ref; if (ref == null) return; final current = Map>.from( ref.read(offlinePendingProgrammerTaskUpdatesProvider), ); current[taskId] = {...?current[taskId], ...fields}; ref.read(offlinePendingProgrammerTaskUpdatesProvider.notifier).state = current; } void _queueLog(Map row) { final ref = _ref; if (ref == null) return; ref.read(offlinePendingProgrammerTaskLogsProvider.notifier).state = List.unmodifiable([ ...ref.read(offlinePendingProgrammerTaskLogsProvider), row, ]); } } // --------------------------------------------------------------------------- // Field-patch helper for optimistic offline updates // --------------------------------------------------------------------------- ProgrammerTask _applyUpdateFields(ProgrammerTask t, Map f) { DateTime? parseOpt(String key) { if (!f.containsKey(key)) return null; final v = f[key]; return v == null ? null : DateTime.parse(v as String).toUtc(); } return ProgrammerTask( id: t.id, title: (f['title'] as String?) ?? t.title, description: f.containsKey('description') ? f['description'] as String? : t.description, category: (f['category'] as String?) ?? t.category, status: (f['status'] as String?) ?? t.status, priority: (f['priority'] as num?)?.toInt() ?? t.priority, assigneeId: f.containsKey('assignee_id') ? f['assignee_id'] as String? : t.assigneeId, creatorId: t.creatorId, projectId: f.containsKey('project_id') ? f['project_id'] as String? : t.projectId, createdAt: t.createdAt, startedAt: f.containsKey('started_at') ? parseOpt('started_at') : t.startedAt, completedAt: f.containsKey('completed_at') ? parseOpt('completed_at') : t.completedAt, cancelledAt: f.containsKey('cancelled_at') ? parseOpt('cancelled_at') : t.cancelledAt, cancellationReason: f.containsKey('cancellation_reason') ? f['cancellation_reason'] as String? : t.cancellationReason, updatedAt: DateTime.now().toUtc(), ); }