Back to "Construction d'un système de flux de travail avec HTMX et ASP.NET Core - Partie 4: Intégration et automatisation du feu de hang"

This is a viewer only at the moment see the article on how this works.

To update the preview hit Ctrl-Alt-R (or ⌘-Alt-R on Mac) or Enter to refresh. The Save icon lets you save the markdown file to disk

This is a preview from the server running through my markdig pipeline

ASP.NET Background Jobs Hangfire Workflow

Construction d'un système de flux de travail avec HTMX et ASP.NET Core - Partie 4: Intégration et automatisation du feu de hang

Wednesday, 15 January 2025

Présentation

DansTroisième partie, nous avons construit un magnifique éditeur visuel.

  • **Mais nos workflows ne s'exécutent que lorsque nous les déclenchons manuellement.**Dans ce dernier post, nous allons rendre les workflows vraiment autonomes en utilisant Hangfire pour:
  • Exécution prévue- Exécuter les workflows selon un calendrier
  • Sondage sur l'API- Surveiller les API externes et déclencher les changements
  • Gestion de l ' État- Suivre les états déclencheurs à travers les exécutions

Tableau de bord

  • Surveiller tous les emplois d'arrière-plan

  • Pourquoi le feu ?

  • Hangfire est parfait pour nos besoins car il:

  • Offres d'emploi dans notre base de données PostgreSQL

  • Fournit un tableau de bord intégré

  • Soutient les emplois récurrents

A une logique de ré-essai automatique

Échelles horizontales

[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; }
}

Le modèle d'État déclencheur

  • Tout d'abord, comprenons notre entité d'état déclencheur (nous l'avons déjà créé dans la deuxième partie) :
  • Cette entité suit tout au sujet d'un déclencheur de workflow :
  • Quand il a couru pour la dernière fois
  • Quelle est sa configuration?

Dans quel état il est (pour les déclencheurs d'état)

Toute erreur qui s'est produite

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; }
}

Flux de travail programmés

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
        };
    }
}

Modèle de configuration

  1. L'emploi de programmeurExecuteScheduledWorkflowsAsync()
  2. Comment ça marche :
  3. Chaque minute, Hangfire appelle
  4. Nous demandons pour les déclencheurs de calendrier activés
  5. Pour chaque déclencheur, vérifiez si suffisamment de temps est passé

Si oui, exécutez le workflow

Mettre à jour l'état du déclencheur avec le dernier temps d'exécution

Sondage sur l'API

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; }
}

Le sondage d'API est plus intéressant - nous surveillons les API externes et déclenchons des flux de travail lorsque le contenu change!

[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);
}

Modèle de configuration

  1. Le travail électoral
  2. Comment ça marche :
  3. Chaque minute, vérifiez tous les déclencheurs du sondage API
  4. Pour chaque déclencheur, vérifiez si suffisamment de temps s'est écoulé depuis le dernier sondage.
  5. Sondage sur l'URL configurée
  6. Calculer un hash du contenu de la réponseAlwaysTriggerComparer avec le hash précédent stocké dans l'état
  7. En cas de changement (ou
  8. est vrai), exécutez le workflow

Passer la réponse de l'API comme données d'entrée au flux de travail

Mettre à jour l'état avec un nouveau hash

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

Exemple de cas d'utilisation

Surveiller les rejets de GitHub :

Ce sondage de l'API GitHub toutes les heures.Program.csQuand une nouvelle version est publiée, le hash de contenu change, et le workflow s'exécute avec les données de publication !

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

builder.Services.AddHangfireServer();

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

Enregistrement des emplois Hangfire

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);

Dans votre

ou configuration de démarrage:/hangfire:

  • **Puis, après le démarrage de l'application, enregistrez des emplois récurrents:**Le tableau de bord Hangfire
  • Hangfire comprend un tableau de bord intégré accessible àEmplois
  • : Voir tous les travaux en file d'attente, de traitement et complétésEmplois récurrents
  • : Gérer nos planificateurs de workflowDemandes de remboursement

: Affichage et réessayer des emplois échoués

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;
    }
}

Serveurs

: Surveiller les serveurs Hangfire

[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();
}

Sécuriser le tableau de bord

<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>

Gestion des déclencheurs via l'interface utilisateur

Ajoutons l'interface utilisateur pour créer et gérer les déclencheurs :

  1. Composante de l'assurance-chômage
  2. Exemple de flux de travail dans le monde réel
  3. Construisons un workflow automatisé complet qui :
  4. Sondages sur l'API GitHub pour les nouvelles versions

Vérifie si la version est plus récente que ce que nous avons vu

{
  "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"
    }
  ]
}

Loge un message

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

(Pourrait envoyer un courriel, envoyer un message à Slack, etc.)

  1. Étape 1: Créer le flux de travail
  2. Étape 2: Créer le déclencheur de sondage d'API
  3. Maintenant, chaque heure, Hangfire:
  4. Sondage sur l'API GitHub

Comparer le hash contenu avec le sondage précédent

En cas de modification, exécutez le workflow

Le flux de travail analyse le JSON et enregistre les informations de sortie

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

Surveillance et observation

Exploitation forestière

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();

Toutes les exécutions de workflow sont enregistrées :

métriques

  • Nous pouvons ajouter des métriques Prométheus:
  • Alertes
  • Mettre en place des alertes pour :
  • Défaillance des workflows (Statut == Échec)

Les flux de travail prennent trop de temps

Défauts de vote de l'API

Déclencheurs qui n'ont pas tiré dans les délais prévus

Considérations de performance

// 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();

Charge de la base de données

Avec de nombreux workflows de sondages fréquemment, la charge de base de données peut être importante:

Solution : requêtes par lots

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");
}

Limite du taux d'API

Lors d'un sondage sur les API externes :

Solution: Décollage exponentiel

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

Caractéristiques avancées

Déclencheurs conditionnels

// 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}");
    }
}

Ne déclencher que si certaines conditions sont remplies:

Dépendances de déclenchement

[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);
}

Déclenchements en chaîne - l'achèvement d'un workflow en déclenche un autre :

Essais d'emplois en feu de bois

✅ **Testez vos tâches à l'unité :**Le présent règlement entre en vigueur le vingtième jour suivant celui de sa publication au Journal officiel de l'Union européenne. ✅ **Nous avons construit un système d'automatisation complet !**Nos workflows peuvent maintenant : ✅ Exécuter sur les horaires- Intervalles horaires, quotidiens ou personnalisés ✅ API de sondage- Surveiller les services externes pour les changements ✅ État de la piste- Souviens-toi de ce qu'on a déjà vu. ✅ Auto-réessayer- Manipulation des défaillances transitoires

Moniteur

  • Tableau de bord pour tous les emplois

  • Échelle- Poignées Hangfire équilibre de charge

  • La série complèteNous avons construit un système de flux de travail de qualité entreprise à partir de zéro:

  • Première partie: Introduction et architecture

  • Deuxième partie: Moteur de flux de travail de base

Troisième partie

  • : Éditeur visuel de flux de travail
  • Quatrième partie
  • : Intégration Hangfire (ce post)
  • Vous avez maintenant :
  • Un puissant moteur de flux de travail

Un magnifique éditeur visuel

Exécution automatisée

  • Surveillance de l'APIPleine observabilité
  • **Qu'est-ce qu'il y a ?**Améliorations possibles:
  • Hooks sur le Web: Déclencher les workflows via les paramètres HTTP
  • Nœuds d'email: Envoyer des e-mails à partir de workflows
  • Nœuds de la base de données: Bases de données de requêtes
  • Nœuds AI: Intégrer avec les LLM

Sous-débits de travail

: Compiler les workflows ensemble

  • Marché des flux de travailMostlylucid.SchedulerService/Jobs/
  • : Partager des modèles de flux de travailMostlylucid.Workflow.Shared/
  • Code sourceMostlylucid.Workflow.Engine/

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

logo

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