263 lines
9.8 KiB
C#
263 lines
9.8 KiB
C#
using Microsoft.Extensions.Logging;
|
|
|
|
namespace ClawdDotNet.Core.Tasks;
|
|
|
|
/// <summary>
|
|
/// Führt eine fällige Aufgabe aus. Kapselt, wie ein Assignee zu einem Lauf wird — der
|
|
/// Scanner selbst kennt die Engine nicht und bleibt so ohne sie testbar.
|
|
/// </summary>
|
|
public interface ITaskDispatcher
|
|
{
|
|
/// <summary>Führt die Aufgabe aus. Gibt zurück, ob der Lauf regulär abschloss.</summary>
|
|
Task<bool> DispatchAsync(TaskItem task, CancellationToken ct);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Entscheidet, ob ein marktabhängiger Termin (<see cref="TaskItem.OnlyWhenMarketOpen"/>)
|
|
/// jetzt laufen darf. Platzhalter, bis C1 (Marktkalender) den echten Kalender liefert.
|
|
/// </summary>
|
|
public interface IMarketCalendar
|
|
{
|
|
bool IsOpen(DateTime nowUtc);
|
|
}
|
|
|
|
/// <summary>Bis C1 kommt: der Markt gilt als immer offen.</summary>
|
|
public sealed class AlwaysOpenMarketCalendar : IMarketCalendar
|
|
{
|
|
public bool IsOpen(DateTime nowUtc) => true;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Der Taktgeber des Taskboards (A1). Ein einziger Takt prüft, was fällig ist —
|
|
/// keine langen Delays (B6). Persistiert wird nur der Last-Fired-Marker; ein
|
|
/// fehlgeschlagener Lauf bleibt der einzige Versuch für diesen Termin (kein Retry-Sturm).
|
|
///
|
|
/// Der heikle Kern ist streng gegen drei Invarianten gebaut, die die Repository-Tests
|
|
/// und die Scanner-Tests absichern:
|
|
/// 1. Nie zwei Läufe auf denselben Termin — das atomare <see cref="ITaskRepository.TryClaimAsync"/>.
|
|
/// 2. Kein Dispatch bei offenem Blocker — blockierte Aufgaben sind nicht claimbar.
|
|
/// 3. Doppelter Takt = ein Lauf — konkurrierende Claims um denselben Occurrence-Key.
|
|
/// </summary>
|
|
public sealed class TaskScanner : IAsyncDisposable
|
|
{
|
|
private readonly ITaskRepository _repo;
|
|
private readonly ITaskDispatcher _dispatcher;
|
|
private readonly IMarketCalendar _market;
|
|
private readonly TimeProvider _clock;
|
|
private readonly ILogger _logger;
|
|
private readonly TimeSpan _tick;
|
|
private readonly TimeSpan _lease;
|
|
|
|
/// <summary>Aufgaben, deren unauflösbare Zeitzone bereits gemeldet wurde — der Takt
|
|
/// läuft jede Minute, die Meldung soll nicht mitlaufen.</summary>
|
|
private readonly HashSet<string> _warnedTimeZones = [];
|
|
|
|
private readonly CancellationTokenSource _cts = new();
|
|
private Task? _loop;
|
|
|
|
public TaskScanner(
|
|
ITaskRepository repo,
|
|
ITaskDispatcher dispatcher,
|
|
ILoggerFactory loggerFactory,
|
|
IMarketCalendar? market = null,
|
|
TimeProvider? clock = null,
|
|
TimeSpan? tick = null,
|
|
TimeSpan? lease = null)
|
|
{
|
|
_repo = repo;
|
|
_dispatcher = dispatcher;
|
|
_market = market ?? new AlwaysOpenMarketCalendar();
|
|
_clock = clock ?? TimeProvider.System;
|
|
_logger = loggerFactory.CreateLogger("ClawdDotNet.Core.Tasks.Scanner");
|
|
_tick = tick ?? TimeSpan.FromSeconds(60);
|
|
_lease = lease ?? TimeSpan.FromMinutes(15);
|
|
}
|
|
|
|
public void Start()
|
|
{
|
|
_loop ??= RunLoopAsync(_cts.Token);
|
|
}
|
|
|
|
/// <summary>Läuft die Taktschleife gerade?</summary>
|
|
public bool IsRunning => _loop is { IsCompleted: false };
|
|
|
|
/// <summary>
|
|
/// Die Schleife lief und ist beendet — anders als „noch nicht gestartet". Der
|
|
/// Unterschied zählt für die Zustandsmeldung an den Watchdog: Ein Scanner, der noch
|
|
/// auf die Startabgleichung wartet, ist in Ordnung; einer, dessen Schleife
|
|
/// ausgestiegen ist, bedeutet, dass keine Aufgabe mehr läuft.
|
|
/// </summary>
|
|
public bool HasStopped => _loop is { IsCompleted: true };
|
|
|
|
private async Task RunLoopAsync(CancellationToken ct)
|
|
{
|
|
// PeriodicTimer über den TimeProvider — im Test steuerbar, im Betrieb driftfrei.
|
|
using var timer = new PeriodicTimer(_tick, _clock);
|
|
while (await timer.WaitForNextTickAsync(ct))
|
|
{
|
|
try
|
|
{
|
|
await ScanOnceAsync(ct);
|
|
}
|
|
catch (Exception ex) when (ex is not OperationCanceledException)
|
|
{
|
|
_logger.LogError(ex, "Ein Scanner-Takt ist gescheitert");
|
|
}
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Ein Durchlauf: fällige Aufgaben beanspruchen und ausführen. Gibt die Zahl der
|
|
/// tatsächlich angestoßenen Läufe zurück. Öffentlich, damit Tests einen Takt
|
|
/// deterministisch auslösen können.
|
|
/// </summary>
|
|
public async Task<int> ScanOnceAsync(CancellationToken ct)
|
|
{
|
|
var now = _clock.GetUtcNow().UtcDateTime;
|
|
var leaseCutoff = now - _lease;
|
|
|
|
var candidates = await _repo.ListAsync(new TaskQuery { Limit = 500 }, ct);
|
|
|
|
var claimed = new List<(TaskItem Task, string Token)>();
|
|
foreach (var task in candidates)
|
|
{
|
|
if (task.Status != TaskItemStatus.Todo)
|
|
continue; // backlog (Halte-Status), blockiert, laufend, erledigt bleiben außen vor
|
|
|
|
// Ein Mensch ist kein Modell-Lauf — solche Aufgaben rührt der Scanner nicht an.
|
|
if (TaskAssignee.KindOf(task.Assignee) == TaskAssigneeKind.Human)
|
|
continue;
|
|
|
|
// Eine Aufgabe mit unauflösbarer Zeitzone feuert nie. Das einmal melden,
|
|
// sonst sucht man den Fehler bei der Aufgabe statt beim System.
|
|
if (TaskSchedule.UnresolvableTimeZone(task) is { } badZone
|
|
&& _warnedTimeZones.Add(task.Id))
|
|
{
|
|
_logger.LogWarning(
|
|
"Aufgabe {TaskId} ({Title}) hat die Zeitzone '{TimeZone}', die auf diesem "
|
|
+ "System nicht auflösbar ist — sie wird nicht ausgeführt. Zeitzone in "
|
|
+ "IANA-Schreibweise eintragen (z.B. 'Europe/Berlin') und sicherstellen, "
|
|
+ "dass tzdata und ICU vorhanden sind.",
|
|
task.Id, task.Title, badZone);
|
|
}
|
|
|
|
var occ = TaskSchedule.DueOccurrence(task, now);
|
|
if (occ is null)
|
|
continue;
|
|
|
|
if (task.OnlyWhenMarketOpen && !_market.IsOpen(now))
|
|
continue;
|
|
|
|
var token = Guid.NewGuid().ToString("N");
|
|
if (await _repo.TryClaimAsync(task.Id, occ, token, now, leaseCutoff, ct))
|
|
claimed.Add((task, token));
|
|
}
|
|
|
|
// Gleichzeitig anstoßen — Läufe desselben Agenten serialisiert ohnehin das
|
|
// Agent-Gate der Engine; verschiedene Agenten laufen echt parallel.
|
|
await Task.WhenAll(claimed.Select(c => ProcessAsync(c.Task, c.Token, ct)));
|
|
|
|
return claimed.Count;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Führt einen Task sofort aus (manueller „Jetzt ausführen"-Knopf), unabhängig vom
|
|
/// Termin. Beansprucht ihn atomar wie ein Takt — läuft er schon oder ist er kein
|
|
/// <c>todo</c>, gibt die Methode <c>false</c> zurück, statt ihn doppelt anzustoßen.
|
|
/// </summary>
|
|
public async Task<bool> RunTaskNowAsync(string taskId, CancellationToken ct)
|
|
{
|
|
var task = await _repo.GetAsync(taskId, ct);
|
|
if (task is null)
|
|
return false;
|
|
|
|
var now = _clock.GetUtcNow().UtcDateTime;
|
|
var token = Guid.NewGuid().ToString("N");
|
|
if (!await _repo.TryClaimAsync(task.Id, now.ToString("O"), token, now, now - _lease, ct))
|
|
return false;
|
|
|
|
await ProcessAsync(task, token, ct);
|
|
return true;
|
|
}
|
|
|
|
private async Task ProcessAsync(TaskItem task, string token, CancellationToken ct)
|
|
{
|
|
bool ok;
|
|
try
|
|
{
|
|
ok = await _dispatcher.DispatchAsync(task, ct);
|
|
}
|
|
catch (OperationCanceledException) when (ct.IsCancellationRequested)
|
|
{
|
|
throw;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogError(ex, "Dispatch für Aufgabe {TaskId} ist gescheitert", task.Id);
|
|
ok = false;
|
|
}
|
|
|
|
var now = _clock.GetUtcNow().UtcDateTime;
|
|
|
|
// Endstatus:
|
|
// - Fehlschlag → zurück auf todo (der Marker steht schon beim Claim, derselbe Termin
|
|
// wird nicht wiederholt).
|
|
// - Wiederkehrend (Cron/every/tool_job) → zurück auf todo, damit der nächste Termin
|
|
// feuern kann. Sonst liefe ein Cron-Task nur ein einziges Mal.
|
|
// - Einmalig → done bzw. in_review (A2).
|
|
var final = !ok || task.IsRecurring
|
|
? TaskItemStatus.Todo
|
|
: task.RequireApproval
|
|
? TaskItemStatus.InReview
|
|
: TaskItemStatus.Done;
|
|
|
|
await _repo.CompleteClaimAsync(task.Id, token, final, now, ct);
|
|
|
|
if (ok && final == TaskItemStatus.Done)
|
|
await UnblockDependentsAsync(task.Id, now, ct);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Auto-Dispatch: Wird eine Aufgabe fertig, werden die auf sie wartenden Aufgaben
|
|
/// freigegeben, sobald <b>alle</b> ihre Blocker erledigt sind.
|
|
/// </summary>
|
|
private async Task UnblockDependentsAsync(string completedId, DateTime now, CancellationToken ct)
|
|
{
|
|
var dependents = await _repo.ListBlockedByAsync(completedId, ct);
|
|
foreach (var dependent in dependents)
|
|
{
|
|
if (dependent.Status != TaskItemStatus.Blocked)
|
|
continue;
|
|
|
|
var allDone = true;
|
|
foreach (var blockerId in dependent.BlockedBy)
|
|
{
|
|
var blocker = await _repo.GetAsync(blockerId, ct);
|
|
if (blocker is null || blocker.Status != TaskItemStatus.Done)
|
|
{
|
|
allDone = false;
|
|
break;
|
|
}
|
|
}
|
|
|
|
if (allDone)
|
|
{
|
|
await _repo.SetStatusAsync(dependent.Id, TaskItemStatus.Todo, now, ct);
|
|
_logger.LogInformation(
|
|
"Aufgabe {TaskId} freigegeben — alle Blocker erledigt", dependent.Id);
|
|
}
|
|
}
|
|
}
|
|
|
|
public async ValueTask DisposeAsync()
|
|
{
|
|
await _cts.CancelAsync();
|
|
if (_loop is not null)
|
|
{
|
|
try { await _loop; }
|
|
catch (OperationCanceledException) { }
|
|
}
|
|
_cts.Dispose();
|
|
}
|
|
}
|