< Summary

Information
Class: Elsa.AI.Host.Services.AIOrchestrator
Assembly: Elsa.AI.Host
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.AI.Host/Services/AIOrchestrator.cs
Line coverage
97%
Covered lines: 393
Uncovered lines: 11
Coverable lines: 404
Total lines: 623
Line coverage: 97.2%
Branch coverage
85%
Covered branches: 186
Total branches: 218
Branch coverage: 85.3%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
ExecuteChatAsync()98.68%7676100%
CreateEvent(...)100%22100%
ReadProviderEventsAsync()100%6695.23%
SelectProvider(...)80%212086.95%
FindAgentProviderName(...)75%44100%
RecordChatAuditAsync()100%44100%
InvokeProviderToolAsync()100%44100%
RecordToolAuditEventsAsync()100%1140%
CreateToolAuditEvent(...)100%11100%
LimitResolvedContext(...)100%88100%
TruncateContext(...)100%11100%
LimitToolResult(...)100%22100%
CreateTruncatedPayload(...)100%11100%
GetUtf8Size(...)100%11100%
Truncate(...)78.57%151481.25%
TryReadAssistantContent(...)100%66100%
TryReadToolResult(...)44.44%363693.33%
CreateMessage(...)100%22100%
HasReconnectUserMessage(...)100%44100%
IsCompletedReconnect(...)100%66100%
HasUserMessage(...)50%22100%
GetNextSequence(...)50%22100%
NormalizeMessage(...)100%11100%
BelongsToTenant(...)100%11100%
BelongsToUser(...)50%22100%
NormalizeTenantId(...)100%22100%
TrySaveConversationAsync()93.75%1616100%
get_Id()100%11100%
get_ToolCallId()100%11100%
get_Event()100%11100%
get_Provider()100%11100%
.ctor(...)100%11100%
InvokeAsync(...)100%11100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.AI.Host/Services/AIOrchestrator.cs

#LineLine coverage
 1using System.Text;
 2using System.Text.Json;
 3using Elsa.AI.Abstractions.Contracts;
 4using Elsa.AI.Abstractions.Models;
 5using Elsa.AI.Host.Context;
 6using Elsa.AI.Host.Options;
 7using Elsa.AI.Host.Streaming;
 8using Microsoft.Extensions.Logging;
 9using Microsoft.Extensions.Options;
 10
 11namespace Elsa.AI.Host.Services;
 12
 4213public class AIOrchestrator(
 4214    IEnumerable<IAIProvider> providers,
 4215    IAIToolRegistry toolRegistry,
 4216    IAIConversationStore conversationStore,
 4217    AIContextResolver contextResolver,
 4218    AIStreamEventMapper streamEventMapper,
 4219    IAIAuditSink auditSink,
 4220    ILogger<AIOrchestrator> logger,
 4221    IOptions<AIHostOptions> options) : IAIOrchestrator
 22{
 23    public async IAsyncEnumerable<AIStreamEvent> ExecuteChatAsync(AIChatRequest request, [System.Runtime.CompilerService
 24    {
 4325        var conversationId = request.ConversationId ?? Guid.NewGuid().ToString("N");
 4326        var sequence = 0L;
 4327        var providerSelection = SelectProvider(request);
 4328        var provider = providerSelection.Provider;
 4329        var conversationPersistenceEnabled = options.Value.ConversationPersistenceEnabled;
 4330        AIConversation? conversation = null;
 4331        Exception? preparationError = null;
 4332        if (conversationPersistenceEnabled)
 33        {
 34            try
 35            {
 4236                conversation = await conversationStore.FindAsync(conversationId, cancellationToken);
 4137            }
 138            catch (Exception e) when (e is not OperationCanceledException)
 39            {
 140                preparationError = e;
 141            }
 42        }
 43
 4344        if (conversation != null && (!BelongsToTenant(conversation, request.TenantId) || !BelongsToUser(conversation, re
 45        {
 246            conversation = null;
 247            conversationId = Guid.NewGuid().ToString("N");
 48        }
 4349        var messages = conversation?.Messages.ToList() ?? [];
 4350        if (request.IsReconnect && messages.Count == 0)
 151            conversationId = Guid.NewGuid().ToString("N");
 52
 4353        if (request.IsReconnect && IsCompletedReconnect(conversation, request.Message))
 54        {
 255            var nextSequence = GetNextSequence(messages);
 256            if (conversation!.Status == AIConversationStatus.Failed)
 57            {
 258                var lastAssistantContent = conversation.Messages.LastOrDefault(x => x.Role == AIMessageRole.Assistant)?.
 159                if (!string.IsNullOrEmpty(lastAssistantContent))
 160                    yield return CreateEvent("conversation.error", conversationId, nextSequence++, new JsonObject
 161                    {
 162                        ["content"] = lastAssistantContent
 163                    });
 64            }
 65
 266            yield return CreateEvent("conversation.completed", conversationId, nextSequence);
 267            yield break;
 68        }
 69
 4170        var providerSessionId = conversation?.ProviderSessionId;
 71
 4172        if (preparationError == null && provider != null && string.IsNullOrWhiteSpace(providerSessionId))
 73        {
 74            try
 75            {
 2576                var session = await provider.CreateSessionAsync(new CreateAISessionRequest
 2577                {
 2578                    ConversationId = conversationId,
 2579                    Agent = request.Agent,
 2580                    TenantId = request.TenantId,
 2581                    ProviderConfiguration = providerSelection.Configuration
 2582                }, cancellationToken);
 2483                providerSessionId = session.ProviderSessionId ?? session.Id;
 2484            }
 185            catch (Exception e) when (e is not OperationCanceledException)
 86            {
 187                preparationError = e;
 188            }
 89        }
 90
 4191        var isDuplicateReconnectMessage = request.IsReconnect && HasReconnectUserMessage(conversation, request.Message);
 4192        if (request.IsReconnect && messages.Count > 0)
 393            sequence = GetNextSequence(messages);
 94
 4195        yield return CreateEvent("conversation.started", conversationId, sequence++);
 96
 4197        var userMessage = isDuplicateReconnectMessage
 598            ? messages.Last(x => x.Role == AIMessageRole.User && string.Equals(NormalizeMessage(x.Content), NormalizeMes
 4199            : CreateMessage(conversationId, AIMessageRole.User, request.Message, sequence++);
 100
 41101        if (!isDuplicateReconnectMessage)
 38102            messages.Add(userMessage);
 103
 41104        await TrySaveConversationAsync(conversationId, request, AIConversationStatus.Active, messages, conversation, pro
 41105        await RecordChatAuditAsync("chat.started", request, conversationId, provider?.Name, cancellationToken);
 106
 41107        IReadOnlyCollection<AIResolvedContext> context = [];
 41108        IReadOnlyCollection<AIToolDefinition> tools = [];
 109
 41110        if (preparationError == null)
 111        {
 112            try
 113            {
 39114                context = LimitResolvedContext(await contextResolver.ResolveAsync(request, cancellationToken));
 38115                tools = await toolRegistry.ListAsync(new AIToolQuery
 38116                {
 38117                    Agent = request.Agent,
 38118                    ActorId = request.UserId,
 38119                    TenantId = request.TenantId,
 38120                    UserPermissions = request.UserPermissions
 38121                }, cancellationToken);
 38122            }
 1123            catch (Exception e) when (e is not OperationCanceledException)
 124            {
 1125                preparationError = e;
 1126            }
 127        }
 128
 41129        if (preparationError != null)
 130        {
 131            const string content = "Weaver could not prepare AI context or tools for this request.";
 3132            logger.LogWarning(preparationError, "Failed to prepare AI chat context or tools for conversation {Conversati
 3133            yield return CreateEvent("conversation.error", conversationId, sequence++, new JsonObject
 3134            {
 3135                ["content"] = content
 3136            });
 3137            messages.Add(CreateMessage(conversationId, AIMessageRole.Assistant, content, sequence - 1));
 3138            await TrySaveConversationAsync(conversationId, request, AIConversationStatus.Failed, messages, conversation,
 3139            await RecordChatAuditAsync("chat.failed", request, conversationId, provider?.Name, cancellationToken);
 3140            yield return CreateEvent("conversation.completed", conversationId, sequence);
 3141            yield break;
 142        }
 143
 38144        if (provider == null)
 145        {
 146            const string content = "Weaver is ready, but no AI provider is configured.";
 11147            yield return CreateEvent("assistant.delta", conversationId, sequence++, new JsonObject
 11148            {
 11149                ["content"] = content
 11150            });
 11151            messages.Add(CreateMessage(conversationId, AIMessageRole.Assistant, content, sequence - 1));
 152        }
 153        else
 154        {
 27155            var assistantContent = new StringBuilder();
 27156            var providerHistory = isDuplicateReconnectMessage
 27157                ? messages.ToList()
 57158                : messages.Where(x => x.Id != userMessage.Id).ToList();
 27159            var turnRequest = new AITurnRequest
 27160            {
 27161                ConversationId = conversationId,
 27162                ProviderSessionId = providerSessionId,
 27163                Message = isDuplicateReconnectMessage ? "" : request.Message,
 27164                Messages = providerHistory,
 27165                Context = context,
 411166                Tools = tools.Where(x => x.IsEnabled).ToList(),
 27167                Agent = request.Agent,
 27168                ProviderConfiguration = providerSelection.Configuration
 27169            };
 27170            var toolInvoker = new HostToolInvoker(this, request, conversationId);
 27171            Exception? providerTurnError = null;
 172
 119173            await foreach (var providerRead in ReadProviderEventsAsync(provider.ExecuteTurnAsync(turnRequest, toolInvoke
 174            {
 33175                if (providerRead.Error != null)
 176                {
 1177                    providerTurnError = providerRead.Error;
 1178                    break;
 179                }
 180
 32181                var providerEvent = providerRead.Event!;
 32182                var streamEvent = streamEventMapper.Map(conversationId, providerEvent) with { Sequence = sequence++ };
 32183                yield return streamEvent;
 184
 32185                if (TryReadAssistantContent(providerEvent, out var content))
 21186                    assistantContent.Append(content);
 187
 32188                if (TryReadToolResult(providerEvent, out var toolResultMessage))
 7189                    messages.Add(CreateMessage(conversationId, AIMessageRole.Tool, toolResultMessage.Summary, streamEven
 7190                    {
 7191                        ["toolCallId"] = toolResultMessage.ToolCallId,
 7192                        ["toolName"] = toolResultMessage.ToolName,
 7193                        ["status"] = toolResultMessage.Status
 7194                    }));
 32195            }
 196
 27197            if (providerTurnError != null)
 198            {
 199                const string content = "Weaver could not complete the AI provider turn for this request.";
 1200                logger.LogWarning(providerTurnError, "Failed to execute AI provider turn for conversation {ConversationI
 1201                yield return CreateEvent("conversation.error", conversationId, sequence++, new JsonObject
 1202                {
 1203                    ["content"] = content
 1204                });
 1205                messages.Add(CreateMessage(conversationId, AIMessageRole.Assistant, content, sequence - 1));
 1206                await TrySaveConversationAsync(conversationId, request, AIConversationStatus.Failed, messages, conversat
 1207                await RecordChatAuditAsync("chat.failed", request, conversationId, provider.Name, cancellationToken);
 1208                yield return CreateEvent("conversation.completed", conversationId, sequence);
 1209                yield break;
 210            }
 211
 26212            if (assistantContent.Length > 0)
 21213                messages.Add(CreateMessage(conversationId, AIMessageRole.Assistant, assistantContent.ToString(), sequenc
 26214        }
 215
 37216        await TrySaveConversationAsync(conversationId, request, AIConversationStatus.Completed, messages, conversation, 
 37217        await RecordChatAuditAsync("chat.completed", request, conversationId, provider?.Name, cancellationToken);
 37218        yield return CreateEvent("conversation.completed", conversationId, sequence);
 43219    }
 220
 221    private static AIStreamEvent CreateEvent(string type, string conversationId, long sequence, JsonObject? data = null)
 100222        new()
 100223        {
 100224            Type = type,
 100225            ConversationId = conversationId,
 100226            Sequence = sequence,
 100227            Timestamp = DateTimeOffset.UtcNow,
 100228            Data = data ?? []
 100229        };
 230
 231    private static async IAsyncEnumerable<ProviderReadResult> ReadProviderEventsAsync(
 232        IAsyncEnumerable<AIProviderEvent> providerEvents,
 233        [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken)
 234    {
 27235        var enumerator = providerEvents.GetAsyncEnumerator(cancellationToken);
 236        try
 237        {
 32238            while (true)
 239            {
 59240                AIProviderEvent? providerEvent = null;
 59241                Exception? error = null;
 59242                var hasEvent = false;
 243
 244                try
 245                {
 59246                    hasEvent = await enumerator.MoveNextAsync();
 58247                    if (hasEvent)
 32248                        providerEvent = enumerator.Current;
 58249                }
 1250                catch (Exception e) when (e is not OperationCanceledException)
 251                {
 1252                    error = e;
 1253                }
 254
 59255                if (error != null)
 256                {
 1257                    yield return new ProviderReadResult(null, error);
 0258                    yield break;
 259                }
 260
 58261                if (!hasEvent)
 26262                    yield break;
 263
 32264                yield return new ProviderReadResult(providerEvent, null);
 32265            }
 266        }
 267        finally
 268        {
 27269            await enumerator.DisposeAsync();
 270        }
 27271    }
 272
 273    private ProviderSelection SelectProvider(AIChatRequest request)
 274    {
 43275        var providerOptions = options.Value.Providers.ToList();
 45276        var configuredProviders = providerOptions.Where(x => x.Enabled).ToList();
 43277        var availableProviders = providers
 32278            .Where(x => providerOptions.IsProviderEnabled(x.Name))
 43279            .ToList();
 43280        var providerName = request.ProviderName ?? FindAgentProviderName(request.Agent) ?? options.Value.DefaultProvider
 281
 43282        if (!string.IsNullOrWhiteSpace(providerName))
 283        {
 5284            var configuredProvider = configuredProviders.FirstOrDefault(x => string.Equals(x.Name, providerName, StringC
 4285            var provider = configuredProvider != null
 1286                ? availableProviders.FirstOrDefault(x => string.Equals(x.Name, configuredProvider.Name, StringComparison
 1287                  availableProviders.FirstOrDefault(x => string.Equals(x.Name, configuredProvider.Provider, StringCompar
 8288                : availableProviders.FirstOrDefault(x => string.Equals(x.Name, providerName, StringComparison.OrdinalIgn
 289
 4290            return new ProviderSelection(provider, configuredProvider?.ToProviderConfiguration());
 291        }
 292
 39293        if (availableProviders.Count != 1)
 294        {
 13295            if (availableProviders.Count > 1)
 0296                logger.LogWarning(
 0297                    "Multiple AI providers are available ({ProviderNames}) but no default provider name is configured. S
 0298                    string.Join(", ", availableProviders.Select(x => x.Name)));
 299
 13300            return new ProviderSelection(null, null);
 301        }
 302
 26303        var selectedProvider = availableProviders[0];
 26304        var selectedConfiguration = configuredProviders.FirstOrDefault(x => string.Equals(x.Name, selectedProvider.Name,
 26305                                                                            string.Equals(x.Provider, selectedProvider.N
 306
 26307        return new ProviderSelection(selectedProvider, selectedConfiguration?.ToProviderConfiguration());
 308    }
 309
 310    private string? FindAgentProviderName(string? agent) =>
 40311        string.IsNullOrWhiteSpace(agent)
 40312            ? null
 41313            : options.Value.Agents.FirstOrDefault(x => string.Equals(x.Name, agent, StringComparison.OrdinalIgnoreCase))
 314
 315    private async ValueTask RecordChatAuditAsync(string type, AIChatRequest request, string conversationId, string? prov
 316    {
 317        try
 318        {
 82319            await auditSink.RecordAsync(new AIAuditEvent
 82320            {
 82321                Type = type,
 82322                TenantId = request.TenantId,
 82323                ActorId = request.UserId,
 82324                ConversationId = conversationId,
 82325                Timestamp = DateTimeOffset.UtcNow,
 82326                Summary = type switch
 82327                {
 41328                    "chat.started" => "Chat started",
 4329                    "chat.failed" => "Chat failed",
 37330                    _ => "Chat completed"
 82331                },
 82332                Data = new JsonObject
 82333                {
 82334                    ["agent"] = request.Agent,
 82335                    ["provider"] = providerName,
 82336                    ["attachmentCount"] = request.Attachments.Count
 82337                }
 82338            }, cancellationToken);
 80339        }
 2340        catch (Exception e) when (e is not OperationCanceledException)
 341        {
 2342            logger.LogWarning(e, "Failed to record AI chat audit event {AuditEventType} for conversation {ConversationId
 2343        }
 82344    }
 345
 346    private async ValueTask<AIToolResult> InvokeProviderToolAsync(AIProviderToolInvocation invocation, AIChatRequest req
 347    {
 6348        var toolCall = new ToolCall(invocation.Id, invocation.ToolName, invocation.Arguments);
 6349        var tool = await toolRegistry.FindAsync(invocation.ToolName, new AIToolQuery
 6350        {
 6351            Agent = request.Agent,
 6352            ActorId = request.UserId,
 6353            TenantId = request.TenantId,
 6354            UserPermissions = request.UserPermissions
 6355        }, cancellationToken);
 6356        if (tool == null)
 357        {
 1358            var result = new AIToolResult { Status = AIToolInvocationStatus.Failed, Error = $"Tool '{toolCall.Name}' was
 1359            await RecordToolAuditEventsAsync(request, conversationId, toolCall, ["tool.failed"], cancellationToken);
 1360            return result;
 361        }
 362
 5363        using var toolScope = tool;
 364        try
 365        {
 5366            await RecordToolAuditEventsAsync(request, conversationId, toolCall, ["tool.invoked"], cancellationToken);
 5367            var result = await tool.ExecuteAsync(new AIToolExecutionContext
 5368            {
 5369                ConversationId = conversationId,
 5370                TenantId = request.TenantId,
 5371                ActorId = request.UserId,
 5372                Agent = request.Agent,
 5373                Arguments = toolCall.Arguments
 5374            }, cancellationToken);
 4375            await RecordToolAuditEventsAsync(request, conversationId, toolCall, ["tool.completed"], cancellationToken);
 376
 4377            return LimitToolResult(result);
 378        }
 1379        catch (Exception e) when (e is not OperationCanceledException)
 380        {
 1381            logger.LogWarning(e, "AI tool {ToolName} failed for conversation {ConversationId}.", toolCall.Name, conversa
 1382            await RecordToolAuditEventsAsync(request, conversationId, toolCall, ["tool.failed"], cancellationToken);
 1383            return new AIToolResult { Status = AIToolInvocationStatus.Failed, Error = "Tool execution failed." };
 384        }
 6385    }
 386
 387    private async ValueTask RecordToolAuditEventsAsync(AIChatRequest request, string conversationId, ToolCall toolCall, 
 388    {
 389        try
 390        {
 22391            await auditSink.RecordManyAsync(types.Select(type => CreateToolAuditEvent(type, request, conversationId, too
 11392        }
 0393        catch (Exception e) when (e is not OperationCanceledException)
 394        {
 0395            logger.LogWarning(e, "Failed to record AI tool audit events for tool {ToolName}.", toolCall.Name);
 0396        }
 11397    }
 398
 399    private static AIAuditEvent CreateToolAuditEvent(string type, AIChatRequest request, string conversationId, ToolCall
 11400        new()
 11401        {
 11402            Type = type,
 11403            TenantId = request.TenantId,
 11404            ActorId = request.UserId,
 11405            ConversationId = conversationId,
 11406            ToolInvocationId = toolCall.Id,
 11407            Timestamp = DateTimeOffset.UtcNow,
 11408            Summary = $"{toolCall.Name} {type}",
 11409            Data = new JsonObject
 11410            {
 11411                ["toolName"] = toolCall.Name
 11412            }
 11413        };
 414
 415    private IReadOnlyCollection<AIResolvedContext> LimitResolvedContext(IReadOnlyCollection<AIResolvedContext> contexts)
 416    {
 38417        var maxBytes = options.Value.MaxResolvedContextBytes;
 38418        if (maxBytes <= 0)
 1419            return contexts;
 420
 37421        var limited = new List<AIResolvedContext>();
 37422        var usedBytes = 0;
 423
 84424        foreach (var context in contexts)
 425        {
 7426            var contextSize = GetUtf8Size(context);
 7427            if (usedBytes + contextSize <= maxBytes)
 428            {
 2429                usedBytes += contextSize;
 2430                limited.Add(context);
 2431                continue;
 432            }
 433
 5434            if (limited.Count > 0)
 435            {
 1436                logger.LogDebug(
 1437                    "Dropping AI resolved context {ContextKind}/{ReferenceId} because resolved context exceeds the confi
 1438                    context.Kind,
 1439                    context.ReferenceId,
 1440                    maxBytes);
 1441                continue;
 442            }
 443
 4444            limited.Add(TruncateContext(context, maxBytes));
 4445            break;
 446        }
 447
 37448        return limited;
 449    }
 450
 451    private static AIResolvedContext TruncateContext(AIResolvedContext context, int maxBytes) =>
 4452        context with
 4453        {
 4454            Summary = Truncate(context.Summary, maxBytes),
 4455            Data = CreateTruncatedPayload(maxBytes),
 4456            Metadata = CreateTruncatedPayload(maxBytes)
 4457        };
 458
 459    private AIToolResult LimitToolResult(AIToolResult result)
 460    {
 4461        if (GetUtf8Size(result) <= options.Value.MaxToolResultBytes)
 3462            return result;
 463
 1464        return result with
 1465        {
 1466            Summary = Truncate(result.Summary, options.Value.MaxToolResultBytes),
 1467            Data = CreateTruncatedPayload(options.Value.MaxToolResultBytes)
 1468        };
 469    }
 470
 471    private static JsonObject CreateTruncatedPayload(int maxBytes) =>
 9472        new()
 9473        {
 9474            ["truncated"] = true,
 9475            ["maxBytes"] = maxBytes
 9476        };
 477
 478    private static int GetUtf8Size<T>(T value) =>
 11479        JsonSerializer.SerializeToUtf8Bytes(value).Length;
 480
 481    private static string Truncate(string value, int maxBytes)
 482    {
 5483        if (string.IsNullOrEmpty(value))
 0484            return value;
 485
 5486        if (maxBytes <= 0)
 0487            return "";
 488
 5489        if (Encoding.UTF8.GetByteCount(value) <= maxBytes)
 0490            return value;
 491
 5492        var low = 0;
 5493        var high = Math.Min(value.Length, maxBytes);
 34494        while (low < high)
 495        {
 29496            var candidate = (low + high + 1) / 2;
 29497            if (Encoding.UTF8.GetByteCount(value.AsSpan(0, candidate)) <= maxBytes)
 25498                low = candidate;
 499            else
 4500                high = candidate - 1;
 501        }
 502
 5503        if (low > 0 && char.IsHighSurrogate(value[low - 1]))
 1504            low--;
 505
 5506        return value[..low];
 507    }
 508
 509    private static bool TryReadAssistantContent(AIProviderEvent providerEvent, out string content)
 510    {
 32511        content = "";
 32512        if (!string.Equals(providerEvent.Type, "assistant.delta", StringComparison.OrdinalIgnoreCase))
 7513            return false;
 514
 25515        content = providerEvent.Data["content"]?.GetValue<string>() ?? "";
 25516        return !string.IsNullOrEmpty(content);
 517    }
 518
 519    private static bool TryReadToolResult(AIProviderEvent providerEvent, out ToolResultMessage toolResult)
 520    {
 32521        toolResult = default;
 32522        if (!string.Equals(providerEvent.Type, "tool.result", StringComparison.OrdinalIgnoreCase) &&
 32523            !string.Equals(providerEvent.Type, "tool.completed", StringComparison.OrdinalIgnoreCase))
 25524            return false;
 525
 7526        var toolCallId = providerEvent.Data["toolCallId"]?.GetValue<string>() ?? providerEvent.Data["id"]?.GetValue<stri
 7527        var toolName = providerEvent.Data["toolName"]?.GetValue<string>() ?? providerEvent.Data["name"]?.GetValue<string
 7528        if (string.IsNullOrWhiteSpace(toolCallId) || string.IsNullOrWhiteSpace(toolName))
 0529            return false;
 530
 7531        var summary = providerEvent.Data["summary"]?.GetValue<string>() ??
 7532                      providerEvent.Data["content"]?.GetValue<string>() ??
 7533                      providerEvent.Data["result"]?.GetValue<string>() ??
 7534                      "";
 7535        var status = providerEvent.Data["status"]?.GetValue<string>() ?? AIToolInvocationStatus.Completed.ToString();
 7536        toolResult = new ToolResultMessage(toolCallId, toolName, status, summary);
 7537        return true;
 538    }
 539
 540    private static AIMessage CreateMessage(string conversationId, AIMessageRole role, string content, long streamSequenc
 81541        new()
 81542        {
 81543            Id = Guid.NewGuid().ToString("N"),
 81544            ConversationId = conversationId,
 81545            Role = role,
 81546            Content = content,
 81547            CreatedAt = DateTimeOffset.UtcNow,
 81548            StreamSequence = streamSequence,
 81549            Metadata = metadata ?? []
 81550        };
 551
 552    private static bool HasReconnectUserMessage(AIConversation? conversation, string message)
 553    {
 4554        return conversation is { Status: AIConversationStatus.Active } &&
 4555               HasUserMessage(conversation, message);
 556    }
 557
 558    private static bool IsCompletedReconnect(AIConversation? conversation, string message) =>
 6559        conversation is { Status: AIConversationStatus.Completed or AIConversationStatus.Failed } && HasUserMessage(conv
 560
 561    private static bool HasUserMessage(AIConversation conversation, string message) =>
 10562        conversation.Messages.Any(x => x.Role == AIMessageRole.User && string.Equals(NormalizeMessage(x.Content), Normal
 563
 564    private static long GetNextSequence(IReadOnlyCollection<AIMessage> messages) =>
 14565        messages.Count == 0 ? 0 : messages.Max(x => x.StreamSequence) + 1;
 566
 567    private static string NormalizeMessage(string message) =>
 16568        message.ReplaceLineEndings("\n").Trim();
 569
 570    private static bool BelongsToTenant(AIConversation conversation, string? tenantId) =>
 10571        string.Equals(NormalizeTenantId(conversation.TenantId), NormalizeTenantId(tenantId), StringComparison.Ordinal);
 572
 573    private static bool BelongsToUser(AIConversation conversation, string userId) =>
 9574        string.IsNullOrWhiteSpace(conversation.UserId) || string.Equals(conversation.UserId, userId, StringComparison.Or
 575
 576    private static string NormalizeTenantId(string? tenantId) =>
 20577        string.IsNullOrWhiteSpace(tenantId) ? "" : tenantId;
 578
 579    private async ValueTask TrySaveConversationAsync(string conversationId, AIChatRequest request, AIConversationStatus 
 580    {
 82581        if (!options.Value.ConversationPersistenceEnabled)
 2582            return;
 583
 584        try
 585        {
 80586            var now = DateTimeOffset.UtcNow;
 80587            var retentionMode = conversation?.RetentionMode ?? AIRetentionMode.Configured;
 80588            DateTimeOffset? retentionExpiresAt = retentionMode == AIRetentionMode.Configured
 80589                ? now.Add(options.Value.ConversationRetention)
 80590                : null;
 591
 80592            await conversationStore.SaveAsync(new AIConversation
 80593            {
 80594                Id = conversationId,
 80595                TenantId = request.TenantId,
 80596                UserId = request.UserId,
 80597                Title = conversation?.Title,
 80598                Status = status,
 80599                CreatedAt = conversation is null || conversation.CreatedAt == default ? now : conversation.CreatedAt,
 80600                UpdatedAt = now,
 80601                ProviderSessionId = providerSessionId ?? conversation?.ProviderSessionId,
 80602                RetentionMode = retentionMode,
 80603                RetentionExpiresAt = retentionExpiresAt,
 80604                Messages = messages.ToList()
 80605            }, cancellationToken);
 78606        }
 2607        catch (Exception e) when (e is not OperationCanceledException)
 608        {
 2609            logger.LogWarning(e, "Failed to persist AI conversation {ConversationId} with status {ConversationStatus}.",
 2610        }
 82611    }
 612
 40613    private readonly record struct ToolCall(string Id, string Name, JsonObject Arguments);
 28614    private readonly record struct ToolResultMessage(string ToolCallId, string ToolName, string Status, string Summary);
 66615    private readonly record struct ProviderReadResult(AIProviderEvent? Event, Exception? Error);
 95616    private readonly record struct ProviderSelection(IAIProvider? Provider, AIProviderConfiguration? Configuration);
 617
 27618    private class HostToolInvoker(AIOrchestrator orchestrator, AIChatRequest request, string conversationId) : IAIProvid
 619    {
 620        public ValueTask<AIToolResult> InvokeAsync(AIProviderToolInvocation invocation, CancellationToken cancellationTo
 6621            orchestrator.InvokeProviderToolAsync(invocation, request, conversationId, cancellationToken);
 622    }
 623}

Methods/Properties

.ctor(System.Collections.Generic.IEnumerable`1<Elsa.AI.Abstractions.Contracts.IAIProvider>,Elsa.AI.Abstractions.Contracts.IAIToolRegistry,Elsa.AI.Abstractions.Contracts.IAIConversationStore,Elsa.AI.Host.Context.AIContextResolver,Elsa.AI.Host.Streaming.AIStreamEventMapper,Elsa.AI.Abstractions.Contracts.IAIAuditSink,Microsoft.Extensions.Logging.ILogger`1<Elsa.AI.Host.Services.AIOrchestrator>,Microsoft.Extensions.Options.IOptions`1<Elsa.AI.Host.Options.AIHostOptions>)
ExecuteChatAsync()
CreateEvent(System.String,System.String,System.Int64,System.Text.Json.Nodes.JsonObject)
ReadProviderEventsAsync()
SelectProvider(Elsa.AI.Abstractions.Models.AIChatRequest)
FindAgentProviderName(System.String)
RecordChatAuditAsync()
InvokeProviderToolAsync()
RecordToolAuditEventsAsync()
CreateToolAuditEvent(System.String,Elsa.AI.Abstractions.Models.AIChatRequest,System.String,Elsa.AI.Host.Services.AIOrchestrator/ToolCall)
LimitResolvedContext(System.Collections.Generic.IReadOnlyCollection`1<Elsa.AI.Abstractions.Models.AIResolvedContext>)
TruncateContext(Elsa.AI.Abstractions.Models.AIResolvedContext,System.Int32)
LimitToolResult(Elsa.AI.Abstractions.Models.AIToolResult)
CreateTruncatedPayload(System.Int32)
GetUtf8Size(T)
Truncate(System.String,System.Int32)
TryReadAssistantContent(Elsa.AI.Abstractions.Models.AIProviderEvent,System.String&)
TryReadToolResult(Elsa.AI.Abstractions.Models.AIProviderEvent,Elsa.AI.Host.Services.AIOrchestrator/ToolResultMessage&)
CreateMessage(System.String,Elsa.AI.Abstractions.Models.AIMessageRole,System.String,System.Int64,System.Text.Json.Nodes.JsonObject)
HasReconnectUserMessage(Elsa.AI.Abstractions.Models.AIConversation,System.String)
IsCompletedReconnect(Elsa.AI.Abstractions.Models.AIConversation,System.String)
HasUserMessage(Elsa.AI.Abstractions.Models.AIConversation,System.String)
GetNextSequence(System.Collections.Generic.IReadOnlyCollection`1<Elsa.AI.Abstractions.Models.AIMessage>)
NormalizeMessage(System.String)
BelongsToTenant(Elsa.AI.Abstractions.Models.AIConversation,System.String)
BelongsToUser(Elsa.AI.Abstractions.Models.AIConversation,System.String)
NormalizeTenantId(System.String)
TrySaveConversationAsync()
get_Id()
get_ToolCallId()
get_Event()
get_Provider()
.ctor(Elsa.AI.Host.Services.AIOrchestrator,Elsa.AI.Abstractions.Models.AIChatRequest,System.String)
InvokeAsync(Elsa.AI.Abstractions.Models.AIProviderToolInvocation,System.Threading.CancellationToken)