using System.Text.Json; using Microsoft.EntityFrameworkCore; using w4c_workflows.Data; using w4c_workflows.Models; using w4c_workflows.Services.Messaging; using w4c_workflows.Services.Runs; namespace w4c_workflows.Services.Triggers; /// /// Queue-trigger consumer for handler-mode workflows. A handler workflow /// subscribes to a stream (trigger.stream, default wf:{tenant}:events) /// and, on each event, starts a run whose input is the event payload. /// /// v1 consumes each event as its own run; the long-lived-instance semantics /// (one durable instance, state accumulation across events, checkpoint/resume) /// are layered on by the run lifecycle engine in step 9, which reuses the same /// correlation id scheme (handler:{workflowId}:{messageId}). /// public class HandlerStreamConsumer : BackgroundService { private readonly IServiceScopeFactory _scopeFactory; private readonly IEventBus _events; private readonly ILogger _logger; private readonly string _consumer; private readonly int _pollDelayMs; private readonly int _batchSize; private readonly TimeSpan _claimIdle; public HandlerStreamConsumer( IServiceScopeFactory scopeFactory, IEventBus events, IConfiguration config, ILogger logger) { _scopeFactory = scopeFactory; _events = events; _logger = logger; _consumer = $"handler-{Environment.MachineName}-{Guid.NewGuid():N}"[..28]; _pollDelayMs = ParseInt(config["Workflows:HandlerPollDelayMs"], 500); _batchSize = ParseInt(config["Workflows:HandlerBatchSize"], 10); _claimIdle = TimeSpan.FromSeconds(ParseInt(config["Workflows:HandlerClaimIdleSeconds"], 30)); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation("HandlerStreamConsumer started as consumer {Consumer}", _consumer); while (!stoppingToken.IsCancellationRequested) { try { var processed = await ProcessAsync(stoppingToken); if (processed == 0) await Task.Delay(_pollDelayMs, stoppingToken); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { break; } catch (Exception ex) { _logger.LogError(ex, "Handler stream consumer loop error"); try { await Task.Delay(_pollDelayMs, stoppingToken); } catch (OperationCanceledException) { break; } } } _logger.LogInformation("HandlerStreamConsumer stopped"); } private async Task ProcessAsync(CancellationToken ct) { using var scope = _scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var launcher = scope.ServiceProvider.GetRequiredService(); var handlers = await db.Workflows .Where(w => w.Status == WorkflowStatus.Compiled && w.Mode == WorkflowMode.Handler) .ToListAsync(ct); var processed = 0; foreach (var workflow in handlers) { if (ct.IsCancellationRequested) return processed; if (!workflow.TriggerEnabled) continue; // auto-trigger toggled off from the UI var spec = TriggerSpec.Parse(workflow.TriggerJson, out _); if (spec == null) continue; var stream = string.IsNullOrWhiteSpace(spec.Stream) ? Streams.DefaultEvents : spec.Stream; await _events.EnsureGroupAsync(workflow.TenantId, stream, ct); var messages = await _events.ReadGroupAsync(workflow.TenantId, stream, _consumer, _batchSize, ct); if (messages.Count == 0) messages = await _events.ClaimPendingAsync(workflow.TenantId, stream, _consumer, _claimIdle, _batchSize, ct); foreach (var message in messages) { ct.ThrowIfCancellationRequested(); await DispatchAsync(launcher, workflow, stream, message, ct); processed++; } } return processed; } private async Task DispatchAsync(IRunLauncher launcher, Workflow workflow, string stream, StreamMessage message, CancellationToken ct) { var input = SerializeEventInput(message.Fields); var correlation = $"handler:{workflow.Id}:{message.Id}"; await launcher.LaunchAsync( new LaunchRequest(workflow.TenantId, workflow.Id, workflow.TriggerJson, input, correlation), ct); await _events.AckAsync(workflow.TenantId, stream, message.Id, ct); } /// Serializes a stream message's fields as a JSON object (the event payload). internal static string SerializeEventInput(IReadOnlyDictionary fields) => JsonSerializer.Serialize(fields, JsonSerializerOptions.Web); private static int ParseInt(string? text, int fallback) => int.TryParse(text, out var value) && value > 0 ? value : fallback; }