B2 beheben: Chat-Laeufe pro Agent serialisieren

Der Konversationskontext _chatContexts[agentId] ist eine geteilte List<ChatMessage>.
Nur der Lookup lief unter Lock, alle Add-Aufrufe im Schleifenkoerper waren
ungeschuetzt. Da WebView, ToolJob-Wakeups und AgentComm denselben Agenten
gleichzeitig ansprechen koennen, verschraenkten sich ihre Nachrichten zu einer
ungueltigen Tool-Sequenz, die die API mit HTTP 400 ablehnt.

ChatAsync laeuft jetzt hinter einem SemaphoreSlim(1,1) pro Agent; verschiedene
Agenten bleiben unabhaengig. Die Timeout-Uhr startet erst nach dem Eintritt,
damit Wartezeit in der Warteschlange den Lauf nicht aufzehrt.

Zwei Folgeprobleme mit demselben Ursprung:

- _runningChats hielt nur EINE CancellationTokenSource je Agent; der zweite Lauf
  ueberschrieb den ersten. AbortChat brach dadurch nur einen ab, der andere lief
  bis ins Run-Timeout. Jetzt eine Liste, die auch wartende Laeufe erfasst.
- ExecuteToolCallAsync fing OperationCanceledException mit ab und gab sie als
  Tool-Ergebnis zurueck, wodurch der Abbruch erst einen Schritt spaeter griff.
  Cancellation wird nun durchgereicht.

Ausserdem: send_message an den eigenen Agenten wird abgelehnt — es waere mit dem
neuen Gate in einen Deadlock gelaufen.

Neu: GetChatContext(agentId) als Momentaufnahme des Kontexts, fuer Diagnose und
Kontextgroessen-Anzeige.

Build-Fix: Das WinForms-Projekt globbt **/*.cs und kompilierte dadurch die
Test-Quellen mit. tests\** wird jetzt wie src\** ausgeschlossen.

Alle 49 Tests gruen.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Richard
2026-07-27 10:55:04 +02:00
co-authored by Claude Opus 4.8
parent eae13771cf
commit 6bbe9f9a80
8 changed files with 489 additions and 27 deletions
+134 -11
View File
@@ -21,9 +21,23 @@ public sealed class AgentEngine : IAgentMessageRouter
private readonly Dictionary<string, List<ChatEntry>> _chatHistories = new();
private readonly Dictionary<string, List<ChatMessage>> _chatContexts = new();
private readonly Dictionary<string, CancellationTokenSource> _runningChats = new();
private readonly Lock _lock = new();
/// <summary>
/// Alle aktiven Chat-Läufe je Agent — laufende wie wartende. Mehrere Quellen können
/// denselben Agenten gleichzeitig ansprechen (WebView, ToolJob, AgentComm), deshalb
/// eine Liste: AbortChat muss jeden davon erreichen.
/// </summary>
private readonly Dictionary<string, List<CancellationTokenSource>> _runningChats = new();
/// <summary>
/// Serialisiert ChatAsync pro Agent. Der Konversationskontext ist eine geteilte
/// Liste — liefen zwei Chats desselben Agenten gleichzeitig, verschränkten sich ihre
/// Nachrichten zu einer ungültigen Tool-Sequenz, die die API mit HTTP 400 ablehnt.
/// Verschiedene Agenten bleiben unabhängig voneinander.
/// </summary>
private readonly Dictionary<string, SemaphoreSlim> _agentGates = new();
private Func<IReadOnlyList<AgentConfig>>? _agentConfigProvider;
private Func<string, string?>? _agentDirResolver;
private string _instanceId = "";
@@ -227,18 +241,53 @@ public sealed class AgentEngine : IAgentMessageRouter
string instanceId,
CancellationToken externalCt,
string? source = null)
{
// Abbrechbar sein, schon bevor der Lauf an der Reihe ist — sonst hängt eine
// wartende Nachricht auch dann noch, wenn der Benutzer längst abgebrochen hat.
using var runCts = CancellationTokenSource.CreateLinkedTokenSource(externalCt);
RegisterRun(agentConfig.AgentId, runCts);
var gate = GetAgentGate(agentConfig.AgentId);
try
{
await gate.WaitAsync(runCts.Token);
}
catch (OperationCanceledException)
{
UnregisterRun(agentConfig.AgentId, runCts);
return new AgentRunResult(
agentConfig.AgentId, AgentRunStatus.Cancelled, "[Chat abgebrochen]",
0, 0, TimeSpan.Zero);
}
try
{
return await ChatCoreAsync(agentConfig, userMessage, instanceId, runCts.Token, source);
}
finally
{
UnregisterRun(agentConfig.AgentId, runCts);
gate.Release();
}
}
private async Task<AgentRunResult> ChatCoreAsync(
AgentConfig agentConfig,
string userMessage,
string instanceId,
CancellationToken runCt,
string? source)
{
var logger = _loggerFactory.CreateLogger($"ClawdDotNet.Core.Engine.Chat.{agentConfig.AgentId}");
var loopGuard = new LoopGuard(agentConfig.LoopGuard);
var sw = Stopwatch.StartNew();
// Die Timeout-Uhr läuft erst ab hier — Wartezeit in der Warteschlange
// darf den Lauf nicht aufzehren.
using var timeoutCts = new CancellationTokenSource(agentConfig.LoopGuard.Timeout);
using var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(externalCt, timeoutCts.Token);
using var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(runCt, timeoutCts.Token);
var ct = linkedCts.Token;
lock (_lock)
_runningChats[agentConfig.AgentId] = linkedCts;
try
{
var tools = _toolRegistry.GetForAgent(agentConfig);
@@ -376,10 +425,46 @@ public sealed class AgentEngine : IAgentMessageRouter
OnRunCompleted?.Invoke(agentConfig.Model, result);
return result;
}
finally
}
// ─── Nebenläufigkeits-Helfer ───
private SemaphoreSlim GetAgentGate(string agentId)
{
lock (_lock)
{
lock (_lock)
_runningChats.Remove(agentConfig.AgentId);
if (!_agentGates.TryGetValue(agentId, out var gate))
{
gate = new SemaphoreSlim(1, 1);
_agentGates[agentId] = gate;
}
return gate;
}
}
private void RegisterRun(string agentId, CancellationTokenSource cts)
{
lock (_lock)
{
if (!_runningChats.TryGetValue(agentId, out var list))
{
list = new List<CancellationTokenSource>();
_runningChats[agentId] = list;
}
list.Add(cts);
}
}
private void UnregisterRun(string agentId, CancellationTokenSource cts)
{
lock (_lock)
{
if (!_runningChats.TryGetValue(agentId, out var list))
return;
list.Remove(cts);
if (list.Count == 0)
_runningChats.Remove(agentId);
}
}
@@ -397,12 +482,38 @@ public sealed class AgentEngine : IAgentMessageRouter
return _runningChats.ContainsKey(agentId);
}
public void AbortChat(string agentId)
/// <summary>
/// Momentaufnahme des Konversationskontexts eines Agenten — also der Nachrichten,
/// die beim nächsten Schritt tatsächlich an das Modell gehen.
/// Nützlich für Diagnose und Kontextgrößen-Anzeige.
/// </summary>
public IReadOnlyList<ChatMessage> GetChatContext(string agentId)
{
lock (_lock)
return _chatContexts.TryGetValue(agentId, out var ctx)
? ctx.ToList()
: [];
}
/// <summary>
/// Bricht ALLE Chat-Läufe des Agenten ab — laufende wie wartende.
/// </summary>
public void AbortChat(string agentId)
{
List<CancellationTokenSource> toCancel;
lock (_lock)
{
if (_runningChats.TryGetValue(agentId, out var cts))
cts.Cancel();
if (!_runningChats.TryGetValue(agentId, out var list))
return;
toCancel = list.ToList();
}
// Außerhalb des Locks abbrechen: Cancel führt Continuations aus, die
// ihrerseits wieder auf _lock zugreifen können.
foreach (var cts in toCancel)
{
try { cts.Cancel(); }
catch (ObjectDisposedException) { /* Lauf war bereits fertig */ }
}
}
@@ -514,6 +625,12 @@ public sealed class AgentEngine : IAgentMessageRouter
if (configs is null)
return new AgentMessageResult(false, null, "Agent config provider not set.");
// Selbstadressierung würde am Agent-Gate hängen bleiben: Der laufende Chat
// hält es bereits und würde auf sich selbst warten.
if (fromAgentId == toAgentId)
return new AgentMessageResult(false, null,
"Ein Agent kann sich keine Nachricht an sich selbst schicken.");
var targetConfig = configs.FirstOrDefault(a => a.AgentId == toAgentId);
if (targetConfig is null)
return new AgentMessageResult(false, null, $"Agent '{toAgentId}' not found.");
@@ -678,6 +795,12 @@ public sealed class AgentEngine : IAgentMessageRouter
logger.LogWarning("Tool access denied: {Message}", ex.Message);
return JsonSerializer.Serialize(new { error = ex.Message });
}
catch (OperationCanceledException) when (ct.IsCancellationRequested)
{
// Nicht als Tool-Fehler zurückgeben: Sonst läuft die Schleife noch einen
// Schritt weiter und der Abbruch greift erst verzögert.
throw;
}
catch (Exception ex)
{
logger.LogError(ex, "Tool {Tool} threw an exception", toolName);