fix: stream session updates live

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
This commit is contained in:
David Kaya
2026-03-21 16:12:48 +01:00
co-authored by Copilot
parent 0b15af444d
commit 9ba0174f68
6 changed files with 327 additions and 49 deletions
@@ -49,7 +49,6 @@ public sealed class CopilotWorkflowRunner : ITurnWorkflowRunner
List<ChatMessageDto> completedMessages = [];
AgentIdentity? activeAgent = null;
HashSet<string> startedAgents = new(StringComparer.OrdinalIgnoreCase);
HashSet<string> completedAgents = new(StringComparer.OrdinalIgnoreCase);
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, inputMessages).ConfigureAwait(false);
await run.TrySendMessageAsync(new TurnToken(emitEvents: true)).ConfigureAwait(false);
@@ -81,7 +80,7 @@ public sealed class CopilotWorkflowRunner : ITurnWorkflowRunner
await onActivity(activity).ConfigureAwait(false);
}
}
else if (evt is AgentResponseUpdateEvent update && !string.IsNullOrEmpty(update.Update.Text))
else if (evt is AgentResponseUpdateEvent update)
{
AgentIdentity? updateAgent = null;
string authorName = update.ExecutorId;
@@ -94,10 +93,6 @@ public sealed class CopilotWorkflowRunner : ITurnWorkflowRunner
authorName = resolvedUpdateAgent.AgentName;
}
string messageId = update.Update.MessageId ?? $"{command.RequestId}-delta-{fallbackMessageIndex++}";
StreamingSegment segment = GetOrCreateSegment(segments, messageId, authorName);
segment.Content.Append(update.Update.Text);
if (updateAgent.HasValue)
{
activeAgent = updateAgent.Value;
@@ -108,6 +103,15 @@ public sealed class CopilotWorkflowRunner : ITurnWorkflowRunner
onActivity).ConfigureAwait(false);
}
if (string.IsNullOrEmpty(update.Update.Text))
{
continue;
}
string messageId = update.Update.MessageId ?? $"{command.RequestId}-delta-{fallbackMessageIndex++}";
StreamingSegment segment = GetOrCreateSegment(segments, messageId, authorName);
segment.Content.Append(update.Update.Text);
await onDelta(new TurnDeltaEventDto
{
Type = "turn-delta",
@@ -124,14 +128,6 @@ public sealed class CopilotWorkflowRunner : ITurnWorkflowRunner
completed.ExecutorId,
out AgentIdentity completedAgent))
{
if (completedAgents.Add(completedAgent.AgentId))
{
await onActivity(CreateActivityEvent(
command,
activityType: "completed",
agent: completedAgent)).ConfigureAwait(false);
}
if (activeAgent.HasValue
&& string.Equals(activeAgent.Value.AgentId, completedAgent.AgentId, StringComparison.Ordinal))
{
@@ -143,11 +139,6 @@ public sealed class CopilotWorkflowRunner : ITurnWorkflowRunner
List<ChatMessage> allMessages = outputEvent.As<List<ChatMessage>>() ?? [];
List<ChatMessage> newMessages = allMessages.Skip(inputMessages.Count).ToList();
completedMessages = ConvertOutputMessages(command, newMessages, segments);
await EmitCompletedActivitiesForMessages(
command,
completedMessages,
completedAgents,
onActivity).ConfigureAwait(false);
}
}
@@ -214,34 +205,6 @@ public sealed class CopilotWorkflowRunner : ITurnWorkflowRunner
agent: agent)).ConfigureAwait(false);
}
private static async Task EmitCompletedActivitiesForMessages(
RunTurnCommandDto command,
IReadOnlyList<ChatMessageDto> messages,
ISet<string> completedAgents,
Func<AgentActivityEventDto, Task> onActivity)
{
foreach (ChatMessageDto message in messages)
{
if (!AgentIdentityResolver.TryResolveKnownAgentIdentity(
command.Pattern,
message.AuthorName,
out AgentIdentity messageAgent))
{
continue;
}
if (!completedAgents.Add(messageAgent.AgentId))
{
continue;
}
await onActivity(CreateActivityEvent(
command,
activityType: "completed",
agent: messageAgent)).ConfigureAwait(false);
}
}
private static bool TryGetHandoffTarget(
PatternDefinitionDto pattern,
RequestInfoEvent requestInfo,