From 8a903430678573538b5b3421886b2b21e4420faf Mon Sep 17 00:00:00 2001 From: Constantinos Leftheris Date: Thu, 16 Jul 2026 19:29:16 +0300 Subject: [PATCH 1/5] Migrate to Microsoft.Extensions.AI abstractions Refactored codebase to replace custom chat models and Azure.AI.OpenAI usage with Microsoft.Extensions.AI types and OpenAI client integration. Updated service registration, workflow steps, and persistence logic to use new abstractions. Removed obsolete files and updated dependencies to Microsoft.Extensions.AI 10.8.0. Streamlined message and event models for improved compatibility and maintainability. --- .../AgentsFeatureExtensions.cs | 26 +++-- .../Data/DbMessage.cs | 6 +- .../Data/Mappings/DbMessageMap.cs | 14 +-- .../Indice.Features.Agents.Core.csproj | 6 +- .../Models/ChatMessage.cs | 49 +++++--- .../Models/ChatMessageContent.cs | 30 +++++ .../Models/ChatMessageMappings.cs | 21 ---- .../Models/ChatMessagePart.cs | 27 ----- .../Models/ChatMessageRole.cs | 21 ---- .../Models/ChatRequest.cs} | 5 +- .../Models/ChatResponse.cs | 35 ------ .../Models/ChatStreamEvent.cs | 2 +- .../Models/Session.cs | 2 + .../Services/ChatsService.cs | 109 +++++++++++------- .../Services/IChatsService.cs | 3 +- .../Services/ISessionsStore.cs | 1 + .../Services/SessionsStore.cs | 86 +++----------- .../Workflows/Abstractions/IDexRunner.cs | 2 + .../Workflows/DexRunner.cs | 12 +- .../Workflows/Events/AnswerDeltaEvent.cs | 2 +- .../Workflows/{ => Events}/DexStreamEvent.cs | 17 ++- .../Reranking/LlmListwiseReranker.cs | 13 ++- .../SessionStoreChatHistoryProvider.cs | 3 +- .../Workflows/Steps/AnswerComposer.cs | 11 +- .../Workflows/Steps/IntentClassifier.cs | 7 +- .../Workflows/Steps/PurposeResponder.cs | 11 +- .../Workflows/Steps/QueryRewriter.cs | 9 +- .../Endpoints/MyChatsApi.cs | 1 + .../Endpoints/MyChatsHandlers.cs | 7 +- .../Endpoints/MyChatsValidators.cs | 1 + 30 files changed, 237 insertions(+), 302 deletions(-) create mode 100644 src/Indice.Features.Agents.Core/Models/ChatMessageContent.cs delete mode 100644 src/Indice.Features.Agents.Core/Models/ChatMessageMappings.cs delete mode 100644 src/Indice.Features.Agents.Core/Models/ChatMessageRole.cs rename src/{Indice.Features.Agents.Server/Endpoints/MyChatsModels.cs => Indice.Features.Agents.Core/Models/ChatRequest.cs} (83%) delete mode 100644 src/Indice.Features.Agents.Core/Models/ChatResponse.cs rename src/Indice.Features.Agents.Core/Workflows/{ => Events}/DexStreamEvent.cs (80%) diff --git a/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs b/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs index e9a22ff9f..87624b3a8 100644 --- a/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs +++ b/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs @@ -1,5 +1,4 @@ using System.ClientModel; -using Azure.AI.OpenAI; using Indice.Features.Agents.Core; using Indice.Features.Agents.Core.Data; using Indice.Features.Agents.Core.Services; @@ -11,6 +10,7 @@ using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection.Extensions; using Microsoft.Extensions.Options; +using OpenAI; using static Indice.Features.Agents.Core.AgentsOptions; namespace Microsoft.Extensions.DependencyInjection; @@ -20,7 +20,7 @@ public static class AgentsFeatureExtensions { /// /// Registers Dex core services: (bound from the Dex configuration section), - /// (singleton — each pipeline step builds its own role-bound agent from it), + /// (singleton — each pipeline step builds its own role-bound agent from it), /// the embedding generator, the wired to SQL Server, and . /// public static IServiceCollection AddAgentsCore(this IServiceCollection services, IConfiguration configuration, Action? configureAction = null) { @@ -36,17 +36,25 @@ public static IServiceCollection AddAgentsCore(this IServiceCollection services, agents.Value.ConfigureModelOptions?.Invoke(models); }); - services.TryAddSingleton(sp => { + services.AddSingleton(sp => { var opts = sp.GetRequiredService>().Value.AzureOpenAI; - return new AzureOpenAIClient(new Uri(opts.Endpoint!), new ApiKeyCredential(opts.ApiKey!)); + return new OpenAIClient(new ApiKeyCredential(opts.ApiKey!), new OpenAIClientOptions { Endpoint = new Uri(opts.Endpoint!) }); + }); + services.AddKeyedChatClient(nameof(AzureOpenAIDeployments.Reasoning), sp => { + var opts = sp.GetRequiredService>().Value.AzureOpenAI.Deployments; + var innerClient = sp.GetRequiredService(); + return innerClient.GetChatClient(opts.Reasoning).AsIChatClient(); + }); + services.AddKeyedChatClient(nameof(AzureOpenAIDeployments.Fast), sp => { + var opts = sp.GetRequiredService>().Value.AzureOpenAI.Deployments; + var innerClient = sp.GetRequiredService(); + return innerClient.GetChatClient(opts.Fast).AsIChatClient(); }); - services.TryAddSingleton>>(sp => { + services.AddEmbeddingGenerator(sp => { var opts = sp.GetRequiredService>().Value.AzureOpenAI; - var client = sp.GetRequiredService(); - return client - .GetEmbeddingClient(opts.Deployments.Embedding!) - .AsIEmbeddingGenerator(opts.EmbeddingDimensions); + var innerClient = sp.GetRequiredService(); + return innerClient.GetEmbeddingClient(opts.Deployments.Embedding!).AsIEmbeddingGenerator(opts.EmbeddingDimensions); }); services.AddDbContext((sp, options) => { diff --git a/src/Indice.Features.Agents.Core/Data/DbMessage.cs b/src/Indice.Features.Agents.Core/Data/DbMessage.cs index 4d72eb414..ff3293cd7 100644 --- a/src/Indice.Features.Agents.Core/Data/DbMessage.cs +++ b/src/Indice.Features.Agents.Core/Data/DbMessage.cs @@ -1,4 +1,4 @@ -using Indice.Features.Agents.Core.Models; +using Microsoft.Extensions.AI; namespace Indice.Features.Agents.Core.Data; @@ -12,10 +12,10 @@ public class DbMessage public Guid SessionId { get; set; } /// Author role of this message. Persisted as the role's string value (e.g. user). - public ChatMessageRole Role { get; set; } = ChatMessageRole.User; + public ChatRole Role { get; set; } = ChatRole.User; /// Message body. - public ChatMessageContent Content { get; set; } = new (); + public List Contents { get; set; } = []; /// Creation timestamp. public DateTimeOffset CreatedAt { get; set; } diff --git a/src/Indice.Features.Agents.Core/Data/Mappings/DbMessageMap.cs b/src/Indice.Features.Agents.Core/Data/Mappings/DbMessageMap.cs index 5e1b81a78..08bc2a7ec 100644 --- a/src/Indice.Features.Agents.Core/Data/Mappings/DbMessageMap.cs +++ b/src/Indice.Features.Agents.Core/Data/Mappings/DbMessageMap.cs @@ -1,11 +1,7 @@ -using System.Data; -using System.Text.Json; using Indice.Configuration; -using Indice.Features.Agents.Core.Models; using Microsoft.EntityFrameworkCore; -using Microsoft.EntityFrameworkCore.ChangeTracking; -using Microsoft.EntityFrameworkCore.Metadata; using Microsoft.EntityFrameworkCore.Metadata.Builders; +using Microsoft.Extensions.AI; namespace Indice.Features.Agents.Core.Data.Mappings; @@ -18,14 +14,14 @@ public void Configure(EntityTypeBuilder builder) { builder.HasKey(x => x.Id); builder.Property(x => x.Role) .HasConversion( - role => role.ToString().ToLowerInvariant(), - value => Enum.Parse(value, ignoreCase: true)) + role => role.ToString(), + value => new ChatRole(value)) .HasMaxLength(TextSizePresets.S32) .IsRequired(); //TODO: check whats going on here with this version. builder.ComplexProperty(b => b.Content, b => b.ToJson().IsRequired()); - builder.ComplexProperty(b => b.Content, b => b.ToJson().IsRequired()); - //builder.Property(x => x.Content).HasRequiredJsonConversion(); + //builder.ComplexProperty(b => b.Contents, b => b.ToJson().IsRequired()); + builder.Property(x => x.Contents).HasRequiredJsonConversion(); builder.Property(x => x.ModelUsed).HasMaxLength(TextSizePresets.M128); builder.HasIndex(x => new { x.SessionId, x.CreatedAt }); diff --git a/src/Indice.Features.Agents.Core/Indice.Features.Agents.Core.csproj b/src/Indice.Features.Agents.Core/Indice.Features.Agents.Core.csproj index a5c4ab43d..203acee35 100644 --- a/src/Indice.Features.Agents.Core/Indice.Features.Agents.Core.csproj +++ b/src/Indice.Features.Agents.Core/Indice.Features.Agents.Core.csproj @@ -12,16 +12,14 @@ - - - + + - diff --git a/src/Indice.Features.Agents.Core/Models/ChatMessage.cs b/src/Indice.Features.Agents.Core/Models/ChatMessage.cs index 7e2e9f34b..67520f9b9 100644 --- a/src/Indice.Features.Agents.Core/Models/ChatMessage.cs +++ b/src/Indice.Features.Agents.Core/Models/ChatMessage.cs @@ -1,25 +1,38 @@ -using OpenAI.Assistants; - namespace Indice.Features.Agents.Core.Models; -/// A single turn (user or assistant) in a chat session. DTO exposed at the service boundary; mirrors . -public class ChatMessage -{ - /// Message identifier. - public Guid Id { get; init; } +///// A single turn (user or assistant) in a chat session. DTO exposed at the service boundary; mirrors . +//public class ChatMessage +//{ +// /// Identifier of the session this turn belongs to. +// public Guid SessionId { get; init; } + +// /// Identifier of the assistant message persisted for this turn. +// public Guid MessageId { get; init; } - /// Author role of this message. Serializes as the role's lowercase value (e.g. user). - public ChatMessageRole Role { get; init; } = ChatMessageRole.User; +// /// Author role of this message. Serializes as the role's lowercase value (e.g. user). +// public ChatMessageRole Role { get; init; } = ChatMessageRole.User; - /// Message body. - public ChatMessageContent Content { get; init; } = new(); +// /// Message body. +// public ChatMessageContent Content { get; init; } = new(); - /// Creation timestamp. - public DateTimeOffset CreatedAt { get; init; } +// /// Creation timestamp. +// public DateTimeOffset CreatedAt { get; init; } - /// References to chunks. - public List Citations { get; set; } = []; +// /// References to chunks. +// public List Citations { get; set; } = []; - /// References to source documents. - public List Sources { get; set; } = []; -} +// /// References to source documents. +// public List Sources { get; set; } = []; + +// /// True when a pipeline step threw and the workflow halted. Out-of-scope is NOT a failure — its refusal text flows through . +// public bool Failed { get; init; } + +// /// True when the turn was blocked by a session usage limit — carries the predefined limit message, nothing was persisted, and is empty. +// public bool LimitReached { get; init; } + +// /// Questions used in this session so far, for a used/total display. null when the message limit is disabled. +// public int? QuestionsUsed { get; init; } + +// /// Total questions allowed per session, for a used/total display. null when the message limit is disabled. +// public int? QuestionsTotal { get; init; } +//} diff --git a/src/Indice.Features.Agents.Core/Models/ChatMessageContent.cs b/src/Indice.Features.Agents.Core/Models/ChatMessageContent.cs new file mode 100644 index 000000000..017a2f1e1 --- /dev/null +++ b/src/Indice.Features.Agents.Core/Models/ChatMessageContent.cs @@ -0,0 +1,30 @@ +using System.Text.Json.Serialization; + +namespace Indice.Features.Agents.Core.Models; + +/// Represents the content of a chat message. +public class ChatMessageContent +{ + /// Creates a new instance of . + public ChatMessageContent() { + + } + /// Creates a new instance of with a single part. + public ChatMessageContent(string content, string contentType = "text/markdown") { + AddPart(content, contentType); + } + + + /// Parts of the message content. + [JsonPropertyName("parts")] + public List Parts { get; set; } = []; + + /// + /// Adds a new part to the message content. + /// + /// The value of the message part. + /// The content type of the message part. + public void AddPart(string value, string contentType) { + Parts.Add(ChatMessagePart.FromText(value, contentType)); + } +} diff --git a/src/Indice.Features.Agents.Core/Models/ChatMessageMappings.cs b/src/Indice.Features.Agents.Core/Models/ChatMessageMappings.cs deleted file mode 100644 index 830588590..000000000 --- a/src/Indice.Features.Agents.Core/Models/ChatMessageMappings.cs +++ /dev/null @@ -1,21 +0,0 @@ -using AITextContent = Microsoft.Extensions.AI.TextContent; -using AIChatMessage = Microsoft.Extensions.AI.ChatMessage; -using AIChatRole = Microsoft.Extensions.AI.ChatRole; - -namespace Indice.Features.Agents.Core.Models; - -/// Projections from the persisted chat message DTOs onto the framework (Microsoft.Extensions.AI) shapes the MAF pipeline consumes. -internal static class ChatMessageMappings -{ - /// Projects a persisted session turn into the framework message shape. - public static AIChatMessage ToAIChatMessage(this ChatMessage message) => new(ToAIChatRole(message.Role), message.Content.Parts.Select(p => new AITextContent(p.Value) { AdditionalProperties = new() { ["contentType"] = p.ContentType } }).ToArray()); - - /// Maps the persisted onto the framework role the MAF pipeline consumes. - public static AIChatRole ToAIChatRole(ChatMessageRole role) => role switch { - ChatMessageRole.User => AIChatRole.User, - ChatMessageRole.Assistant => AIChatRole.Assistant, - ChatMessageRole.System => AIChatRole.System, - ChatMessageRole.Tool => AIChatRole.Tool, - _ => AIChatRole.User, - }; -} diff --git a/src/Indice.Features.Agents.Core/Models/ChatMessagePart.cs b/src/Indice.Features.Agents.Core/Models/ChatMessagePart.cs index 79247ebc0..c80af8637 100644 --- a/src/Indice.Features.Agents.Core/Models/ChatMessagePart.cs +++ b/src/Indice.Features.Agents.Core/Models/ChatMessagePart.cs @@ -4,33 +4,6 @@ namespace Indice.Features.Agents.Core.Models; -/// Represents the content of a chat message. -public class ChatMessageContent -{ - /// Creates a new instance of . - public ChatMessageContent() { - - } - /// Creates a new instance of with a single part. - public ChatMessageContent(string content, string contentType = "text/markdown") { - AddPart(content, contentType); - } - - - /// Parts of the message content. - [JsonPropertyName("parts")] - public List Parts { get; set; } = []; - - /// - /// Adds a new part to the message content. - /// - /// The value of the message part. - /// The content type of the message part. - public void AddPart(string value, string contentType) { - Parts.Add(ChatMessagePart.FromText(value, contentType)); - } -} - /// Represents a part of a chat message. public class ChatMessagePart { diff --git a/src/Indice.Features.Agents.Core/Models/ChatMessageRole.cs b/src/Indice.Features.Agents.Core/Models/ChatMessageRole.cs deleted file mode 100644 index 6d3f950a2..000000000 --- a/src/Indice.Features.Agents.Core/Models/ChatMessageRole.cs +++ /dev/null @@ -1,21 +0,0 @@ -using System.Text.Json.Serialization; - -namespace Indice.Features.Agents.Core.Models; - -/// Author role of a chat message. Exposed at the service boundary and persisted per turn; serializes as its lowercase value (e.g. user). -[JsonConverter(typeof(JsonStringEnumConverter))] -public enum ChatMessageRole -{ - /// Message authored by the end user. - [JsonStringEnumMemberName("user")] - User, - /// Message authored by the assistant. - [JsonStringEnumMemberName("assistant")] - Assistant, - /// System / developer instruction message. - [JsonStringEnumMemberName("system")] - System, - /// Tool invocation or tool result message. - [JsonStringEnumMemberName("tool")] - Tool -} diff --git a/src/Indice.Features.Agents.Server/Endpoints/MyChatsModels.cs b/src/Indice.Features.Agents.Core/Models/ChatRequest.cs similarity index 83% rename from src/Indice.Features.Agents.Server/Endpoints/MyChatsModels.cs rename to src/Indice.Features.Agents.Core/Models/ChatRequest.cs index 5815e26bb..a4364cc3a 100644 --- a/src/Indice.Features.Agents.Server/Endpoints/MyChatsModels.cs +++ b/src/Indice.Features.Agents.Core/Models/ChatRequest.cs @@ -1,9 +1,8 @@ -namespace Indice.Features.Agents.Server.Endpoints; +namespace Indice.Features.Agents.Core.Models; /// Body accepted by both POST /api/my/chats (creates the session inline) and POST /api/my/chats/{id}/messages. public class ChatRequest { /// The end-user message text. public string Text { get; init; } = string.Empty; -} - +} \ No newline at end of file diff --git a/src/Indice.Features.Agents.Core/Models/ChatResponse.cs b/src/Indice.Features.Agents.Core/Models/ChatResponse.cs deleted file mode 100644 index e1e4a4609..000000000 --- a/src/Indice.Features.Agents.Core/Models/ChatResponse.cs +++ /dev/null @@ -1,35 +0,0 @@ -namespace Indice.Features.Agents.Core.Models; - -/// Response returned by both POST /api/my/chats and POST /api/my/chats/{id}/messages. -public class ChatResponse -{ - /// Identifier of the session this turn belongs to. - public Guid SessionId { get; init; } - - /// Identifier of the assistant message persisted for this turn. - public Guid MessageId { get; init; } - - /// The pipeline's answer — grounded when in-scope, or the polite out-of-scope refusal text when not. - public string? Answer { get; init; } - - /// Citations supporting the answer; empty for out-of-scope responses and on error. - public IReadOnlyList Citations { get; init; } = []; - - /// Links to the source documents that were retrieved and used to compose the answer; empty for out-of-scope responses and on error. - public IReadOnlyList Sources { get; init; } = []; - - /// True when a pipeline step threw and the workflow halted. Out-of-scope is NOT a failure — its refusal text flows through . - public bool Failed { get; init; } - - /// Error message from the step that threw; null when is false. - public string? FailureReason { get; init; } - - /// True when the turn was blocked by a session usage limit — carries the predefined limit message, nothing was persisted, and is empty. - public bool LimitReached { get; init; } - - /// Questions used in this session so far, for a used/total display. null when the message limit is disabled. - public int? QuestionsUsed { get; init; } - - /// Total questions allowed per session, for a used/total display. null when the message limit is disabled. - public int? QuestionsTotal { get; init; } -} diff --git a/src/Indice.Features.Agents.Core/Models/ChatStreamEvent.cs b/src/Indice.Features.Agents.Core/Models/ChatStreamEvent.cs index 6167fa341..237ee567d 100644 --- a/src/Indice.Features.Agents.Core/Models/ChatStreamEvent.cs +++ b/src/Indice.Features.Agents.Core/Models/ChatStreamEvent.cs @@ -33,7 +33,7 @@ public record ChatStreamEvent public Guid? SessionId { get; init; } /// Identifier of the persisted assistant message; populated on complete. - public Guid? MessageId { get; init; } + public string? MessageId { get; init; } /// True when a pipeline step threw and the workflow halted; populated on complete. public bool? Failed { get; init; } diff --git a/src/Indice.Features.Agents.Core/Models/Session.cs b/src/Indice.Features.Agents.Core/Models/Session.cs index 2e0a49c9c..097b933d3 100644 --- a/src/Indice.Features.Agents.Core/Models/Session.cs +++ b/src/Indice.Features.Agents.Core/Models/Session.cs @@ -1,3 +1,5 @@ +using Microsoft.Extensions.AI; + namespace Indice.Features.Agents.Core.Models; /// Detail view of a chat session, including the most recent messages. diff --git a/src/Indice.Features.Agents.Core/Services/ChatsService.cs b/src/Indice.Features.Agents.Core/Services/ChatsService.cs index 0f5b572c3..6e7b12103 100644 --- a/src/Indice.Features.Agents.Core/Services/ChatsService.cs +++ b/src/Indice.Features.Agents.Core/Services/ChatsService.cs @@ -3,7 +3,9 @@ using Indice.Features.Agents.Core.Models; using Indice.Features.Agents.Core.Workflows; using Indice.Features.Agents.Core.Workflows.Abstractions; +using Indice.Features.Agents.Core.Workflows.Events; using Indice.Types; +using Microsoft.Extensions.AI; using Microsoft.Extensions.Options; namespace Indice.Features.Agents.Core.Services; @@ -27,7 +29,7 @@ public ChatsService(ISessionsStore store, IDexRunner runner, IUsageGuardService } /// - public async Task SendAsync(string userId, Guid? sessionId, string text, CancellationToken cancellationToken) { + public async Task SendAsync(string userId, Guid? sessionId, ChatRequest chatRequest, CancellationToken cancellationToken) { await EnsureSessionCreationAllowedAsync(userId, sessionId, cancellationToken); var session = await _store.LoadOrCreateAsync(userId, sessionId, cancellationToken); if (session is null) { @@ -35,31 +37,38 @@ public ChatsService(ISessionsStore store, IDexRunner runner, IUsageGuardService } var turnCheck = _usageGuard.Check(session); if (!turnCheck.Allowed) { - return new ChatResponse { - SessionId = session.Id, - MessageId = Guid.Empty, - Answer = turnCheck.Message, - LimitReached = true, - QuestionsUsed = _sessionOptions.GetQuestionsUsed(session.MessageCount), - QuestionsTotal = _sessionOptions.GetQuestionsTotal(), + return new ChatResponse(new ChatMessage(ChatRole.Assistant, turnCheck.Message)) { + ConversationId = session.Id.ToString(), + ResponseId = Guid.NewGuid().ToString(), + Usage = new () { + AdditionalCounts = new() { + ["questionsUsed"] = _sessionOptions.GetQuestionsUsed(session.MessageCount) ?? 0, + ["questionsTotal"] = _sessionOptions.GetQuestionsTotal() ?? 0 + } + }, + AdditionalProperties = new() { + ["limitReached"] = true + } }; } - var request = new RagRequest { Question = text, SessionId = session.Id }; + var request = new RagRequest { Question = chatRequest.Text, SessionId = session.Id }; var result = await _runner.RunAsync(request, cancellationToken); - var userMessage = new ChatMessage { - Id = Guid.NewGuid(), - Role = ChatMessageRole.User, - Content = new ChatMessageContent { Parts = [new () { Value = request.Question, ContentType = "text" }] }, - CreatedAt = request.TimeStamp, + var userMessage = new ChatMessage(ChatRole.User, request.Question) { + MessageId = Guid.NewGuid().ToString(), + CreatedAt = request.TimeStamp + }; + var assistantText = new TextContent(result.Answer ?? string.Empty) { + Annotations = result.Citations?.Select(c => (AIAnnotation)new CitationAnnotation { + FileId = c.DocumentId.ToString(), + Title = result.Sources.First(x => x.Id == c.DocumentId).SourceTitle, + Url = new Uri(result.Sources.First(x => x.Id == c.DocumentId).SourceUrl), + Snippet = c.Title + // Other metadata like page number, confidence, etc. + }).ToList() ?? [] }; - var assistantText = result.Answer ?? string.Empty; - var assistantMessage = new ChatMessage { - Id = Guid.NewGuid(), - Role = ChatMessageRole.Assistant, - Content = new ChatMessageContent { Parts = [new () { Value = assistantText, ContentType = "text" }] }, + var assistantMessage = new ChatMessage(ChatRole.Assistant, [assistantText]) { + MessageId = Guid.NewGuid().ToString(), CreatedAt = DateTimeOffset.UtcNow, - Citations = result.Citations?.ToList() ?? [], - Sources = result.Sources?.ToList() ?? [], }; var persistedAssistant = await _store.AppendTurnAsync(session.Id, userMessage, assistantMessage, @@ -67,15 +76,23 @@ public ChatsService(ISessionsStore store, IDexRunner runner, IUsageGuardService modelUsed: result.ModelUsed ?? _deployments.Reasoning, cancellationToken); return new ChatResponse { - SessionId = session.Id, - MessageId = persistedAssistant.Id, - Answer = assistantText, - Citations = result.Citations ?? [], - Sources = result.Sources ?? [], - Failed = result.Failed, - FailureReason = result.FailureReason, - QuestionsUsed = _sessionOptions.GetQuestionsUsed(session.MessageCount + 2), - QuestionsTotal = _sessionOptions.GetQuestionsTotal(), + ConversationId = session.Id.ToString(), + ResponseId = persistedAssistant.MessageId, + Messages = [assistantMessage], + ModelId = result.ModelUsed ?? _deployments.Reasoning, + Usage = new() { + InputTokenCount = result.Usage?.InputTokenCount ?? 0, + OutputTokenCount = result.Usage?.OutputTokenCount ?? 0, + AdditionalCounts = new() { + ["questionsUsed"] = _sessionOptions.GetQuestionsUsed(session.MessageCount) ?? 0, + ["questionsTotal"] = _sessionOptions.GetQuestionsTotal() ?? 0 + } + }, + AdditionalProperties = new() { + ["failed"] = result.Failed.ToString().ToLower(), + ["failureReason"] = result.FailureReason, + ["sources"] = result.Sources ?? [], + } }; } @@ -109,7 +126,7 @@ private async IAsyncEnumerable> LimitReachedStream(Sess yield return new SseItem(new ChatStreamEvent { Type = "complete", SessionId = session.Id, - MessageId = Guid.Empty, + MessageId = Guid.Empty.ToString(), Answer = message, Citations = [], Sources = [], @@ -146,15 +163,23 @@ private async IAsyncEnumerable> StreamTurnAsync( /// Persists the user/assistant turn (mirroring ) and builds the terminal complete event. private async Task PersistTurnAsync(Session session, RagRequest request, DexFinalEvent? final, CancellationToken cancellationToken) { - var assistantText = final?.Answer ?? string.Empty; - var userMessage = new ChatMessage { Id = Guid.NewGuid(), Role = ChatMessageRole.User, Content = new ChatMessageContent(request.Question), CreatedAt = request.TimeStamp }; - var assistantMessage = new ChatMessage { - Id = Guid.NewGuid(), - Role = ChatMessageRole.Assistant, - Content = new ChatMessageContent(assistantText), - CreatedAt = final?.TimeStamp ?? DateTimeOffset.UtcNow, - Citations = final?.Citations?.ToList() ?? [], - Sources = final?.Sources?.ToList() ?? [], + var userMessage = new ChatMessage(ChatRole.User, request.Question) { + MessageId = Guid.NewGuid().ToString(), + CreatedAt = request.TimeStamp + }; + + var assistantText = new TextContent(final?.Answer ?? string.Empty) { + Annotations = final?.Citations?.Select(c => (AIAnnotation)new CitationAnnotation { + FileId = c.DocumentId.ToString(), + Title = final?.Sources.First(x => x.Id == c.DocumentId).SourceTitle, + Url = new Uri(final!.Sources.First(x => x.Id == c.DocumentId).SourceUrl), + Snippet = c.Title + // Other metadata like page number, confidence, etc. + }).ToList() ?? [] + }; + var assistantMessage = new ChatMessage(ChatRole.Assistant, [assistantText]) { + MessageId = Guid.NewGuid().ToString(), + CreatedAt = DateTimeOffset.UtcNow, }; var persistedAssistant = await _store.AppendTurnAsync(session.Id, userMessage, assistantMessage, @@ -164,8 +189,8 @@ private async Task PersistTurnAsync(Session session, RagRequest return new ChatStreamEvent { Type = "complete", SessionId = session.Id, - MessageId = persistedAssistant.Id, - Answer = assistantText, + MessageId = persistedAssistant.MessageId, + Answer = assistantText.Text, Citations = final?.Citations ?? [], Sources = final?.Sources ?? [], Failed = final?.Failed ?? false, diff --git a/src/Indice.Features.Agents.Core/Services/IChatsService.cs b/src/Indice.Features.Agents.Core/Services/IChatsService.cs index 2d7592086..cb29f1316 100644 --- a/src/Indice.Features.Agents.Core/Services/IChatsService.cs +++ b/src/Indice.Features.Agents.Core/Services/IChatsService.cs @@ -1,6 +1,7 @@ using System.Net.ServerSentEvents; using Indice.Features.Agents.Core.Models; using Indice.Types; +using Microsoft.Extensions.AI; namespace Indice.Features.Agents.Core.Services; @@ -11,7 +12,7 @@ public interface IChatsService /// Posts a turn. When is null, the session is created inline as part of this call. /// Returns null when is supplied but no session matches . /// - Task SendAsync(string userId, Guid? sessionId, string text, CancellationToken cancellationToken); + Task SendAsync(string userId, Guid? sessionId, ChatRequest chatRequest, CancellationToken cancellationToken); /// /// Streaming counterpart of : posts a turn and returns the live SSE event stream diff --git a/src/Indice.Features.Agents.Core/Services/ISessionsStore.cs b/src/Indice.Features.Agents.Core/Services/ISessionsStore.cs index 47f74435c..c53c34193 100644 --- a/src/Indice.Features.Agents.Core/Services/ISessionsStore.cs +++ b/src/Indice.Features.Agents.Core/Services/ISessionsStore.cs @@ -1,5 +1,6 @@ using Indice.Features.Agents.Core.Models; using Indice.Types; +using Microsoft.Extensions.AI; using static Indice.Features.Agents.Core.AgentsOptions; namespace Indice.Features.Agents.Core.Services; diff --git a/src/Indice.Features.Agents.Core/Services/SessionsStore.cs b/src/Indice.Features.Agents.Core/Services/SessionsStore.cs index bd996bf14..08a3e41d5 100644 --- a/src/Indice.Features.Agents.Core/Services/SessionsStore.cs +++ b/src/Indice.Features.Agents.Core/Services/SessionsStore.cs @@ -2,6 +2,7 @@ using Indice.Features.Agents.Core.Models; using Indice.Types; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.AI; using Microsoft.Extensions.Options; using static Indice.Features.Agents.Core.AgentsOptions; @@ -11,13 +12,11 @@ namespace Indice.Features.Agents.Core.Services; public class SessionsStore : ISessionsStore { private readonly AgentsDbContext _db; - private readonly ISourceLinkGenerator _sourceLinkGenerator; private readonly SessionOptions _sessionOptions; /// Creates a new . - public SessionsStore(AgentsDbContext db, IOptions options, ISourceLinkGenerator sourceLinkGenerator) { + public SessionsStore(AgentsDbContext db, IOptions options) { _db = db ?? throw new ArgumentNullException(nameof(db)); - _sourceLinkGenerator = sourceLinkGenerator ?? throw new ArgumentNullException(nameof(sourceLinkGenerator)); _sessionOptions = options.Value.Session; } @@ -80,37 +79,14 @@ public SessionsStore(AgentsDbContext db, IOptions options, ISourc .Take(messageTake) .OrderBy(m => m.CreatedAt) .Select(m => new ChatMessage { - Id = m.Id, + MessageId = m.Id.ToString(), Role = m.Role, - Content = m.Content, - CreatedAt = m.CreatedAt, - Citations = m.Citations.Select(c => new Citation { - ChunkId = c.ChunkId, - DocumentId = c.Chunk.DocumentId, - Title = c.Chunk.Title, - HeadingPath = c.Chunk.HeadingPath, - Number = c.Number, - Score = c.Score, - }).ToList(), - Sources = m.Citations.Select(c => new SourceDocumentLink { - Id = c.Chunk.DocumentId, - SourceTitle = c.Chunk.Document.Title, - SourceUrl = _sourceLinkGenerator.GenerateLink(c.Chunk.Document.Source), - IsPrivate = c.Chunk.Document.IsPrivate, - ContentHash = c.Chunk.Document.ContentHash, - ContentType = c.Chunk.Document.Blob == null ? "application/markdown" : c.Chunk.Document.Blob.ContentType, - Length = c.Chunk.Document.Blob == null ? -1 : c.Chunk.Document.Blob.ContentLength, - FileName = c.Chunk.Document.Blob == null ? c.Chunk.Document.Title : c.Chunk.Document.Blob.FileName, - }).ToList() - }) + Contents = m.Contents, + CreatedAt = m.CreatedAt + }) .ToList(), }) .FirstOrDefaultAsync(cancellationToken); - if (session is not null) { - foreach (var message in session.Messages) { - message.Sources = message.Sources.DistinctBy(s => s.Id).ToList(); - } - } return session; } @@ -125,33 +101,12 @@ public async Task> GetHistoryAsync(Guid sessionId, Ca .Take(messageTake) .OrderBy(m => m.CreatedAt) .Select(m => new ChatMessage { - Id = m.Id, + MessageId = m.Id.ToString(), Role = m.Role, - Content = m.Content, + Contents = m.Contents, CreatedAt = m.CreatedAt, - Citations = m.Citations.Select(c => new Citation { - ChunkId = c.ChunkId, - DocumentId = c.Chunk.DocumentId, - Title = c.Chunk.Title, - HeadingPath = c.Chunk.HeadingPath, - Number = c.Number, - Score = c.Score, - }).ToList(), - Sources = m.Citations.Select(c => new SourceDocumentLink { - Id = c.Chunk.DocumentId, - SourceTitle = c.Chunk.Document.Title, - SourceUrl = _sourceLinkGenerator.GenerateLink(c.Chunk.Document.Source), - IsPrivate = c.Chunk.Document.IsPrivate, - ContentHash = c.Chunk.Document.ContentHash, - ContentType = c.Chunk.Document.Blob == null ? "application/markdown" : c.Chunk.Document.Blob.ContentType, - Length = c.Chunk.Document.Blob == null ? -1 : c.Chunk.Document.Blob.ContentLength, - FileName = c.Chunk.Document.Blob == null ? c.Chunk.Document.Title : c.Chunk.Document.Blob.FileName, - }).ToList() }) .ToListAsync(cancellationToken); - foreach (var message in messages) { - message.Sources = message.Sources.DistinctBy(s => s.Id).ToList(); - } return messages; } @@ -182,20 +137,20 @@ public async Task AppendTurnAsync(Guid sessionId, ChatMessage userM _db.Add(userRow); _db.Add(assistantRow); - session.LastActivityAt = assistantMessage.CreatedAt; + session.LastActivityAt = assistantMessage.CreatedAt ?? DateTimeOffset.UtcNow; session.TotalPromptTokens += promptTokens; session.TotalCompletionTokens += completionTokens; session.MessageCount += 2; if (session.Title is null && _sessionOptions.TitleAutoGenerate) { - session.Title = DeriveTitle(userMessage.Content); + session.Title = DeriveTitle(userMessage); } await _db.SaveChangesAsync(cancellationToken); return new ChatMessage { - Id = assistantRow.Id, + MessageId = assistantRow.Id.ToString(), Role = assistantRow.Role, - Content = assistantRow.Content, + Contents = assistantRow.Contents, CreatedAt = assistantRow.CreatedAt, }; } @@ -238,20 +193,15 @@ public async Task GetUsageTokensAsync(string userId, DateTimeOffset since, } private static DbMessage ToDb(Guid sessionId, ChatMessage m, int? prompt, int? completion, string? model) => new() { - Id = m.Id == Guid.Empty ? Guid.NewGuid() : m.Id, + Id = string.IsNullOrWhiteSpace(m.MessageId) || !Guid.TryParse(m.MessageId, out var parsedId) ? Guid.NewGuid() : parsedId, SessionId = sessionId, Role = m.Role, - Content = m.Content, - CreatedAt = m.CreatedAt == default ? DateTimeOffset.UtcNow : m.CreatedAt, + Contents = m.Contents.ToList(), + CreatedAt = m.CreatedAt ?? DateTimeOffset.UtcNow, PromptTokens = prompt, CompletionTokens = completion, ModelUsed = model, - MetadataJson = null, - Citations = m.Citations.Select(c => new DbCitation { - ChunkId = c.ChunkId, - Number = c.Number, - Score = c.Score, - }).ToList(), + MetadataJson = null }; private Session ToDto(DbSession s, IReadOnlyList messages) => new() { @@ -267,8 +217,8 @@ public async Task GetUsageTokensAsync(string userId, DateTimeOffset since, Messages = messages, }; - private static string DeriveTitle(ChatMessageContent firstUserMessage) { - var normalized = firstUserMessage.Parts.FirstOrDefault()?.Value.Replace('\r', ' ').Replace('\n', ' ').Trim() ?? string.Empty; + private static string DeriveTitle(ChatMessage firstUserMessage) { + var normalized = firstUserMessage.Text.Replace('\r', ' ').Replace('\n', ' ').Trim() ?? string.Empty; return normalized.Length <= 80 ? normalized : normalized[..80]; } } diff --git a/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexRunner.cs b/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexRunner.cs index 06dc7e936..c7f09b483 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexRunner.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexRunner.cs @@ -1,3 +1,5 @@ +using Indice.Features.Agents.Core.Workflows.Events; + namespace Indice.Features.Agents.Core.Workflows.Abstractions; /// Entry point for executing the Dex RAG pipeline against a single user question. diff --git a/src/Indice.Features.Agents.Core/Workflows/DexRunner.cs b/src/Indice.Features.Agents.Core/Workflows/DexRunner.cs index 265624c8f..bf156cb08 100644 --- a/src/Indice.Features.Agents.Core/Workflows/DexRunner.cs +++ b/src/Indice.Features.Agents.Core/Workflows/DexRunner.cs @@ -11,14 +11,14 @@ namespace Indice.Features.Agents.Core.Workflows; /// public class DexRunner : IDexRunner { - private readonly Workflow? workflow; + private readonly Workflow _workflow; /// /// Creates a new instance. /// /// The workflow instance to execute. - public DexRunner ([FromKeyedServices("Default")] Workflow? workflow) { - this.workflow = workflow; + public DexRunner ([FromKeyedServices("Default")] Workflow workflow) { + _workflow = workflow; } /// Human-friendly progress labels keyed by executor id, surfaced as SSE step events. @@ -35,7 +35,7 @@ public DexRunner ([FromKeyedServices("Default")] Workflow? workflow) { /// public async Task RunAsync(RagRequest request, CancellationToken cancellationToken) { var initial = CreateInitialEnvelope(request); - await using var run = await InProcessExecution.RunAsync(workflow!, initial, cancellationToken: cancellationToken); + await using var run = await InProcessExecution.RunAsync(_workflow, initial, cancellationToken: cancellationToken); PipelineStepContext? final = null; string? failure = null; UsageDetails? usage = null; @@ -83,7 +83,7 @@ public async Task RunAsync(RagRequest request, CancellationToken canc public async IAsyncEnumerable RunStreamingAsync( RagRequest request, [EnumeratorCancellation] CancellationToken cancellationToken) { var initial = CreateInitialEnvelope(request); - await using var run = await InProcessExecution.RunStreamingAsync(workflow!, initial, cancellationToken: cancellationToken); + await using var run = await InProcessExecution.RunStreamingAsync(_workflow!, initial, cancellationToken: cancellationToken); PipelineStepContext? final = null; string? failure = null; UsageDetails? usage = null; @@ -131,7 +131,7 @@ public async IAsyncEnumerable RunStreamingAsync( /// Validates a workflow is registered and builds the initial pipeline envelope from . private PipelineStepContext CreateInitialEnvelope(RagRequest request) { - if (workflow is null) { + if (_workflow is null) { throw new InvalidOperationException( "No RAG workflow registered. Call services.AddDefaultDexPipeline() or register a Microsoft.Agents.AI.Workflows.Workflow manually."); } diff --git a/src/Indice.Features.Agents.Core/Workflows/Events/AnswerDeltaEvent.cs b/src/Indice.Features.Agents.Core/Workflows/Events/AnswerDeltaEvent.cs index 7c96c8876..34264d6d6 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Events/AnswerDeltaEvent.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Events/AnswerDeltaEvent.cs @@ -8,7 +8,7 @@ namespace Indice.Features.Agents.Core.Workflows.Events; /// the caller as an SSE delta event. The non-streaming run path ignores it. /// /// Creates a new carrying a single answer text delta. -public class AnswerDeltaEvent(string delta) : WorkflowEvent(delta) +public class AnswerDeltaEvent(string executorId, string delta) : ExecutorEvent(executorId, delta) { /// The incremental answer text produced by this streaming update. public string Delta => (string)Data!; diff --git a/src/Indice.Features.Agents.Core/Workflows/DexStreamEvent.cs b/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs similarity index 80% rename from src/Indice.Features.Agents.Core/Workflows/DexStreamEvent.cs rename to src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs index 9a6fed4b1..eb958e7bc 100644 --- a/src/Indice.Features.Agents.Core/Workflows/DexStreamEvent.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs @@ -1,7 +1,7 @@ using Indice.Features.Agents.Core.Models; using Microsoft.Extensions.AI; -namespace Indice.Features.Agents.Core.Workflows; +namespace Indice.Features.Agents.Core.Workflows.Events; /// /// Base type for the real-time events yields as @@ -40,6 +40,21 @@ public sealed class DexDeltaEvent : DexStreamEvent public string Text { get; } } +/// Signals that a pipeline step threw an exception. The pipeline will halt and yield a single with set to true. +public sealed class DexErrorEvent : DexStreamEvent +{ + /// Creates a new . + public DexErrorEvent(string stepId, string errorMessage) { + StepId = stepId; + ErrorMessage = errorMessage; + } + /// The executor id of the step that threw (e.g. Retriever). + public string StepId { get; } + + /// The error message from the step that threw. + public string ErrorMessage { get; } +} + /// /// Terminal event yielded once after the run completes. Mirrors the fields of : /// the full answer, citations, failure state, and reasoning-model token totals. diff --git a/src/Indice.Features.Agents.Core/Workflows/Reranking/LlmListwiseReranker.cs b/src/Indice.Features.Agents.Core/Workflows/Reranking/LlmListwiseReranker.cs index 183d5ca22..176d0b20f 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Reranking/LlmListwiseReranker.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Reranking/LlmListwiseReranker.cs @@ -1,11 +1,12 @@ using System.Text; -using Azure.AI.OpenAI; using Indice.Features.Agents.Core.Workflows.Abstractions; using Indice.Features.Agents.Core.Workflows.Prompts; using Microsoft.Agents.AI; using Microsoft.Extensions.AI; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using OpenAI.Chat; +using static Indice.Features.Agents.Core.AgentsOptions; namespace Indice.Features.Agents.Core.Workflows.Reranking; @@ -19,15 +20,15 @@ public class LlmListwiseReranker : ILlmReranker private readonly int _snippetLength; /// Creates a new . - public LlmListwiseReranker(AzureOpenAIClient openAIClient, IOptions options, - IOptions models, IPromptTemplateRenderer prompts) { + public LlmListwiseReranker([FromKeyedServices(nameof(AzureOpenAIDeployments.Fast))] IChatClient chatClient, + IOptions options, + IOptions models, + IPromptTemplateRenderer prompts) { var opts = options.Value; _snippetLength = opts.Retrieval.RerankSnippetLength; var chatOptions = models.Value.BaseFastModelOptions.Clone(); chatOptions.Instructions = prompts.Render("Reranker"); - _agent = openAIClient - .GetChatClient(opts.AzureOpenAI.Deployments.Fast!) - .AsAIAgent(options: new ChatClientAgentOptions() { + _agent = chatClient.AsAIAgent(options: new ChatClientAgentOptions() { ChatOptions = chatOptions, Name = "DexReranker", }); diff --git a/src/Indice.Features.Agents.Core/Workflows/SessionStoreChatHistoryProvider.cs b/src/Indice.Features.Agents.Core/Workflows/SessionStoreChatHistoryProvider.cs index 27a54e97d..e3ef00452 100644 --- a/src/Indice.Features.Agents.Core/Workflows/SessionStoreChatHistoryProvider.cs +++ b/src/Indice.Features.Agents.Core/Workflows/SessionStoreChatHistoryProvider.cs @@ -1,4 +1,3 @@ -using Indice.Features.Agents.Core.Models; using Indice.Features.Agents.Core.Services; using Microsoft.Agents.AI; @@ -46,7 +45,7 @@ public static void SetSessionId(AgentSession agentSession, Guid sessionId) return []; } var history = await _store.GetHistoryAsync(sessionId, cancellationToken); - return history.Select(message => message.ToAIChatMessage()); + return history; } /// diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/AnswerComposer.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/AnswerComposer.cs index e7c56c1e5..e4d237f58 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/AnswerComposer.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/AnswerComposer.cs @@ -1,13 +1,14 @@ using System.Text; -using Azure.AI.OpenAI; using Indice.Features.Agents.Core.Workflows.Events; using Indice.Features.Agents.Core.Workflows.Prompts; using Indice.Features.Agents.Core.Workflows.State; using Microsoft.Agents.AI; using Microsoft.Agents.AI.Workflows; using Microsoft.Extensions.AI; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using OpenAI.Chat; +using static Indice.Features.Agents.Core.AgentsOptions; namespace Indice.Features.Agents.Core.Workflows.Steps; @@ -23,7 +24,7 @@ public sealed class AnswerComposer : Executor, private readonly string _model; /// Creates a new . - public AnswerComposer(AzureOpenAIClient openAIClient, IOptions options, + public AnswerComposer([FromKeyedServices(nameof(AzureOpenAIDeployments.Reasoning))] IChatClient chatClient, IOptions options, IOptions models, IPromptTemplateRenderer prompts, UserClaimsAIContextProvider userClaimsProvider, SessionStoreChatHistoryProvider historyProvider) : base("AnswerComposer") { @@ -35,9 +36,7 @@ public AnswerComposer(AzureOpenAIClient openAIClient, IOptions op strictGrounding = _options.Pipeline.StrictGrounding, }); - _agent = openAIClient - .GetChatClient(_model) - .AsAIAgent( + _agent = chatClient.AsAIAgent( options: new ChatClientAgentOptions() { ChatOptions = chatOptions, AIContextProviders = [userClaimsProvider], @@ -63,7 +62,7 @@ public override async ValueTask> HandleAs await foreach (var update in _agent.RunStreamingAsync(prompt, agentSession, cancellationToken: cancellationToken)) { if (!string.IsNullOrEmpty(update.Text)) { answer.Append(update.Text); - await context.AddEventAsync(new AnswerDeltaEvent(update.Text), cancellationToken); + await context.AddEventAsync(new AnswerDeltaEvent(Id, update.Text), cancellationToken); } foreach (var usageContent in update.Contents.OfType()) { (usage ??= new UsageDetails()).Add(usageContent.Details); diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/IntentClassifier.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/IntentClassifier.cs index ca31cc699..c40c191d8 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/IntentClassifier.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/IntentClassifier.cs @@ -1,10 +1,10 @@ -using Azure.AI.OpenAI; using Indice.Features.Agents.Core.Workflows.Events; using Indice.Features.Agents.Core.Workflows.Prompts; using Indice.Features.Agents.Core.Workflows.State; using Microsoft.Agents.AI; using Microsoft.Agents.AI.Workflows; using Microsoft.Extensions.AI; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using OpenAI.Chat; using static Indice.Features.Agents.Core.AgentsOptions; @@ -23,7 +23,7 @@ public sealed class IntentClassifier : ExecutorCreates a new . - public IntentClassifier(AzureOpenAIClient openAIClient, IOptions options, + public IntentClassifier([FromKeyedServices(nameof(AzureOpenAIDeployments.Reasoning))] IChatClient chatClient, IOptions options, IOptions models, IPromptTemplateRenderer prompts, SessionStoreChatHistoryProvider historyProvider) : base("IntentClassifier") { _options = options.Value; @@ -33,8 +33,7 @@ public IntentClassifier(AzureOpenAIClient openAIClient, IOptions categories = _options.Taxonomy.Categories, languages = _options.Taxonomy.Languages, }); - _agent = openAIClient.GetChatClient(_model) - .AsAIAgent(options: new ChatClientAgentOptions() { + _agent = chatClient.AsAIAgent(options: new ChatClientAgentOptions() { ChatOptions = chatOptions, Name = "DexIntentClassifier", ChatHistoryProvider = historyProvider, diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/PurposeResponder.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/PurposeResponder.cs index f084a7e93..c82a0ccaa 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/PurposeResponder.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/PurposeResponder.cs @@ -1,14 +1,14 @@ using System.Text; -using Azure; -using Azure.AI.OpenAI; using Indice.Features.Agents.Core.Workflows.Events; using Indice.Features.Agents.Core.Workflows.Prompts; using Indice.Features.Agents.Core.Workflows.State; using Microsoft.Agents.AI; using Microsoft.Agents.AI.Workflows; using Microsoft.Extensions.AI; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using OpenAI.Chat; +using static Indice.Features.Agents.Core.AgentsOptions; namespace Indice.Features.Agents.Core.Workflows.Steps; @@ -24,7 +24,7 @@ internal class PurposeResponder : Executor, Pi /// Creates a new . - public PurposeResponder(AzureOpenAIClient openAIClient, IOptions options, + public PurposeResponder([FromKeyedServices(nameof(AzureOpenAIDeployments.Reasoning))] IChatClient chatClient, IOptions options, IOptions models, IPromptTemplateRenderer prompts, UserClaimsAIContextProvider userClaimsProvider, SessionStoreChatHistoryProvider historyProvider) : base("PurposeResponder") { @@ -36,8 +36,7 @@ public PurposeResponder(AzureOpenAIClient openAIClient, IOptions strictGrounding = _options.Pipeline.StrictGrounding, }); - _agent = openAIClient - .GetChatClient(_model) + _agent = chatClient .AsAIAgent( options: new ChatClientAgentOptions() { ChatOptions = chatOptions, @@ -64,7 +63,7 @@ public override async ValueTask> HandleAs await foreach (var update in _agent.RunStreamingAsync(prompt, agentSession, cancellationToken: cancellationToken)) { if (!string.IsNullOrEmpty(update.Text)) { answer.Append(update.Text); - await context.AddEventAsync(new AnswerDeltaEvent(update.Text), cancellationToken); + await context.AddEventAsync(new AnswerDeltaEvent(Id, update.Text), cancellationToken); } foreach (var usageContent in update.Contents.OfType()) { (usage ??= new UsageDetails()).Add(usageContent.Details); diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/QueryRewriter.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/QueryRewriter.cs index 0fb1fe975..d76c7ae61 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/QueryRewriter.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/QueryRewriter.cs @@ -1,12 +1,12 @@ -using System.ClientModel; -using Azure.AI.OpenAI; using Indice.Features.Agents.Core.Workflows.Prompts; using Indice.Features.Agents.Core.Workflows.State; using Microsoft.Agents.AI; using Microsoft.Agents.AI.Workflows; using Microsoft.Extensions.AI; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using OpenAI.Chat; +using static Indice.Features.Agents.Core.AgentsOptions; namespace Indice.Features.Agents.Core.Workflows.Steps; @@ -22,15 +22,14 @@ public sealed class QueryRewriter : Executor, private readonly string _model; /// Creates a new . - public QueryRewriter(AzureOpenAIClient openAIClient, IOptions options, + public QueryRewriter([FromKeyedServices(nameof(AzureOpenAIDeployments.Fast))] IChatClient chatClient, IOptions options, IOptions models, IPromptTemplateRenderer prompts, SessionStoreChatHistoryProvider historyProvider) : base("QueryRewriter") { _options = options.Value; _model = _options.AzureOpenAI.Deployments.Fast!; var chatOptions = models.Value.BaseFastModelOptions.Clone(); chatOptions.Instructions = prompts.Render("QueryRewriter"); - _agent = openAIClient - .GetChatClient(_model) + _agent = chatClient .AsAIAgent(options: new ChatClientAgentOptions() { ChatOptions = chatOptions, Name = "DexQueryRewriter", diff --git a/src/Indice.Features.Agents.Server/Endpoints/MyChatsApi.cs b/src/Indice.Features.Agents.Server/Endpoints/MyChatsApi.cs index fb090d65b..0b5cac761 100644 --- a/src/Indice.Features.Agents.Server/Endpoints/MyChatsApi.cs +++ b/src/Indice.Features.Agents.Server/Endpoints/MyChatsApi.cs @@ -1,3 +1,4 @@ +using Indice.Features.Agents.Core.Models; using Indice.Features.Agents.Server; using Indice.Features.Agents.Server.Endpoints; using Indice.Security; diff --git a/src/Indice.Features.Agents.Server/Endpoints/MyChatsHandlers.cs b/src/Indice.Features.Agents.Server/Endpoints/MyChatsHandlers.cs index 4a23ce2de..582f743f2 100644 --- a/src/Indice.Features.Agents.Server/Endpoints/MyChatsHandlers.cs +++ b/src/Indice.Features.Agents.Server/Endpoints/MyChatsHandlers.cs @@ -5,6 +5,7 @@ using Indice.Types; using Microsoft.AspNetCore.Http; using Microsoft.AspNetCore.Http.HttpResults; +using Microsoft.Extensions.AI; namespace Indice.Features.Agents.Server.Endpoints; @@ -16,15 +17,15 @@ internal static class MyChatsHandlers /// POST /api/my/chats — creates a session with the first question. public static async Task> Create(ChatRequest request, ClaimsPrincipal user, IChatsService chats, CancellationToken cancellationToken) { var userId = user.FindSubjectId()!; - var response = await chats.SendAsync(userId, sessionId: null, request.Text, cancellationToken); - return TypedResults.CreatedAtRoute(response, nameof(GetChatSession), new { chatId = response!.SessionId }); + var response = await chats.SendAsync(userId, sessionId: null, request, cancellationToken); + return TypedResults.CreatedAtRoute(response, nameof(GetChatSession), new { chatId = Guid.Parse(response!.ConversationId!) }); } /// POST /api/my/chats/{chatId}/messages — posts a follow-up turn. public static async Task, NotFound>> SendMessage(Guid chatId, ChatRequest request, ClaimsPrincipal user, IChatsService chats, CancellationToken cancellationToken) { var userId = user.FindSubjectId()!; - var response = await chats.SendAsync(userId, chatId, request.Text, cancellationToken); + var response = await chats.SendAsync(userId, chatId, request, cancellationToken); return response is null ? TypedResults.NotFound() : TypedResults.Ok(response); } diff --git a/src/Indice.Features.Agents.Server/Endpoints/MyChatsValidators.cs b/src/Indice.Features.Agents.Server/Endpoints/MyChatsValidators.cs index 01917e43f..7e3bba186 100644 --- a/src/Indice.Features.Agents.Server/Endpoints/MyChatsValidators.cs +++ b/src/Indice.Features.Agents.Server/Endpoints/MyChatsValidators.cs @@ -1,4 +1,5 @@ using FluentValidation; +using Indice.Features.Agents.Core.Models; namespace Indice.Features.Agents.Server.Endpoints; From d7c6ddddf31a5c2a75907e3dff3d9cf80a4869e3 Mon Sep 17 00:00:00 2001 From: Constantinos Leftheris Date: Thu, 16 Jul 2026 19:50:11 +0300 Subject: [PATCH 2/5] Switch to AzureOpenAIClient and update dependencies Replaced OpenAIClient with AzureOpenAIClient for chat and embedding services. Updated all related service registrations. Added Azure.AI.OpenAI 2.9.0-beta.1 NuGet package. --- .../AgentsFeatureExtensions.cs | 9 +++++---- .../Indice.Features.Agents.Core.csproj | 2 ++ 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs b/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs index 87624b3a8..45e988fb9 100644 --- a/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs +++ b/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs @@ -1,4 +1,5 @@ using System.ClientModel; +using Azure.AI.OpenAI; using Indice.Features.Agents.Core; using Indice.Features.Agents.Core.Data; using Indice.Features.Agents.Core.Services; @@ -38,22 +39,22 @@ public static IServiceCollection AddAgentsCore(this IServiceCollection services, services.AddSingleton(sp => { var opts = sp.GetRequiredService>().Value.AzureOpenAI; - return new OpenAIClient(new ApiKeyCredential(opts.ApiKey!), new OpenAIClientOptions { Endpoint = new Uri(opts.Endpoint!) }); + return new AzureOpenAIClient(new Uri(opts.Endpoint!), new ApiKeyCredential(opts.ApiKey!)); }); services.AddKeyedChatClient(nameof(AzureOpenAIDeployments.Reasoning), sp => { var opts = sp.GetRequiredService>().Value.AzureOpenAI.Deployments; - var innerClient = sp.GetRequiredService(); + var innerClient = sp.GetRequiredService(); return innerClient.GetChatClient(opts.Reasoning).AsIChatClient(); }); services.AddKeyedChatClient(nameof(AzureOpenAIDeployments.Fast), sp => { var opts = sp.GetRequiredService>().Value.AzureOpenAI.Deployments; - var innerClient = sp.GetRequiredService(); + var innerClient = sp.GetRequiredService(); return innerClient.GetChatClient(opts.Fast).AsIChatClient(); }); services.AddEmbeddingGenerator(sp => { var opts = sp.GetRequiredService>().Value.AzureOpenAI; - var innerClient = sp.GetRequiredService(); + var innerClient = sp.GetRequiredService(); return innerClient.GetEmbeddingClient(opts.Deployments.Embedding!).AsIEmbeddingGenerator(opts.EmbeddingDimensions); }); diff --git a/src/Indice.Features.Agents.Core/Indice.Features.Agents.Core.csproj b/src/Indice.Features.Agents.Core/Indice.Features.Agents.Core.csproj index 203acee35..ad9cf3c61 100644 --- a/src/Indice.Features.Agents.Core/Indice.Features.Agents.Core.csproj +++ b/src/Indice.Features.Agents.Core/Indice.Features.Agents.Core.csproj @@ -16,6 +16,8 @@ + + From f9d05e8bcb7b9720aa11b331817baa7a78fa289e Mon Sep 17 00:00:00 2001 From: Constantinos Leftheris Date: Fri, 17 Jul 2026 12:24:19 +0300 Subject: [PATCH 3/5] Preserve type for "failed" and remove unused using Changed "failed" in AdditionalProperties to retain its original type in ChatsService.cs. Removed an unused using directive from OutOfScopeResponder.cs for cleaner code. --- src/Indice.Features.Agents.Core/Services/ChatsService.cs | 2 +- .../Workflows/Steps/OutOfScopeResponder.cs | 1 - 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/src/Indice.Features.Agents.Core/Services/ChatsService.cs b/src/Indice.Features.Agents.Core/Services/ChatsService.cs index 6e7b12103..f7c8e8d4a 100644 --- a/src/Indice.Features.Agents.Core/Services/ChatsService.cs +++ b/src/Indice.Features.Agents.Core/Services/ChatsService.cs @@ -89,7 +89,7 @@ public ChatsService(ISessionsStore store, IDexRunner runner, IUsageGuardService } }, AdditionalProperties = new() { - ["failed"] = result.Failed.ToString().ToLower(), + ["failed"] = result.Failed, ["failureReason"] = result.FailureReason, ["sources"] = result.Sources ?? [], } diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/OutOfScopeResponder.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/OutOfScopeResponder.cs index 3224113e6..c1944b3a4 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/OutOfScopeResponder.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/OutOfScopeResponder.cs @@ -1,5 +1,4 @@ using Indice.Features.Agents.Core.Models; -using Indice.Features.Agents.Core.Workflows; using Indice.Features.Agents.Core.Workflows.State; using Microsoft.Agents.AI.Workflows; From c7b6f65c9edd1e25cfa8a00b3ea3590fd9d71209 Mon Sep 17 00:00:00 2001 From: Constantinos Leftheris Date: Fri, 17 Jul 2026 18:24:07 +0300 Subject: [PATCH 4/5] Refactor pipeline to use record-based state and chat client Refactored codebase to remove custom pipeline envelope and state types in favor of ConversationState, IntentState, and direct payload records. Replaced IDexRunner/DexRunner with IDexChatClient/DexChatClient, aligning with Microsoft.Extensions.AI chat abstractions. Updated all pipeline step executors and workflow context to use new state and payload handling. Revised service registration and ChatsService to use new chat client and response types. Improved maintainability and reduced boilerplate by adopting the new chat pipeline model. --- .../AgentsFeatureExtensions.cs | 6 +- .../Services/ChatsService.cs | 136 +++++-------- .../Workflows/Abstractions/IDexChatClient.cs | 8 + .../Workflows/Abstractions/IDexRunner.cs | 17 -- .../Workflows/DefaultDexPipeline.cs | 6 +- .../Workflows/DexChatClient.cs | 188 ++++++++++++++++++ .../Workflows/DexRunner.cs | 144 -------------- .../Workflows/Events/DexStreamEvent.cs | 4 +- .../Workflows/RagPipelineInput.cs | 9 - .../Workflows/RagRequest.cs | 17 -- .../Workflows/RagResult.cs | 33 --- .../Workflows/RerankOutput.cs | 11 - .../Workflows/RetrievalOutput.cs | 11 - .../Workflows/State/ConversationState.cs | 13 ++ .../Workflows/State/IntentState.cs | 4 + .../Workflows/State/PipelineStepContext.cs | 24 --- .../Workflows/State/RagState.cs | 18 -- .../Workflows/State/StateExtensions.cs | 32 +++ .../Workflows/Steps/AnswerComposer.cs | 17 +- .../Workflows/Steps/IntentClassifier.cs | 20 +- .../Workflows/Steps/OutOfScopeResponder.cs | 12 +- .../Workflows/Steps/PurposeResponder.cs | 16 +- .../Workflows/Steps/QueryRewriter.cs | 20 +- .../Workflows/Steps/Reranker.cs | 25 ++- .../Workflows/Steps/Retriever.cs | 25 ++- 25 files changed, 377 insertions(+), 439 deletions(-) create mode 100644 src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexChatClient.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexRunner.cs create mode 100644 src/Indice.Features.Agents.Core/Workflows/DexChatClient.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/DexRunner.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/RagPipelineInput.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/RagRequest.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/RagResult.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/RerankOutput.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/RetrievalOutput.cs create mode 100644 src/Indice.Features.Agents.Core/Workflows/State/ConversationState.cs create mode 100644 src/Indice.Features.Agents.Core/Workflows/State/IntentState.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/State/PipelineStepContext.cs delete mode 100644 src/Indice.Features.Agents.Core/Workflows/State/RagState.cs create mode 100644 src/Indice.Features.Agents.Core/Workflows/State/StateExtensions.cs diff --git a/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs b/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs index 45e988fb9..94a8bc4de 100644 --- a/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs +++ b/src/Indice.Features.Agents.Core/AgentsFeatureExtensions.cs @@ -22,7 +22,7 @@ public static class AgentsFeatureExtensions /// /// Registers Dex core services: (bound from the Dex configuration section), /// (singleton — each pipeline step builds its own role-bound agent from it), - /// the embedding generator, the wired to SQL Server, and . + /// the embedding generator, the wired to SQL Server, and . /// public static IServiceCollection AddAgentsCore(this IServiceCollection services, IConfiguration configuration, Action? configureAction = null) { var optionsBuilder = services.AddOptions().BindConfiguration("Dex"); @@ -74,12 +74,12 @@ public static IServiceCollection AddAgentsCore(this IServiceCollection services, services.TryAddTransient(); services.TryAddTransient(); services.TryAddSingleton(); - services.TryAddTransient(); + services.TryAddTransient(); services.TryAddSingleton(sp => () => null ); services.TryAddSingleton(); - services.AddDefaultDexPipeline(); + services.AddAgentsDefaultPipeline(); return services; } diff --git a/src/Indice.Features.Agents.Core/Services/ChatsService.cs b/src/Indice.Features.Agents.Core/Services/ChatsService.cs index f7c8e8d4a..63aedf1d5 100644 --- a/src/Indice.Features.Agents.Core/Services/ChatsService.cs +++ b/src/Indice.Features.Agents.Core/Services/ChatsService.cs @@ -1,9 +1,7 @@ using System.Net.ServerSentEvents; using System.Runtime.CompilerServices; using Indice.Features.Agents.Core.Models; -using Indice.Features.Agents.Core.Workflows; using Indice.Features.Agents.Core.Workflows.Abstractions; -using Indice.Features.Agents.Core.Workflows.Events; using Indice.Types; using Microsoft.Extensions.AI; using Microsoft.Extensions.Options; @@ -14,15 +12,15 @@ namespace Indice.Features.Agents.Core.Services; public class ChatsService : IChatsService { private readonly ISessionsStore _store; - private readonly IDexRunner _runner; + private readonly IDexChatClient _dexClient; private readonly IUsageGuardService _usageGuard; private readonly AgentsOptions.AzureOpenAIDeployments _deployments; private readonly AgentsOptions.SessionOptions _sessionOptions; /// Creates a new . - public ChatsService(ISessionsStore store, IDexRunner runner, IUsageGuardService usageGuard, IOptions options) { + public ChatsService(ISessionsStore store, IDexChatClient dexClient, IUsageGuardService usageGuard, IOptions options) { _store = store; - _runner = runner; + _dexClient = dexClient; _usageGuard = usageGuard; _deployments = options.Value.AzureOpenAI.Deployments; _sessionOptions = options.Value.Session; @@ -51,49 +49,21 @@ public ChatsService(ISessionsStore store, IDexRunner runner, IUsageGuardService } }; } - var request = new RagRequest { Question = chatRequest.Text, SessionId = session.Id }; - var result = await _runner.RunAsync(request, cancellationToken); - var userMessage = new ChatMessage(ChatRole.User, request.Question) { + var userMessage = new ChatMessage(ChatRole.User, chatRequest.Text) { MessageId = Guid.NewGuid().ToString(), - CreatedAt = request.TimeStamp + CreatedAt = DateTimeOffset.UtcNow }; - var assistantText = new TextContent(result.Answer ?? string.Empty) { - Annotations = result.Citations?.Select(c => (AIAnnotation)new CitationAnnotation { - FileId = c.DocumentId.ToString(), - Title = result.Sources.First(x => x.Id == c.DocumentId).SourceTitle, - Url = new Uri(result.Sources.First(x => x.Id == c.DocumentId).SourceUrl), - Snippet = c.Title - // Other metadata like page number, confidence, etc. - }).ToList() ?? [] - }; - var assistantMessage = new ChatMessage(ChatRole.Assistant, [assistantText]) { - MessageId = Guid.NewGuid().ToString(), - CreatedAt = DateTimeOffset.UtcNow, - }; - - var persistedAssistant = await _store.AppendTurnAsync(session.Id, userMessage, assistantMessage, + var result = await _dexClient.GetResponseAsync(userMessage, new ChatOptions { ConversationId = session.Id.ToString() }, cancellationToken); + + var persistedAssistant = await _store.AppendTurnAsync(session.Id, userMessage, result.Messages.First(), promptTokens: result.Usage?.InputTokenCount ?? 0, completionTokens: result.Usage?.OutputTokenCount ?? 0, - modelUsed: result.ModelUsed ?? _deployments.Reasoning, cancellationToken); - - return new ChatResponse { - ConversationId = session.Id.ToString(), - ResponseId = persistedAssistant.MessageId, - Messages = [assistantMessage], - ModelId = result.ModelUsed ?? _deployments.Reasoning, - Usage = new() { - InputTokenCount = result.Usage?.InputTokenCount ?? 0, - OutputTokenCount = result.Usage?.OutputTokenCount ?? 0, - AdditionalCounts = new() { - ["questionsUsed"] = _sessionOptions.GetQuestionsUsed(session.MessageCount) ?? 0, - ["questionsTotal"] = _sessionOptions.GetQuestionsTotal() ?? 0 - } - }, - AdditionalProperties = new() { - ["failed"] = result.Failed, - ["failureReason"] = result.FailureReason, - ["sources"] = result.Sources ?? [], - } + modelUsed: result.ModelId ?? _deployments.Reasoning, cancellationToken); + result.Usage ??= new UsageDetails(); + result.Usage.AdditionalCounts = new() { + ["questionsUsed"] = _sessionOptions.GetQuestionsUsed(session.MessageCount) ?? 0, + ["questionsTotal"] = _sessionOptions.GetQuestionsTotal() ?? 0 }; + return result; } /// @@ -140,66 +110,60 @@ private async IAsyncEnumerable> LimitReachedStream(Sess /// Maps the runner's stream to SSE items, then persists the turn and emits the terminal complete event. private async IAsyncEnumerable> StreamTurnAsync( Session session, string text, [EnumeratorCancellation] CancellationToken cancellationToken) { - var request = new RagRequest { Question = text, SessionId = session.Id }; - DexFinalEvent? final = null; - await foreach (var evt in _runner.RunStreamingAsync(request, cancellationToken)) { - switch (evt) { - case DexStepEvent step: - yield return new SseItem(new ChatStreamEvent { Type = "step", Step = step.Label }, eventType: "step"); + + + var userMessage = new ChatMessage(ChatRole.User, text) { + MessageId = Guid.NewGuid().ToString(), + CreatedAt = DateTimeOffset.UtcNow + }; + ChatResponseUpdate? final = null; + UsageContent? usageContent = null; + string? failure = null; + await foreach (var evt in _dexClient.GetStreamingResponseAsync(userMessage, new ChatOptions { ConversationId = session.Id.ToString() }, cancellationToken)) { + switch (evt.RawRepresentation) { + case "Step": + yield return new SseItem(new ChatStreamEvent { Type = "step", Step = evt.Text }, eventType: "step"); + break; + case "Delta": + yield return new SseItem(new ChatStreamEvent { Type = "delta", Text = evt.Text }, eventType: "delta"); + break; + case "Usage": + usageContent = evt.Contents.FirstOrDefault() as UsageContent; break; - case DexDeltaEvent delta: - yield return new SseItem(new ChatStreamEvent { Type = "delta", Text = delta.Text }, eventType: "delta"); + case "Failure": + failure = evt.Text; break; - case DexFinalEvent f: - final = f; + case "Final": + final = evt; break; } } - // The runner yields exactly one terminal DexFinalEvent on success or failure; a mid-stream cancellation - // throws out of the foreach above (the run just stops), so we only reach here on a completed run. - var complete = await PersistTurnAsync(session, request, final, cancellationToken); - yield return new SseItem(complete, eventType: "complete"); - } - - /// Persists the user/assistant turn (mirroring ) and builds the terminal complete event. - private async Task PersistTurnAsync(Session session, RagRequest request, DexFinalEvent? final, CancellationToken cancellationToken) { - var userMessage = new ChatMessage(ChatRole.User, request.Question) { - MessageId = Guid.NewGuid().ToString(), - CreatedAt = request.TimeStamp - }; - var assistantText = new TextContent(final?.Answer ?? string.Empty) { - Annotations = final?.Citations?.Select(c => (AIAnnotation)new CitationAnnotation { - FileId = c.DocumentId.ToString(), - Title = final?.Sources.First(x => x.Id == c.DocumentId).SourceTitle, - Url = new Uri(final!.Sources.First(x => x.Id == c.DocumentId).SourceUrl), - Snippet = c.Title - // Other metadata like page number, confidence, etc. - }).ToList() ?? [] - }; - var assistantMessage = new ChatMessage(ChatRole.Assistant, [assistantText]) { - MessageId = Guid.NewGuid().ToString(), - CreatedAt = DateTimeOffset.UtcNow, + var assistantMessage = new ChatMessage(ChatRole.Assistant, final?.Text ?? "") { + MessageId = final?.MessageId ?? Guid.NewGuid().ToString(), + CreatedAt = DateTimeOffset.UtcNow }; var persistedAssistant = await _store.AppendTurnAsync(session.Id, userMessage, assistantMessage, - promptTokens: final?.Usage?.InputTokenCount ?? 0, completionTokens: final?.Usage?.OutputTokenCount ?? 0, - modelUsed: final?.ModelUsed ?? _deployments.Reasoning, cancellationToken); + promptTokens: usageContent?.Details?.InputTokenCount ?? 0, completionTokens: usageContent?.Details?.OutputTokenCount ?? 0, + modelUsed: final?.ModelId ?? _deployments.Reasoning, cancellationToken); - return new ChatStreamEvent { + var finalEvent = new ChatStreamEvent { Type = "complete", SessionId = session.Id, MessageId = persistedAssistant.MessageId, - Answer = assistantText.Text, - Citations = final?.Citations ?? [], - Sources = final?.Sources ?? [], - Failed = final?.Failed ?? false, - FailureReason = final?.FailureReason, + Answer = assistantMessage.Text, + Citations = /*final?.Citations ??*/ [], + Sources = /*final?.Sources ??*/ [], + Failed = failure != null, + FailureReason = failure, QuestionsUsed = _sessionOptions.GetQuestionsUsed(session.MessageCount + 2), QuestionsTotal = _sessionOptions.GetQuestionsTotal(), }; + yield return new SseItem(finalEvent, eventType: "complete"); } + /// public Task GetAsync(string userId, Guid sessionId, CancellationToken cancellationToken) => _store.GetAsync(sessionId, userId, cancellationToken); diff --git a/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexChatClient.cs b/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexChatClient.cs new file mode 100644 index 000000000..c037452e1 --- /dev/null +++ b/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexChatClient.cs @@ -0,0 +1,8 @@ +using Microsoft.Extensions.AI; + +namespace Indice.Features.Agents.Core.Workflows.Abstractions; + +/// Entry point for executing the Dex RAG pipeline against a single user question. +public interface IDexChatClient : IChatClient +{ +} diff --git a/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexRunner.cs b/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexRunner.cs deleted file mode 100644 index c7f09b483..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/Abstractions/IDexRunner.cs +++ /dev/null @@ -1,17 +0,0 @@ -using Indice.Features.Agents.Core.Workflows.Events; - -namespace Indice.Features.Agents.Core.Workflows.Abstractions; - -/// Entry point for executing the Dex RAG pipeline against a single user question. -public interface IDexRunner -{ - /// Runs the configured pipeline against and projects the final envelope into . - Task RunAsync(RagRequest request, CancellationToken cancellationToken); - - /// - /// Runs the configured pipeline against , yielding real-time - /// s as it executes: a as each step starts, - /// s as the answer streams, and a single terminal . - /// - IAsyncEnumerable RunStreamingAsync(RagRequest request, CancellationToken cancellationToken); -} diff --git a/src/Indice.Features.Agents.Core/Workflows/DefaultDexPipeline.cs b/src/Indice.Features.Agents.Core/Workflows/DefaultDexPipeline.cs index a275336e2..7c1bca5b8 100644 --- a/src/Indice.Features.Agents.Core/Workflows/DefaultDexPipeline.cs +++ b/src/Indice.Features.Agents.Core/Workflows/DefaultDexPipeline.cs @@ -15,7 +15,7 @@ public static class DefaultDexPipelineExtensions /// Registers the five default steps, the default , and a scoped /// wiring them in order. Call after AddDex(...). /// - public static IServiceCollection AddDefaultDexPipeline(this IServiceCollection services) { + public static IServiceCollection AddAgentsDefaultPipeline(this IServiceCollection services) { services.TryAddTransient(); services.TryAddTransient(); services.TryAddTransient(); @@ -38,8 +38,8 @@ public static IServiceCollection AddDefaultDexPipeline(this IServiceCollection s var builder = new WorkflowBuilder(intent); builder.AddSwitch(intent, sw => sw - .AddCase>(env => env!.Payload.Intent.Category == "purpose_of_agent", purposeResponder) - .AddCase>(env => env!.Payload.Intent.IsInScope, rewrite) + .AddCase(env => env!.Intent.Category == "purpose_of_agent", purposeResponder) + .AddCase(env => env!.Intent.IsInScope, rewrite) .WithDefault(outOfScopeReply)); builder.AddEdge(rewrite, retrieve); builder.AddEdge(retrieve, rerank); diff --git a/src/Indice.Features.Agents.Core/Workflows/DexChatClient.cs b/src/Indice.Features.Agents.Core/Workflows/DexChatClient.cs new file mode 100644 index 000000000..c5d99d993 --- /dev/null +++ b/src/Indice.Features.Agents.Core/Workflows/DexChatClient.cs @@ -0,0 +1,188 @@ +using System.Runtime.CompilerServices; +using Indice.Features.Agents.Core.Workflows.Abstractions; +using Indice.Features.Agents.Core.Workflows.Events; +using Indice.Features.Agents.Core.Workflows.State; +using Microsoft.Agents.AI.Workflows; +using Microsoft.Extensions.AI; +using Microsoft.Extensions.DependencyInjection; + +namespace Indice.Features.Agents.Core.Workflows; + +/// +public class DexChatClient : IDexChatClient +{ + private readonly Workflow _workflow; + private readonly IServiceProvider _serviceProvider; + + /// + /// Creates a new instance. + /// + /// The workflow instance to execute. + /// + public DexChatClient ([FromKeyedServices("Default")] Workflow workflow, IServiceProvider serviceProvider) { + _workflow = workflow; + _serviceProvider = serviceProvider; + } + + /// Human-friendly progress labels keyed by executor id, surfaced as SSE step events. + private static readonly IReadOnlyDictionary StepLabels = new Dictionary(StringComparer.Ordinal) { + ["IntentClassifier"] = "Classifying intent", + ["QueryRewriter"] = "Rewriting query", + ["Retriever"] = "Retrieving relevant context", + ["Reranker"] = "Ranking results", + ["AnswerComposer"] = "Composing answer", + ["PurposeResponder"] = "Answering", + ["OutOfScopeResponder"] = "Preparing response", + }; + + /// + public async Task GetResponseAsync(IEnumerable messages, ChatOptions? options = null, CancellationToken cancellationToken = default) { + if (messages.Count() != 1) { + throw new ArgumentException("DexChatClient only supports a single user message per request. No batching allowed.", nameof(messages)); + } + var message = messages.First(); + var state = new ConversationState(message, options?.ConversationId ?? Guid.NewGuid().ToString()); + await using var run = await InProcessExecution.RunAsync(_workflow, state, sessionId: state.ConversationId, cancellationToken: cancellationToken); + RagPipelineOutput? final = null; + string? failure = null; + UsageDetails usage = new UsageDetails(); + string? modelUsed = null; + foreach (var evt in run.NewEvents) { + switch (evt) { + // The terminal executors (compose / out-of-scope) are registered via WithOutputFrom, so their + // returned envelope is yielded as a WorkflowOutputEvent — MAF's dedicated terminal-output channel. + case WorkflowOutputEvent { Data: RagPipelineOutput env }: + final = env; + break; + // Each LLM step reports its own call usage; fold into a single run total. + case UsageEvent usageEvent: + usage.Add(usageEvent.Details); + modelUsed = usageEvent.Model; + break; + // A throwing step halts the run; MAF surfaces the original exception here, followed by a + // WorkflowErrorEvent wrapping it — keep the first (richer) message. + case ExecutorFailedEvent failed: + failure ??= $"{failed.ExecutorId}: {failed.Data?.Message ?? "unknown error"}"; + break; + case WorkflowErrorEvent error: + failure ??= (error.Data as Exception)?.Message ?? "Workflow failed without exception details."; + break; + } + } + // Cancellation never surfaces as a failure event — the run just stops emitting events — so check the + // caller's token explicitly and report a cancellation rather than a pipeline failure. + cancellationToken.ThrowIfCancellationRequested(); + if (final is null && failure is null) { + throw new InvalidOperationException("Workflow completed without emitting a final RagPipelineOutput envelope."); + } + + + var assistantText = new TextContent(final?.Answer ?? string.Empty) { + Annotations = final?.Citations?.Select(c => (AIAnnotation)new CitationAnnotation { + FileId = c.DocumentId.ToString(), + Title = final.Sources.First(x => x.Id == c.DocumentId).SourceTitle, + Url = new Uri(final.Sources.First(x => x.Id == c.DocumentId).SourceUrl), + Snippet = c.Title + // Other metadata like page number, confidence, etc. + }).ToList() ?? [] + }; + var assistantMessage = new ChatMessage(ChatRole.Assistant, [assistantText]) { + MessageId = Guid.NewGuid().ToString(), + CreatedAt = DateTimeOffset.UtcNow, + }; + return new ChatResponse { + ConversationId = state.ConversationId, + ResponseId = assistantMessage.MessageId, + Messages = [assistantMessage], + ModelId = modelUsed, + Usage = new() { + InputTokenCount = usage.InputTokenCount, + OutputTokenCount = usage.OutputTokenCount, + }, + AdditionalProperties = new() { + ["failed"] = failure is not null, + ["failureReason"] = failure, + ["sources"] = final?.Sources ?? [], + } + }; + } + + /// + public async IAsyncEnumerable GetStreamingResponseAsync(IEnumerable messages, ChatOptions? options = null, [EnumeratorCancellation]CancellationToken cancellationToken = default) { + if (messages.Count() != 1) { + throw new ArgumentException("DexChatClient only supports a single user message per request. No batching allowed.", nameof(messages)); + } + var message = messages.First(); + var state = new ConversationState(message, options?.ConversationId ?? Guid.NewGuid().ToString()); + await using var run = await InProcessExecution.RunStreamingAsync(_workflow, state, sessionId: state.ConversationId, cancellationToken: cancellationToken); + RagPipelineOutput? final = null; + string? failure = null; + UsageDetails usage = new UsageDetails(); + string? modelUsed = null; + var messageId = Guid.NewGuid().ToString(); + await foreach (var evt in run.WatchStreamAsync().WithCancellation(cancellationToken)) { + switch (evt) { + // One progress event per step start; unmapped executor ids are skipped. + case ExecutorInvokedEvent invoked when StepLabels.TryGetValue(invoked.ExecutorId, out var label): + yield return new ChatResponseUpdate(ChatRole.Assistant, label) { ConversationId = state.ConversationId, MessageId = messageId, RawRepresentation = "Step" } ; + break; + // Each LLM step reports its own call usage; fold into a single run total. + case UsageEvent usageEvent: + usage.Add(usageEvent.Details); + modelUsed = usageEvent.Model; + yield return new ChatResponseUpdate(ChatRole.Assistant, [new UsageContent(usage)]) { ConversationId = state.ConversationId, MessageId = messageId, RawRepresentation = "Usage" }; + break; + // Answer text deltas emitted by AnswerComposer as the reasoning model streams. + case AnswerDeltaEvent delta when delta.Delta.Length > 0: + yield return new ChatResponseUpdate(ChatRole.Assistant, delta.Delta) { ConversationId = state.ConversationId, MessageId = messageId, RawRepresentation = "Delta" }; + break; + // Terminal output from compose / out-of-scope (registered via WithOutputFrom). + case WorkflowOutputEvent { Data: RagPipelineOutput env }: + final = env; + break; + // A throwing step halts the run; keep the first (richer) message. + case ExecutorFailedEvent failed: + failure ??= $"{failed.ExecutorId}: {failed.Data?.Message ?? "unknown error"}"; + yield return new ChatResponseUpdate(ChatRole.Assistant, failure) { ConversationId = state.ConversationId, MessageId = messageId, RawRepresentation = "Failure" }; + break; + case WorkflowErrorEvent error: + failure ??= (error.Data as Exception)?.Message ?? "Workflow failed without exception details."; + yield return new ChatResponseUpdate(ChatRole.Assistant, failure) { ConversationId = state.ConversationId, MessageId = messageId, RawRepresentation = "Error" }; + break; + } + } + // Cancellation just stops the stream rather than raising a failure event — surface it as cancellation. + cancellationToken.ThrowIfCancellationRequested(); + + + yield return new ChatResponseUpdate(ChatRole.Assistant, final?.Answer ?? "") { + RawRepresentation = new DexFinalEvent { + Answer = final?.Answer, + Citations = final?.Citations ?? [], + Sources = final?.Sources ?? [], + Failed = failure is not null, + FailureReason = failure, + Usage = usage, + ModelUsed = modelUsed, + }, + ConversationId = state.ConversationId, + ResponseId = messageId, + MessageId = messageId, + ModelId = modelUsed, + AdditionalProperties = new() { + ["failed"] = failure is not null, + ["failureReason"] = failure, + ["sources"] = final?.Sources ?? [], + } + }; + + } + + /// + public object? GetService(Type serviceType, object? serviceKey = null) => serviceKey is null ? _serviceProvider.GetService(serviceType) : _serviceProvider.GetKeyedService(serviceType, serviceKey); + + /// + public void Dispose() { + + } +} diff --git a/src/Indice.Features.Agents.Core/Workflows/DexRunner.cs b/src/Indice.Features.Agents.Core/Workflows/DexRunner.cs deleted file mode 100644 index bf156cb08..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/DexRunner.cs +++ /dev/null @@ -1,144 +0,0 @@ -using System.Runtime.CompilerServices; -using Microsoft.Agents.AI.Workflows; -using Microsoft.Extensions.AI; -using Microsoft.Extensions.DependencyInjection; -using Indice.Features.Agents.Core.Workflows.State; -using Indice.Features.Agents.Core.Workflows.Events; -using Indice.Features.Agents.Core.Workflows.Abstractions; - -namespace Indice.Features.Agents.Core.Workflows; - -/// -public class DexRunner : IDexRunner -{ - private readonly Workflow _workflow; - - /// - /// Creates a new instance. - /// - /// The workflow instance to execute. - public DexRunner ([FromKeyedServices("Default")] Workflow workflow) { - _workflow = workflow; - } - - /// Human-friendly progress labels keyed by executor id, surfaced as SSE step events. - private static readonly IReadOnlyDictionary StepLabels = new Dictionary(StringComparer.Ordinal) { - ["IntentClassifier"] = "Classifying intent", - ["QueryRewriter"] = "Rewriting query", - ["Retriever"] = "Retrieving relevant context", - ["Reranker"] = "Ranking results", - ["AnswerComposer"] = "Composing answer", - ["PurposeResponder"] = "Answering", - ["OutOfScopeResponder"] = "Preparing response", - }; - - /// - public async Task RunAsync(RagRequest request, CancellationToken cancellationToken) { - var initial = CreateInitialEnvelope(request); - await using var run = await InProcessExecution.RunAsync(_workflow, initial, cancellationToken: cancellationToken); - PipelineStepContext? final = null; - string? failure = null; - UsageDetails? usage = null; - string? modelUsed = null; - foreach (var evt in run.NewEvents) { - switch (evt) { - // The terminal executors (compose / out-of-scope) are registered via WithOutputFrom, so their - // returned envelope is yielded as a WorkflowOutputEvent — MAF's dedicated terminal-output channel. - case WorkflowOutputEvent { Data: PipelineStepContext env }: - final = env; - break; - // Each LLM step reports its own call usage; fold into a single run total. - case UsageEvent usageEvent: - (usage ??= new UsageDetails()).Add(usageEvent.Details); - modelUsed = usageEvent.Model; - break; - // A throwing step halts the run; MAF surfaces the original exception here, followed by a - // WorkflowErrorEvent wrapping it — keep the first (richer) message. - case ExecutorFailedEvent failed: - failure ??= $"{failed.ExecutorId}: {failed.Data?.Message ?? "unknown error"}"; - break; - case WorkflowErrorEvent error: - failure ??= (error.Data as Exception)?.Message ?? "Workflow failed without exception details."; - break; - } - } - // Cancellation never surfaces as a failure event — the run just stops emitting events — so check the - // caller's token explicitly and report a cancellation rather than a pipeline failure. - cancellationToken.ThrowIfCancellationRequested(); - if (final is null && failure is null) { - throw new InvalidOperationException("Workflow completed without emitting a final RagPipelineOutput envelope."); - } - return new RagResult { - Answer = final?.Payload?.Answer, - Citations = final?.Payload?.Citations ?? [], - Sources = final?.Payload?.Sources ?? [], - Failed = failure is not null, - FailureReason = failure, - Usage = usage, - ModelUsed = modelUsed, - }; - } - - /// - public async IAsyncEnumerable RunStreamingAsync( - RagRequest request, [EnumeratorCancellation] CancellationToken cancellationToken) { - var initial = CreateInitialEnvelope(request); - await using var run = await InProcessExecution.RunStreamingAsync(_workflow!, initial, cancellationToken: cancellationToken); - PipelineStepContext? final = null; - string? failure = null; - UsageDetails? usage = null; - string? modelUsed = null; - await foreach (var evt in run.WatchStreamAsync().WithCancellation(cancellationToken)) { - switch (evt) { - // One progress event per step start; unmapped executor ids are skipped. - case ExecutorInvokedEvent invoked when StepLabels.TryGetValue(invoked.ExecutorId, out var label): - yield return new DexStepEvent(invoked.ExecutorId, label); - break; - // Each LLM step reports its own call usage; fold into a single run total. - case UsageEvent usageEvent: - (usage ??= new UsageDetails()).Add(usageEvent.Details); - modelUsed = usageEvent.Model; - break; - // Answer text deltas emitted by AnswerComposer as the reasoning model streams. - case AnswerDeltaEvent delta when delta.Delta.Length > 0: - yield return new DexDeltaEvent(delta.Delta); - break; - // Terminal output from compose / out-of-scope (registered via WithOutputFrom). - case WorkflowOutputEvent { Data: PipelineStepContext env }: - final = env; - break; - // A throwing step halts the run; keep the first (richer) message. - case ExecutorFailedEvent failed: - failure ??= $"{failed.ExecutorId}: {failed.Data?.Message ?? "unknown error"}"; - break; - case WorkflowErrorEvent error: - failure ??= (error.Data as Exception)?.Message ?? "Workflow failed without exception details."; - break; - } - } - // Cancellation just stops the stream rather than raising a failure event — surface it as cancellation. - cancellationToken.ThrowIfCancellationRequested(); - yield return new DexFinalEvent { - Answer = final?.Payload?.Answer, - Citations = final?.Payload?.Citations ?? [], - Sources = final?.Payload?.Sources ?? [], - Failed = failure is not null, - FailureReason = failure, - Usage = usage, - ModelUsed = modelUsed, - }; - } - - /// Validates a workflow is registered and builds the initial pipeline envelope from . - private PipelineStepContext CreateInitialEnvelope(RagRequest request) { - if (_workflow is null) { - throw new InvalidOperationException( - "No RAG workflow registered. Call services.AddDefaultDexPipeline() or register a Microsoft.Agents.AI.Workflows.Workflow manually."); - } - var initialState = new RagState { - Question = request.Question, - SessionId = request.SessionId, - }; - return PipelineStepContext.From(new RagPipelineInput(), initialState); - } -} diff --git a/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs b/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs index eb958e7bc..adc31d6c4 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs @@ -4,7 +4,7 @@ namespace Indice.Features.Agents.Core.Workflows.Events; /// -/// Base type for the real-time events yields as +/// Base type for the real-time events yields as /// the pipeline executes: per-step progress (), answer text deltas /// (), and a single terminal . /// @@ -56,7 +56,7 @@ public DexErrorEvent(string stepId, string errorMessage) { } /// -/// Terminal event yielded once after the run completes. Mirrors the fields of : +/// Terminal event yielded once after the run completes. /// the full answer, citations, failure state, and reasoning-model token totals. /// public sealed class DexFinalEvent : DexStreamEvent diff --git a/src/Indice.Features.Agents.Core/Workflows/RagPipelineInput.cs b/src/Indice.Features.Agents.Core/Workflows/RagPipelineInput.cs deleted file mode 100644 index e1a34f0cc..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/RagPipelineInput.cs +++ /dev/null @@ -1,9 +0,0 @@ -namespace Indice.Features.Agents.Core.Workflows; - -/// -/// Marker payload for the first edge of a Dex RAG pipeline. The question itself lives on -/// RagState.Question (seeded by DexRunner), so the initial payload carries no data. -/// -public class RagPipelineInput -{ -} diff --git a/src/Indice.Features.Agents.Core/Workflows/RagRequest.cs b/src/Indice.Features.Agents.Core/Workflows/RagRequest.cs deleted file mode 100644 index c377bcc1b..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/RagRequest.cs +++ /dev/null @@ -1,17 +0,0 @@ -namespace Indice.Features.Agents.Core.Workflows; - -/// The input to . -public class RagRequest -{ - /// The end-user question being asked. - public string Question { get; init; } = string.Empty; - - /// Timestamp of when the request was created. - public DateTimeOffset TimeStamp { get; init; } = DateTimeOffset.UtcNow; - - /// - /// The chat session this question belongs to. The pipeline's - /// loads the windowed conversation history for it during the run. - /// - public Guid SessionId { get; init; } -} diff --git a/src/Indice.Features.Agents.Core/Workflows/RagResult.cs b/src/Indice.Features.Agents.Core/Workflows/RagResult.cs deleted file mode 100644 index 4f55a5535..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/RagResult.cs +++ /dev/null @@ -1,33 +0,0 @@ -using Indice.Features.Agents.Core.Models; -using Indice.Features.Agents.Core.Workflows.Abstractions; -using Microsoft.Extensions.AI; - -namespace Indice.Features.Agents.Core.Workflows; - -/// -/// The output of . Projected by DexRunner from the final -/// envelope's payload, the workflow's failure events, and the accumulated token usage. -/// -public class RagResult -{ - /// The pipeline's answer — grounded when intent was in-scope; the polite refusal text from OutOfScopeResponder when not; null only when a step threw (). - public string? Answer { get; init; } - - /// Citations accumulated across retrieval/rerank/compose, surfaced from the final payload. - public IReadOnlyList Citations { get; init; } = []; - - /// Links to the source documents that were retrieved and used to compose the answer, surfaced from the final payload. - public IReadOnlyList Sources { get; init; } = []; - - /// True when a step threw and the workflow halted (surfaced via MAF's ExecutorFailedEvent). Out-of-scope is NOT a failure — it flows through OutOfScopeResponder and produces a regular . - public bool Failed { get; init; } - - /// Error message from the step that threw, prefixed with its executor id; null when is false. - public string? FailureReason { get; init; } - - /// Total reasoning-model token usage across this run, folded from the steps' UsageEvents. Persisted to the session, not returned to the caller; null when no reasoning call ran. - public UsageDetails? Usage { get; init; } - - /// The reasoning-model deployment the tokens were billed against; null when no reasoning call ran. - public string? ModelUsed { get; init; } -} diff --git a/src/Indice.Features.Agents.Core/Workflows/RerankOutput.cs b/src/Indice.Features.Agents.Core/Workflows/RerankOutput.cs deleted file mode 100644 index 9fa11fbb1..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/RerankOutput.cs +++ /dev/null @@ -1,11 +0,0 @@ -namespace Indice.Features.Agents.Core.Workflows; - -/// Output payload of Reranker. -public class RerankOutput -{ - /// The classified intent, forwarded from upstream. - public Intent Intent { get; init; } = new(); - - /// Top-N candidates reordered by reranker score; their Score reflects the rerank outcome. - public IReadOnlyList RerankedCandidates { get; init; } = Array.Empty(); -} diff --git a/src/Indice.Features.Agents.Core/Workflows/RetrievalOutput.cs b/src/Indice.Features.Agents.Core/Workflows/RetrievalOutput.cs deleted file mode 100644 index 41467526f..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/RetrievalOutput.cs +++ /dev/null @@ -1,11 +0,0 @@ -namespace Indice.Features.Agents.Core.Workflows; - -/// Output payload of Retriever. -public class RetrievalOutput -{ - /// The classified intent, forwarded from upstream. - public Intent Intent { get; init; } = new(); - - /// Top candidates union'd and deduped across all rewritten queries, ordered by cosine similarity. - public IReadOnlyList Candidates { get; init; } = Array.Empty(); -} diff --git a/src/Indice.Features.Agents.Core/Workflows/State/ConversationState.cs b/src/Indice.Features.Agents.Core/Workflows/State/ConversationState.cs new file mode 100644 index 000000000..89101d51d --- /dev/null +++ b/src/Indice.Features.Agents.Core/Workflows/State/ConversationState.cs @@ -0,0 +1,13 @@ +using Microsoft.Agents.AI.Workflows; +using Microsoft.Extensions.AI; + +namespace Indice.Features.Agents.Core.Workflows.State; + +/// +/// Immutable read-only context carried alongside the typed payload through every pipeline edge. +/// Seeded once by DexRunner from the incoming RagRequest; steps forward it untouched +/// via PipelineEnvelope.Next — never mutate. +/// +/// The chat session this run belongs to. History-aware steps stamp it on their per-run +/// The chat message being processed. +public record ConversationState(ChatMessage Message, string ConversationId); \ No newline at end of file diff --git a/src/Indice.Features.Agents.Core/Workflows/State/IntentState.cs b/src/Indice.Features.Agents.Core/Workflows/State/IntentState.cs new file mode 100644 index 000000000..cb09f143f --- /dev/null +++ b/src/Indice.Features.Agents.Core/Workflows/State/IntentState.cs @@ -0,0 +1,4 @@ +namespace Indice.Features.Agents.Core.Workflows.State; + +/// Immutable read-only context carried alongside the typed payload through every pipeline edge. +public record IntentState(Intent Intent, RetrievalFilters Filters); diff --git a/src/Indice.Features.Agents.Core/Workflows/State/PipelineStepContext.cs b/src/Indice.Features.Agents.Core/Workflows/State/PipelineStepContext.cs deleted file mode 100644 index b462019a1..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/State/PipelineStepContext.cs +++ /dev/null @@ -1,24 +0,0 @@ -namespace Indice.Features.Agents.Core.Workflows.State; - -/// -/// The message that flows along every pipeline edge. Pairs a typed with the -/// immutable cross-cutting . Steps construct new envelopes via ; -/// they never mutate the input. -/// -/// The payload type for this edge (changes from step to step). -public class PipelineStepContext -{ - /// The typed payload carried by this edge. - public TPayload Payload { get; init; } = default!; - - /// The cross-cutting state accumulated so far (immutable). - public RagState State { get; init; } = new(); - - /// Creates an envelope from a payload (and optionally a starting state). - public static PipelineStepContext From(TPayload payload, RagState? state = null) - => new() { Payload = payload, State = state ?? new RagState() }; - - /// Creates a downstream envelope carrying a new payload and the current (or replaced) state. - public PipelineStepContext Next(TNext payload, RagState? state = null) - => new() { Payload = payload, State = state ?? State }; -} diff --git a/src/Indice.Features.Agents.Core/Workflows/State/RagState.cs b/src/Indice.Features.Agents.Core/Workflows/State/RagState.cs deleted file mode 100644 index 6700f74dd..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/State/RagState.cs +++ /dev/null @@ -1,18 +0,0 @@ -namespace Indice.Features.Agents.Core.Workflows.State; - -/// -/// Immutable read-only context carried alongside the typed payload through every pipeline edge. -/// Seeded once by DexRunner from the incoming RagRequest; steps forward it untouched -/// via PipelineEnvelope.Next — never mutate. -/// -public class RagState -{ - /// The current user question that initiated this pipeline run. - public string Question { get; init; } = string.Empty; - - /// - /// The chat session this run belongs to. History-aware steps stamp it on their per-run - /// AgentSession so the can load the conversation. - /// - public Guid SessionId { get; init; } -} diff --git a/src/Indice.Features.Agents.Core/Workflows/State/StateExtensions.cs b/src/Indice.Features.Agents.Core/Workflows/State/StateExtensions.cs new file mode 100644 index 000000000..ed35c5819 --- /dev/null +++ b/src/Indice.Features.Agents.Core/Workflows/State/StateExtensions.cs @@ -0,0 +1,32 @@ +using Microsoft.Agents.AI.Workflows; + +namespace Indice.Features.Agents.Core.Workflows.State; + +/// Extension methods for . +public static class IWorkflowContextStateExtensions +{ + /// The scope name used to store the in the workflow context. + public const string ConversationScope = "ConversationScope"; + /// Reads the from the workflow context. + public static async Task GetConversationStateAsync(this IWorkflowContext context, CancellationToken cancellationToken = default) { + return await context.ReadStateAsync(nameof(ConversationState), scopeName: ConversationScope, cancellationToken: cancellationToken) ?? + throw new InvalidOperationException("ConversationState not found in workflow context."); + } + + /// Writes the to the workflow context. + public static async Task SetConversationStateAsync(this IWorkflowContext context, ConversationState state, CancellationToken cancellationToken = default) { + await context.QueueStateUpdateAsync(nameof(ConversationState), state, scopeName: ConversationScope, cancellationToken: cancellationToken); + } + + /// Reads the from the workflow context. + public static async Task GetIntentStateAsync(this IWorkflowContext context, CancellationToken cancellationToken = default) { + return await context.ReadStateAsync(nameof(IntentState), scopeName: ConversationScope, cancellationToken: cancellationToken) ?? + throw new InvalidOperationException("IntentState not found in workflow context."); + } + + /// Writes the to the workflow context. + public static async Task SetIntentStateAsync(this IWorkflowContext context, IntentState state, CancellationToken cancellationToken = default) { + await context.QueueStateUpdateAsync(nameof(IntentState), state, scopeName: ConversationScope, cancellationToken: cancellationToken); + } +} + diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/AnswerComposer.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/AnswerComposer.cs index e4d237f58..b851e9ddc 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/AnswerComposer.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/AnswerComposer.cs @@ -17,7 +17,7 @@ namespace Indice.Features.Agents.Core.Workflows.Steps; /// in the provided context and cite chunk IDs in [#chunkId] form. Projects the candidates into /// records on the output payload. /// -public sealed class AnswerComposer : Executor, PipelineStepContext> +public sealed class AnswerComposer : Executor { private readonly AIAgent _agent; private readonly AgentsOptions _options; @@ -35,7 +35,6 @@ public AnswerComposer([FromKeyedServices(nameof(AzureOpenAIDeployments.Reasoning chatOptions.Instructions = prompts.Render("AnswerComposer", new { strictGrounding = _options.Pipeline.StrictGrounding, }); - _agent = chatClient.AsAIAgent( options: new ChatClientAgentOptions() { ChatOptions = chatOptions, @@ -46,13 +45,13 @@ public AnswerComposer([FromKeyedServices(nameof(AzureOpenAIDeployments.Reasoning } /// - public override async ValueTask> HandleAsync(PipelineStepContext envelope, + public override async ValueTask HandleAsync(RerankOutput message, IWorkflowContext context, CancellationToken cancellationToken = default) { - - var candidates = envelope.Payload.RerankedCandidates; - var prompt = BuildPrompt(envelope.State.Question, candidates); + var state = await context.GetConversationStateAsync(cancellationToken); + var candidates = message.RerankedCandidates; + var prompt = BuildPrompt(state.Message.Text, candidates); var agentSession = await _agent.CreateSessionAsync(cancellationToken); - SessionStoreChatHistoryProvider.SetSessionId(agentSession, envelope.State.SessionId); + SessionStoreChatHistoryProvider.SetSessionId(agentSession, Guid.Parse(state.ConversationId)); // Stream the answer: emit each text delta as a workflow event (surfaced as an SSE `delta` by the // streaming runner; ignored by the non-streaming runner) while accumulating the full text. Token @@ -83,11 +82,11 @@ public override async ValueTask> HandleAs }) .ToList(); var sources = candidates.Select(c => c.Source).DistinctBy(x => x.Id).ToList(); - return envelope.Next(new RagPipelineOutput { + return new RagPipelineOutput { Answer = answer.ToString(), Citations = citations, Sources = sources, - }); + }; } private static string BuildPrompt(string question, IReadOnlyList candidates) { diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/IntentClassifier.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/IntentClassifier.cs index c40c191d8..0cb6449f5 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/IntentClassifier.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/IntentClassifier.cs @@ -6,7 +6,6 @@ using Microsoft.Extensions.AI; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; -using OpenAI.Chat; using static Indice.Features.Agents.Core.AgentsOptions; namespace Indice.Features.Agents.Core.Workflows.Steps; @@ -16,7 +15,7 @@ namespace Indice.Features.Agents.Core.Workflows.Steps; /// On out-of-scope the workflow routes to via a conditional edge; otherwise /// downstream steps receive validated and . /// -public sealed class IntentClassifier : Executor, PipelineStepContext> +public sealed class IntentClassifier : Executor { private readonly AIAgent _agent; private readonly AgentsOptions _options; @@ -41,13 +40,16 @@ public IntentClassifier([FromKeyedServices(nameof(AzureOpenAIDeployments.Reasoni } /// - public override async ValueTask> HandleAsync( - PipelineStepContext envelope, + public override async ValueTask HandleAsync( + ConversationState message, IWorkflowContext context, CancellationToken cancellationToken = default) { - var question = envelope.State.Question; + var question = message.Message.Text; + var conversationId = message.ConversationId; + await context.SetConversationStateAsync(message, cancellationToken); var agentSession = await _agent.CreateSessionAsync(cancellationToken); - SessionStoreChatHistoryProvider.SetSessionId(agentSession, envelope.State.SessionId); + + SessionStoreChatHistoryProvider.SetSessionId(agentSession, Guid.Parse(conversationId)); var response = await _agent.RunAsync(question, agentSession, cancellationToken: cancellationToken); if (response.Usage is not null) { await context.AddEventAsync(new UsageEvent(response.Usage, _model), cancellationToken); @@ -64,11 +66,11 @@ public override async ValueTask> HandleAsync( IsInScope = result.IsInScope, OutOfScopeReason = result.OutOfScopeReason, }; - - return envelope.Next(new IntentOutput { + await context.SetIntentStateAsync(new IntentState(intent, new RetrievalFilters { Category = category, Language = language }), cancellationToken); + return new IntentOutput { Intent = intent, Filters = new RetrievalFilters { Category = category, Language = language }, - }); + }; } private sealed class IntentResult diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/OutOfScopeResponder.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/OutOfScopeResponder.cs index c1944b3a4..9cde2a1d7 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/OutOfScopeResponder.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/OutOfScopeResponder.cs @@ -9,20 +9,20 @@ namespace Indice.Features.Agents.Core.Workflows.Steps; /// Projects the classifier's into a final /// envelope with no citations. /// -public sealed class OutOfScopeResponder : Executor, PipelineStepContext> +public sealed class OutOfScopeResponder : Executor { /// Creates a new . public OutOfScopeResponder() : base("OutOfScopeResponder") { } /// - public override ValueTask> HandleAsync( - PipelineStepContext envelope, + public override ValueTask HandleAsync( + IntentOutput intentResult, IWorkflowContext context, CancellationToken cancellationToken = default) { - var reason = envelope.Payload.Intent.OutOfScopeReason ?? "Sorry, that question is outside the scope of what I can answer here."; - return ValueTask.FromResult(envelope.Next(new RagPipelineOutput { + var reason = intentResult.Intent.OutOfScopeReason ?? "Sorry, that question is outside the scope of what I can answer here."; + return ValueTask.FromResult(new RagPipelineOutput { Answer = reason, Citations = Array.Empty(), - })); + }); } } diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/PurposeResponder.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/PurposeResponder.cs index c82a0ccaa..3f0abebab 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/PurposeResponder.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/PurposeResponder.cs @@ -16,7 +16,7 @@ namespace Indice.Features.Agents.Core.Workflows.Steps; /// Terminal branch of the pipeline when IntentClassifier decides the \ /// question is a general question about the capabilities of the agent. /// -internal class PurposeResponder : Executor, PipelineStepContext> +internal class PurposeResponder : Executor { private readonly AIAgent _agent; private readonly AgentsOptions _options; @@ -47,13 +47,13 @@ public PurposeResponder([FromKeyedServices(nameof(AzureOpenAIDeployments.Reasoni } /// - public override async ValueTask> HandleAsync( - PipelineStepContext envelope,IWorkflowContext context, + public override async ValueTask HandleAsync( + IntentOutput intentResult, IWorkflowContext context, CancellationToken cancellationToken = default) { - - var prompt = envelope.State.Question; + var state = await context.GetConversationStateAsync(cancellationToken); + var prompt = state.Message.Text; var agentSession = await _agent.CreateSessionAsync(cancellationToken); - SessionStoreChatHistoryProvider.SetSessionId(agentSession, envelope.State.SessionId); + SessionStoreChatHistoryProvider.SetSessionId(agentSession, Guid.Parse(state.ConversationId)); // Stream the answer: emit each text delta as a workflow event (surfaced as an SSE `delta` by the // streaming runner; ignored by the non-streaming runner) while accumulating the full text. Token @@ -73,9 +73,9 @@ public override async ValueTask> HandleAs await context.AddEventAsync(new UsageEvent(usage, _model), cancellationToken); } - return envelope.Next(new RagPipelineOutput { + return new RagPipelineOutput { Answer = answer.ToString() - }); + }; } } diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/QueryRewriter.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/QueryRewriter.cs index d76c7ae61..82db34c98 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/QueryRewriter.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/QueryRewriter.cs @@ -15,7 +15,7 @@ namespace Indice.Features.Agents.Core.Workflows.Steps; /// broaden retrieval recall. Disabled by DefaultPipelineOptions.EnableQueryRewrite = false; on any LLM /// failure, falls back to the original question. /// -public sealed class QueryRewriter : Executor, PipelineStepContext> +public sealed class QueryRewriter : Executor { private readonly AIAgent _agent; private readonly AgentsOptions _options; @@ -38,17 +38,17 @@ public QueryRewriter([FromKeyedServices(nameof(AzureOpenAIDeployments.Fast))] IC } /// - public override async ValueTask> HandleAsync(PipelineStepContext envelope, + public override async ValueTask HandleAsync(IntentOutput intentResult, IWorkflowContext context, CancellationToken cancellationToken = default) { - var question = envelope.State.Question; + var state = await context.GetConversationStateAsync(cancellationToken); var expansion = _options.Retrieval.QueryExpansion; var enabled = _options.Pipeline.EnableQueryRewrite && expansion > 1; - var queries = new List { question }; + var queries = new List { state.Message.Text }; if (enabled) { var agentSession = await _agent.CreateSessionAsync(cancellationToken); - SessionStoreChatHistoryProvider.SetSessionId(agentSession, envelope.State.SessionId); - var prompt = $"Question: {question}\n\nProduce {expansion - 1} alternative rewrite(s)."; + SessionStoreChatHistoryProvider.SetSessionId(agentSession, Guid.Parse(state.ConversationId)); + var prompt = $"Question: {state.Message.Text}\n\nProduce {expansion - 1} alternative rewrite(s)."; var response = await _agent.RunAsync(prompt, agentSession, cancellationToken: cancellationToken); foreach (var q in response.Result.Queries) { if (!string.IsNullOrWhiteSpace(q) && !queries.Contains(q, StringComparer.OrdinalIgnoreCase)) { @@ -57,11 +57,11 @@ public override async ValueTask> HandleA if (queries.Count >= expansion) break; } } - return envelope.Next(new QueryRewriteOutput { - Intent = envelope.Payload.Intent, - Filters = envelope.Payload.Filters, + return new QueryRewriteOutput { + Intent = intentResult.Intent, + Filters = intentResult.Filters, RewrittenQueries = queries, - }); + }; } private sealed class RewriteResult diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/Reranker.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/Reranker.cs index e0ba2dcc2..cb9419bae 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/Reranker.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/Reranker.cs @@ -10,7 +10,7 @@ namespace Indice.Features.Agents.Core.Workflows.Steps; /// the registered ; bypasses the LLM when reranking is disabled or already /// at/below the target size. /// -public sealed class Reranker : Executor, PipelineStepContext> +public sealed class Reranker : Executor { private readonly ILlmReranker _reranker; private readonly AgentsOptions _options; @@ -22,18 +22,23 @@ public Reranker(ILlmReranker reranker, IOptions options) : base(" } /// - public override async ValueTask> HandleAsync(PipelineStepContext envelope, IWorkflowContext context, + public override async ValueTask HandleAsync(RetrievalOutput retrievalOutput, IWorkflowContext context, CancellationToken cancellationToken = default) { var topResults = _options.Retrieval.NumberOfResults; - var candidates = envelope.Payload.Candidates; - + var candidates = retrievalOutput.Candidates; + var state = await context.GetConversationStateAsync(cancellationToken); + var intentState = await context.GetIntentStateAsync(cancellationToken); IReadOnlyList reranked = !_options.Pipeline.EnableRerank || candidates.Count <= topResults ? candidates.OrderByDescending(c => c.Score).Take(topResults).ToList() - : await _reranker.RerankAsync(envelope.State.Question, candidates, topResults, cancellationToken); - - return envelope.Next(new RerankOutput { - Intent = envelope.Payload.Intent, - RerankedCandidates = reranked, - }); + : await _reranker.RerankAsync(state.Message.Text, candidates, topResults, cancellationToken); + return new RerankOutput( + intentState.Intent, + reranked + ); } } + +/// Output payload of Reranker. +/// The classified intent, forwarded from upstream. +/// Top-N candidates reordered by reranker score; their Score reflects the rerank outcome. +public record RerankOutput(Intent Intent, IReadOnlyList RerankedCandidates); diff --git a/src/Indice.Features.Agents.Core/Workflows/Steps/Retriever.cs b/src/Indice.Features.Agents.Core/Workflows/Steps/Retriever.cs index acdb579db..024125fc2 100644 --- a/src/Indice.Features.Agents.Core/Workflows/Steps/Retriever.cs +++ b/src/Indice.Features.Agents.Core/Workflows/Steps/Retriever.cs @@ -10,7 +10,7 @@ namespace Indice.Features.Agents.Core.Workflows.Steps; /// Embeds each rewritten query and retrieves the top-K most cosine-similar chunks via . /// Unions and deduplicates across rewrites by ChunkId, keeping the highest score. /// -public sealed class Retriever : Executor, PipelineStepContext> +public sealed class Retriever : Executor { private readonly IDocumentsService _documentsService; private readonly IEmbeddingGenerator> _embedder; @@ -27,13 +27,13 @@ public Retriever( } /// - public override async ValueTask> HandleAsync( - PipelineStepContext envelope, IWorkflowContext context, CancellationToken cancellationToken = default) { - var filters = envelope.Payload.Filters; + public override async ValueTask HandleAsync( + QueryRewriteOutput queryRewriteOutput, IWorkflowContext context, CancellationToken cancellationToken = default) { + var filters = queryRewriteOutput.Filters; var topK = _options.Retrieval.NumberOfCandidates; var relevantAnswers = new Dictionary(); - foreach (var query in envelope.Payload.RewrittenQueries) { + foreach (var query in queryRewriteOutput.RewrittenQueries) { // Embedding dimensions are configured once on the generator registration (AddAgentsCore). var vector = await _embedder.GenerateVectorAsync(query, cancellationToken: cancellationToken); var hits = await _documentsService.SearchAsync(vector, filters, topK, _options.Retrieval.MinScore, cancellationToken); @@ -46,9 +46,16 @@ public override async ValueTask> HandleAsyn var candidates = relevantAnswers.Values .OrderByDescending(c => c.Score) .ToList(); - return envelope.Next(new RetrievalOutput { - Intent = envelope.Payload.Intent, - Candidates = candidates, - }); + return new RetrievalOutput( + queryRewriteOutput.Intent, + candidates + ); } } + + +/// Output payload of Retriever. +/// The classified intent, forwarded from upstream. +/// Top candidates union'd and deduped across all rewritten queries, ordered by cosine similarity. +public record RetrievalOutput(Intent Intent, IReadOnlyList Candidates); + From 7091e123a80e438e9ecae7eb9c962c4c321d463b Mon Sep 17 00:00:00 2001 From: Constantinos Leftheris Date: Fri, 17 Jul 2026 18:30:39 +0300 Subject: [PATCH 5/5] Refactor chat event citation/source handling Citations and sources are now extracted from `final?.AdditionalProperties` in `ChatsService.cs` instead of defaulting to empty arrays. In `DexChatClient.cs`, the final chat response's `RawRepresentation` is simplified, and citations/sources are added to `AdditionalProperties`. All event class definitions are removed from `DexStreamEvent.cs`, reflecting a shift to a new event handling mechanism. --- .../Services/ChatsService.cs | 4 +- .../Workflows/DexChatClient.cs | 12 +-- .../Workflows/Events/DexStreamEvent.cs | 84 ------------------- 3 files changed, 5 insertions(+), 95 deletions(-) delete mode 100644 src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs diff --git a/src/Indice.Features.Agents.Core/Services/ChatsService.cs b/src/Indice.Features.Agents.Core/Services/ChatsService.cs index 63aedf1d5..f0426148c 100644 --- a/src/Indice.Features.Agents.Core/Services/ChatsService.cs +++ b/src/Indice.Features.Agents.Core/Services/ChatsService.cs @@ -153,8 +153,8 @@ private async IAsyncEnumerable> StreamTurnAsync( SessionId = session.Id, MessageId = persistedAssistant.MessageId, Answer = assistantMessage.Text, - Citations = /*final?.Citations ??*/ [], - Sources = /*final?.Sources ??*/ [], + Citations = final?.AdditionalProperties?["citations"] as IReadOnlyList ?? [], + Sources = final?.AdditionalProperties?["sources"] as IReadOnlyList ?? [], Failed = failure != null, FailureReason = failure, QuestionsUsed = _sessionOptions.GetQuestionsUsed(session.MessageCount + 2), diff --git a/src/Indice.Features.Agents.Core/Workflows/DexChatClient.cs b/src/Indice.Features.Agents.Core/Workflows/DexChatClient.cs index c5d99d993..7940cf803 100644 --- a/src/Indice.Features.Agents.Core/Workflows/DexChatClient.cs +++ b/src/Indice.Features.Agents.Core/Workflows/DexChatClient.cs @@ -1,4 +1,5 @@ using System.Runtime.CompilerServices; +using Indice.Features.Agents.Core.Models; using Indice.Features.Agents.Core.Workflows.Abstractions; using Indice.Features.Agents.Core.Workflows.Events; using Indice.Features.Agents.Core.Workflows.State; @@ -156,15 +157,7 @@ public async IAsyncEnumerable GetStreamingResponseAsync(IEnu yield return new ChatResponseUpdate(ChatRole.Assistant, final?.Answer ?? "") { - RawRepresentation = new DexFinalEvent { - Answer = final?.Answer, - Citations = final?.Citations ?? [], - Sources = final?.Sources ?? [], - Failed = failure is not null, - FailureReason = failure, - Usage = usage, - ModelUsed = modelUsed, - }, + RawRepresentation = "Final", ConversationId = state.ConversationId, ResponseId = messageId, MessageId = messageId, @@ -173,6 +166,7 @@ public async IAsyncEnumerable GetStreamingResponseAsync(IEnu ["failed"] = failure is not null, ["failureReason"] = failure, ["sources"] = final?.Sources ?? [], + ["citations"] = final?.Citations ?? [], } }; diff --git a/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs b/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs deleted file mode 100644 index adc31d6c4..000000000 --- a/src/Indice.Features.Agents.Core/Workflows/Events/DexStreamEvent.cs +++ /dev/null @@ -1,84 +0,0 @@ -using Indice.Features.Agents.Core.Models; -using Microsoft.Extensions.AI; - -namespace Indice.Features.Agents.Core.Workflows.Events; - -/// -/// Base type for the real-time events yields as -/// the pipeline executes: per-step progress (), answer text deltas -/// (), and a single terminal . -/// -public abstract class DexStreamEvent -{ - /// Timestamp of when the event was created. - public DateTimeOffset TimeStamp { get; set; } = DateTimeOffset.UtcNow; -} - -/// Signals that a pipeline step has started executing. -public sealed class DexStepEvent : DexStreamEvent -{ - /// Creates a new . - public DexStepEvent(string stepId, string label) { - StepId = stepId; - Label = label; - } - - /// The executor id of the step (e.g. Retriever). - public string StepId { get; } - - /// Human-friendly progress label for the step (e.g. Retrieving relevant context). - public string Label { get; } -} - -/// Carries a single incremental chunk of the answer as the composer streams it. -public sealed class DexDeltaEvent : DexStreamEvent -{ - /// Creates a new . - public DexDeltaEvent(string text) => Text = text; - - /// The incremental answer text. - public string Text { get; } -} - -/// Signals that a pipeline step threw an exception. The pipeline will halt and yield a single with set to true. -public sealed class DexErrorEvent : DexStreamEvent -{ - /// Creates a new . - public DexErrorEvent(string stepId, string errorMessage) { - StepId = stepId; - ErrorMessage = errorMessage; - } - /// The executor id of the step that threw (e.g. Retriever). - public string StepId { get; } - - /// The error message from the step that threw. - public string ErrorMessage { get; } -} - -/// -/// Terminal event yielded once after the run completes. -/// the full answer, citations, failure state, and reasoning-model token totals. -/// -public sealed class DexFinalEvent : DexStreamEvent -{ - /// The full grounded answer (or the out-of-scope refusal); null only when a step threw. - public string? Answer { get; init; } - - /// Citations supporting the answer; empty for out-of-scope responses and on error. - public IReadOnlyList Citations { get; init; } = []; - - /// Links to the source documents that were retrieved and used to compose the answer; empty for out-of-scope responses and on error. - public IReadOnlyList Sources { get; init; } = []; - - /// True when a step threw and the workflow halted. Out-of-scope is NOT a failure. - public bool Failed { get; init; } - - /// Error message from the step that threw, prefixed with its executor id; null when not failed. - public string? FailureReason { get; init; } - - /// Total reasoning-model token usage across this run; null when no reasoning call ran. - public UsageDetails? Usage { get; init; } - - /// The reasoning-model deployment the tokens were billed against; null when no reasoning call ran. - public string? ModelUsed { get; init; } -}