Bygga ett arbetsflödessystem med HTMX och ASP.NET Core - Del 4: Hangfire Integration och Automation (Svenska (Swedish))

Bygga ett arbetsflödessystem med HTMX och ASP.NET Core - Del 4: Hangfire Integration och Automation

Wednesday, 15 January 2025

//

12 minute read

Inledning

TillHäfte 3Vi byggde en vacker bildredigerare.

  • **Men våra arbetsflöden körs bara när vi manuellt utlöser dem.**I det här sista inlägget kommer vi att göra arbetsflöden verkligt autonoma med Hangfire för:
  • Schemalagd körning- Kör arbetsflöden enligt ett schema
  • API- röstning- Övervaka externa API:er och utlösa ändringar
  • Statlig förvaltning- Spåra utlösande stater över avrättningar

Dashboard

  • Övervaka alla bakgrundsjobb

  • Varför Hangfire?

  • Hangfire är perfekt för våra behov eftersom det:

  • Lagrar jobb i vår befintliga PostgreSQL-databas

  • Ger en inbyggd instrumentbrädan

  • Stöder återkommande jobb

Har automatisk försökslogik

Vågar horisontellt

[Table("workflow_trigger_states")]
public class WorkflowTriggerStateEntity
{
    public int Id { get; set; }
    public int WorkflowDefinitionId { get; set; }

    // Type: "Schedule", "ApiPoll", "Webhook"
    public string TriggerType { get; set; } = string.Empty;

    // Configuration as JSON
    public string ConfigJson { get; set; } = "{}";

    // Current state as JSON (stores last poll time, content hash, etc.)
    public string StateJson { get; set; } = "{}";

    public bool IsEnabled { get; set; } = true;
    public DateTime? LastCheckedAt { get; set; }
    public DateTime? LastFiredAt { get; set; }
    public int FireCount { get; set; } = 0;
    public string? LastError { get; set; }
}

Utlösarstatsmodellen

  • För det första, låt oss förstå vår utlösande stat enhet (vi skapade redan detta i del 2):
  • Denna enhet spårar allt om ett arbetsflöde trigger:
  • När det senast sprang
  • Vad dess konfiguration är

Vilket tillstånd det är i (för statiska triggers)

Eventuella fel som uppstod

public class ScheduleTriggerConfig
{
    public string IntervalType { get; set; } = "minutes"; // minutes, hours, days
    public int IntervalValue { get; set; } = 60;
    public Dictionary<string, object>? InputData { get; set; }
}

Schemalagda arbetsflöden

public class WorkflowSchedulerJob
{
    private readonly MostlylucidDbContext _context;
    private readonly WorkflowExecutionService _executionService;
    private readonly ILogger<WorkflowSchedulerJob> _logger;

    [AutomaticRetry(Attempts = 3)]
    public async Task ExecuteScheduledWorkflowsAsync()
    {
        _logger.LogInformation("Checking for scheduled workflows");

        // Get all enabled schedule triggers
        var triggers = await _context.WorkflowTriggerStates
            .Include(t => t.WorkflowDefinition)
            .Where(t => t.IsEnabled && t.TriggerType == "Schedule")
            .ToListAsync();

        foreach (var trigger in triggers)
        {
            try
            {
                var config = JsonSerializer.Deserialize<ScheduleTriggerConfig>(
                    trigger.ConfigJson);

                if (config == null) continue;

                // Check if it's time to run
                if (!ShouldRunScheduledWorkflow(trigger, config))
                    continue;

                _logger.LogInformation(
                    "Executing scheduled workflow {WorkflowId}",
                    trigger.WorkflowDefinition.WorkflowId);

                // Execute the workflow
                await _executionService.ExecuteWorkflowAsync(
                    trigger.WorkflowDefinition.WorkflowId,
                    config.InputData,
                    "Scheduler");

                // Update trigger state
                trigger.LastCheckedAt = DateTime.UtcNow;
                trigger.LastFiredAt = DateTime.UtcNow;
                trigger.FireCount++;

                var state = JsonSerializer.Deserialize<Dictionary<string, object>>(
                                trigger.StateJson) ?? new();
                state["lastRun"] = DateTime.UtcNow.ToString("O");
                trigger.StateJson = JsonSerializer.Serialize(state);

                await _context.SaveChangesAsync();
            }
            catch (Exception ex)
            {
                _logger.LogError(ex,
                    "Error executing scheduled workflow {TriggerId}",
                    trigger.Id);
                trigger.LastError = ex.Message;
                await _context.SaveChangesAsync();
            }
        }
    }

    private bool ShouldRunScheduledWorkflow(
        WorkflowTriggerStateEntity trigger,
        ScheduleTriggerConfig config)
    {
        // First run?
        if (!trigger.LastFiredAt.HasValue)
            return true;

        var timeSinceLastRun = DateTime.UtcNow - trigger.LastFiredAt.Value;

        return config.IntervalType.ToLower() switch
        {
            "minutes" => timeSinceLastRun.TotalMinutes >= config.IntervalValue,
            "hours" => timeSinceLastRun.TotalHours >= config.IntervalValue,
            "days" => timeSinceLastRun.TotalDays >= config.IntervalValue,
            _ => false
        };
    }
}

Inställningsmodell

  1. Jobbet SchemaläggareExecuteScheduledWorkflowsAsync()
  2. Hur det fungerar:
  3. Varje minut ringer Hangfire.
  4. Vi frågar efter aktiverade schema triggers
  5. För varje utlösare, kontrollera om tillräckligt med tid har passerat

Om ja, kör arbetsflödet

Uppdatera utlösningstillståndet med senaste körningstid

API-mätning

public class ApiPollTriggerConfig
{
    public string Url { get; set; } = string.Empty;
    public int IntervalSeconds { get; set; } = 300; // 5 minutes
    public bool AlwaysTrigger { get; set; } = false;
    public Dictionary<string, string>? Headers { get; set; }
}

API-röstning är mer intressant - vi övervakar externa API:er och utlöser arbetsflöden när innehållet ändras!

[AutomaticRetry(Attempts = 3)]
public async Task PollApiTriggersAsync()
{
    _logger.LogInformation("Polling API triggers");

    var triggers = await _context.WorkflowTriggerStates
        .Include(t => t.WorkflowDefinition)
        .Where(t => t.IsEnabled && t.TriggerType == "ApiPoll")
        .ToListAsync();

    foreach (var trigger in triggers)
    {
        try
        {
            var config = JsonSerializer.Deserialize<ApiPollTriggerConfig>(
                trigger.ConfigJson);

            if (config == null) continue;

            // Check if it's time to poll
            if (trigger.LastCheckedAt.HasValue)
            {
                var timeSinceLastCheck = DateTime.UtcNow - trigger.LastCheckedAt.Value;
                if (timeSinceLastCheck.TotalSeconds < config.IntervalSeconds)
                    continue;
            }

            _logger.LogInformation("Polling API for workflow {WorkflowId}",
                trigger.WorkflowDefinition.WorkflowId);

            // Poll the API
            using var httpClient = new HttpClient();
            var response = await httpClient.GetAsync(config.Url);
            var content = await response.Content.ReadAsStringAsync();

            // Get previous state
            var state = JsonSerializer.Deserialize<Dictionary<string, object>>(
                            trigger.StateJson) ?? new();

            var previousHash = state.GetValueOrDefault("contentHash")?.ToString();
            var currentHash = ComputeHash(content);

            // Has content changed?
            if (previousHash != currentHash || config.AlwaysTrigger)
            {
                _logger.LogInformation(
                    "API content changed, triggering workflow {WorkflowId}",
                    trigger.WorkflowDefinition.WorkflowId);

                // Pass response as input to workflow
                var inputData = new Dictionary<string, object>
                {
                    ["apiResponse"] = content,
                    ["statusCode"] = (int)response.StatusCode,
                    ["previousHash"] = previousHash ?? string.Empty,
                    ["currentHash"] = currentHash
                };

                // Execute the workflow
                await _executionService.ExecuteWorkflowAsync(
                    trigger.WorkflowDefinition.WorkflowId,
                    inputData,
                    $"ApiPoll:{config.Url}");

                trigger.LastFiredAt = DateTime.UtcNow;
                trigger.FireCount++;

                // Update state
                state["contentHash"] = currentHash;
                state["lastContent"] = content.Length > 1000
                    ? content.Substring(0, 1000)
                    : content;
                state["lastPoll"] = DateTime.UtcNow.ToString("O");
            }

            trigger.LastCheckedAt = DateTime.UtcNow;
            trigger.StateJson = JsonSerializer.Serialize(state);
            trigger.LastError = null;

            await _context.SaveChangesAsync();
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "Error polling API trigger {TriggerId}",
                trigger.Id);
            trigger.LastError = ex.Message;
            trigger.LastCheckedAt = DateTime.UtcNow;
            await _context.SaveChangesAsync();
        }
    }
}

private string ComputeHash(string content)
{
    using var sha256 = System.Security.Cryptography.SHA256.Create();
    var bytes = System.Text.Encoding.UTF8.GetBytes(content);
    var hash = sha256.ComputeHash(bytes);
    return Convert.ToBase64String(hash);
}

Inställningsmodell

  1. Det rannsakande arbetet
  2. Hur det fungerar:
  3. Varje minut, kontrollera alla API-undersökning triggers
  4. För varje utlösare, kontrollera om tillräckligt med tid har gått sedan förra undersökningen
  5. Välj den inställda webbadressen@ info: whatsthis
  6. Beräkna en hash av svarsinnehålletAlwaysTriggerJämför med tidigare hasch som lagrats i tillstånd
  7. Om den ändras (eller
  8. är sant), kör arbetsflödet

Skicka API-svaret som indata till arbetsflödet

Uppdatera tillståndet med ny hash

{
  "triggerType": "ApiPoll",
  "config": {
    "url": "https://api.github.com/repos/dotnet/aspnetcore/releases/latest",
    "intervalSeconds": 3600,
    "alwaysTrigger": false
  }
}

Exempel på användningsfall

Övervaka GitHub-utgåvor:

Detta röstar GitHub API varje timme.Program.csNär en ny utgåva publiceras ändras innehållets hash och arbetsflödet utförs med utgivningsdata!

// Add Hangfire services
builder.Services.AddHangfire(config =>
{
    config.UsePostgreSqlStorage(
        builder.Configuration.GetConnectionString("DefaultConnection"));
});

builder.Services.AddHangfireServer();

// Register our job
builder.Services.AddScoped<WorkflowSchedulerJob>();

Registrera Hangfire-jobb

app.UseHangfireDashboard("/hangfire");

// Register recurring jobs
RecurringJob.AddOrUpdate<WorkflowSchedulerJob>(
    "scheduled-workflows",
    job => job.ExecuteScheduledWorkflowsAsync(),
    Cron.Minutely);

RecurringJob.AddOrUpdate<WorkflowSchedulerJob>(
    "api-poll-triggers",
    job => job.PollApiTriggersAsync(),
    Cron.Minutely);

I din

eller startinställning:/hangfire:

  • **Sedan, efter att appen startar, registrera återkommande jobb:**Hangfiretavlan
  • Hangfire inkluderar en inbyggd instrumentbrädan tillgänglig påArbetstillfällen
  • : Se alla köade, bearbetnings- och färdiga jobbÅterkommande arbeten
  • : Hantera våra arbetsflöde schemaläggareRestriktioner

: Visa och försök igen misslyckades jobb

app.UseHangfireDashboard("/hangfire", new DashboardOptions
{
    Authorization = new[]
    {
        new HangfireAuthorizationFilter()
    }
});

public class HangfireAuthorizationFilter : IDashboardAuthorizationFilter
{
    public bool Authorize(DashboardContext context)
    {
        var httpContext = context.GetHttpContext();

        // Only allow authenticated users
        return httpContext.User.Identity?.IsAuthenticated == true;
    }
}

Serverer

: Övervaka Hangfire servrar

[HttpPost("workflow/{id}/triggers")]
public async Task<IActionResult> CreateTrigger(
    string id,
    [FromBody] TriggerCreateRequest request)
{
    var workflow = await _context.WorkflowDefinitions
        .FirstOrDefaultAsync(w => w.WorkflowId == id);

    if (workflow == null)
        return NotFound();

    var trigger = new WorkflowTriggerStateEntity
    {
        WorkflowDefinitionId = workflow.Id,
        TriggerType = request.Type,
        ConfigJson = JsonSerializer.Serialize(request.Config),
        StateJson = "{}",
        IsEnabled = true
    };

    await _context.WorkflowTriggerStates.AddAsync(trigger);
    await _context.SaveChangesAsync();

    return Json(new { success = true, triggerId = trigger.Id });
}

public class TriggerCreateRequest
{
    public string Type { get; set; } = string.Empty; // Schedule, ApiPoll
    public object Config { get; set; } = new();
}

Att säkra brädan

<div class="card bg-base-100 shadow-xl">
    <div class="card-body">
        <h2 class="card-title">⏰ Add Trigger</h2>

        <div class="form-control">
            <label class="label">Trigger Type</label>
            <select class="select select-bordered" x-model="triggerType">
                <option value="Schedule">Schedule</option>
                <option value="ApiPoll">API Poll</option>
            </select>
        </div>

        <!-- Schedule Config -->
        <template x-if="triggerType === 'Schedule'">
            <div class="space-y-4">
                <div class="form-control">
                    <label class="label">Interval</label>
                    <div class="flex gap-2">
                        <input type="number"
                               x-model="scheduleConfig.intervalValue"
                               class="input input-bordered flex-1" />
                        <select x-model="scheduleConfig.intervalType"
                                class="select select-bordered">
                            <option value="minutes">Minutes</option>
                            <option value="hours">Hours</option>
                            <option value="days">Days</option>
                        </select>
                    </div>
                </div>
            </div>
        </template>

        <!-- API Poll Config -->
        <template x-if="triggerType === 'ApiPoll'">
            <div class="space-y-4">
                <div class="form-control">
                    <label class="label">API URL</label>
                    <input type="url"
                           x-model="apiConfig.url"
                           class="input input-bordered"
                           placeholder="https://api.example.com/data" />
                </div>

                <div class="form-control">
                    <label class="label">Poll Interval (seconds)</label>
                    <input type="number"
                           x-model="apiConfig.intervalSeconds"
                           class="input input-bordered"
                           value="300" />
                </div>
            </div>
        </template>

        <button @click="createTrigger()" class="btn btn-primary mt-4">
            Create Trigger
        </button>
    </div>
</div>

Hantera triggers via UI

Låt oss lägga till UI för att skapa och hantera triggers:

  1. UI- komponent
  2. Verkligt exempel på arbetsflöde
  3. Låt oss bygga ett komplett automatiserat arbetsflöde som:
  4. Undersökningar GitHub API för nya utgåvor

Kontrollerar om versionen är nyare än vad vi har sett

{
  "name": "GitHub Release Monitor",
  "startNodeId": "parse-data",
  "nodes": [
    {
      "id": "parse-data",
      "type": "Transform",
      "name": "Extract Version",
      "inputs": {
        "operation": "json_parse",
        "data": "{{apiResponse}}"
      }
    },
    {
      "id": "log-release",
      "type": "Log",
      "name": "Log New Release",
      "inputs": {
        "message": "New release: {{tag_name}} - {{name}}",
        "level": "info"
      }
    }
  ],
  "connections": [
    {
      "sourceNodeId": "parse-data",
      "targetNodeId": "log-release"
    }
  ]
}

Loggar ett meddelande

{
  "type": "ApiPoll",
  "config": {
    "url": "https://api.github.com/repos/dotnet/aspnetcore/releases/latest",
    "intervalSeconds": 3600
  }
}

(Kan skicka ett e-postmeddelande, posta till Slack, etc.)

  1. Steg 1: Skapa arbetsflödet
  2. Steg 2: Skapa API Polling Trigger
  3. Nu, varje timme, Hangfire kommer:
  4. Prövning av GitHub- API:et

Jämför innehållet hash med föregående undersökning

Om den ändras, kör arbetsflödet

Arbetsflödet tolkar JSON och loggar release info

_logger.LogInformation(
    "Workflow {WorkflowId} execution {ExecutionId} completed in {Duration}ms with status {Status}",
    execution.WorkflowId,
    execution.Id,
    execution.DurationMs,
    execution.Status);

Övervakning och observerbarhet

Loggning

private static readonly Counter WorkflowExecutions = Metrics
    .CreateCounter("workflow_executions_total",
        "Total workflow executions",
        new CounterConfiguration
        {
            LabelNames = new[] { "workflow_id", "status" }
        });

// In execution service
WorkflowExecutions
    .WithLabels(workflow.Id, execution.Status.ToString())
    .Inc();

Alla arbetsflöden loggas:

Tröskelvärden

  • Vi kan lägga till Prometheus-mått:
  • Registreringar
  • Inrätta registreringar för
  • Misslyckades med arbetsflöden (Status == Misslyckades)

Arbetsflödena tar för lång tid

Misslyckades med API- röstning

Utlösare som inte har skjutit i förväntad tidsram

Prestandaöverväganden

// Instead of querying per trigger
var triggers = await _context.WorkflowTriggerStates
    .Include(t => t.WorkflowDefinition)
    .Where(t => t.IsEnabled && t.TriggerType == "ApiPoll")
    .AsNoTracking() // Read-only
    .ToListAsync();

Databasbelastning

Med många arbetsflöden omröstningar ofta, kan databasbelastningen vara betydande:

Lösning: Batch-frågor

catch (HttpRequestException ex) when (ex.StatusCode == HttpStatusCode.TooManyRequests)
{
    // Back off
    var retryAfter = response.Headers.RetryAfter?.Delta ?? TimeSpan.FromMinutes(5);
    state["backoffUntil"] = DateTime.UtcNow.Add(retryAfter).ToString("O");
}

Begränsning av API-frekvens

Vid val av externa API:er:

Lösning: Exponentiell backoff

public class ConditionalTriggerConfig : ApiPollTriggerConfig
{
    public string? Condition { get; set; } // e.g., "{{stars}} > 1000"
}

Avancerade funktioner

Villkorliga triggers

// After workflow completes
if (execution.Status == WorkflowExecutionStatus.Completed)
{
    var dependentTriggers = await _context.WorkflowTriggerStates
        .Where(t => t.TriggerType == "WorkflowComplete" &&
                    t.ConfigJson.Contains(execution.WorkflowId))
        .ToListAsync();

    foreach (var trigger in dependentTriggers)
    {
        await _executionService.ExecuteWorkflowAsync(
            trigger.WorkflowDefinition.WorkflowId,
            execution.OutputData,
            $"Triggered by {execution.WorkflowId}");
    }
}

Utlöser endast om vissa villkor är uppfyllda:

Utlösande av beroenden

[Fact]
public async Task ExecuteScheduledWorkflows_ShouldExecuteWhenIntervalPassed()
{
    // Arrange
    var mockContext = CreateMockContext();
    var mockExecutionService = new Mock<IWorkflowExecutionService>();
    var job = new WorkflowSchedulerJob(mockContext.Object,
        mockExecutionService.Object, Mock.Of<ILogger>());

    // Act
    await job.ExecuteScheduledWorkflowsAsync();

    // Assert
    mockExecutionService.Verify(s => s.ExecuteWorkflowAsync(
        It.IsAny<string>(),
        It.IsAny<Dictionary<string, object>>(),
        "Scheduler",
        It.IsAny<CancellationToken>()), Times.Once);
}

Kedjeutlösare - ett arbetsflödes färdigställande utlöser ett annat:

Testa Hangfire-jobb

✅ **Enheten testar dina jobb:**Slutsatser ✅ **Vi har byggt ett komplett automationssystem!**Våra arbetsflöden kan nu: ✅ Kör på scheman- Varje timme, dagligen, eller anpassade intervall ✅ Prövning av API:er- Övervaka externa tjänster för förändringar ✅ Spårtillstånd- Kom ihåg vad vi har sett förut. ✅ Automatiskt försök- Hantera övergående fel

Övervaka

  • Dashboard för alla jobb

  • Skala- Hangfire hanterar lastbalansering

  • Hela serienVi har byggt ett arbetsflödessystem av företagsklass från grunden:

  • Häfte 1: Introduktion och arkitektur

  • Häfte 2: Huvudarbetsflödesmotor

Häfte 3

  • : Visual workflow editor
  • Häfte 4
  • : Integrering av hangfire (detta inlägg)
  • Du har nu följande:
  • En kraftfull arbetsflödesmotor

En vacker bildredigerare

Automatiskt genomförande

  • API-övervakningFullständig observerbarhet
  • **Vad är nästa?**Möjliga förbättringar:
  • Webbhooks: Trigger arbetsflöden via HTTP-slutpunkter
  • E- postnoder: Skicka e-post från arbetsflöden
  • Databasnoder: Frågedatabaser
  • AI-noder: Integrera med LLMs

Delarbetsflöden

: Komponera arbetsflöden tillsammans

  • Marknaden för arbetsflödenMostlylucid.SchedulerService/Jobs/
  • : Mallar för delade arbetsflödenMostlylucid.Workflow.Shared/
  • KällkodMostlylucid.Workflow.Engine/

Thank you for following this series! Happy workflow building! 🎉

Finding related posts...
logo

© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.