惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Jina AI
Jina AI
C
Cybersecurity and Infrastructure Security Agency CISA
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
T
Threat Research - Cisco Blogs
L
LINUX DO - 热门话题
Simon Willison's Weblog
Simon Willison's Weblog
L
Lohrmann on Cybersecurity
S
Schneier on Security
T
The Exploit Database - CXSecurity.com
Know Your Adversary
Know Your Adversary
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
Cyberwarzone
Cyberwarzone
T
Threatpost
Hugging Face - Blog
Hugging Face - Blog
博客园_首页
Scott Helme
Scott Helme
WordPress大学
WordPress大学
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
W
WeLiveSecurity
L
LINUX DO - 最新话题
G
GRAHAM CLULEY
酷 壳 – CoolShell
酷 壳 – CoolShell
S
SegmentFault 最新的问题
Vercel News
Vercel News
Microsoft Azure Blog
Microsoft Azure Blog
有赞技术团队
有赞技术团队
Cisco Talos Blog
Cisco Talos Blog
V2EX - 技术
V2EX - 技术
Apple Machine Learning Research
Apple Machine Learning Research
H
Help Net Security
F
Fortinet All Blogs
The Hacker News
The Hacker News
IT之家
IT之家
Forbes - Security
Forbes - Security
月光博客
月光博客
S
Security @ Cisco Blogs
SecWiki News
SecWiki News
博客园 - 聂微东
GbyAI
GbyAI
S
Security Affairs
H
Heimdal Security Blog
人人都是产品经理
人人都是产品经理
大猫的无限游戏
大猫的无限游戏
AWS News Blog
AWS News Blog
T
Tenable Blog
P
Privacy International News Feed
Microsoft Security Blog
Microsoft Security Blog
C
Cyber Attacks, Cyber Crime and Cyber Security
AI
AI

DEV Community

Authentication Security Deep Dive: From Brute Force to Salted Hashing (With Java Examples) Why AI Systems Don’t Fail — They Drift Spilling beans for how i learn for exam😁"Reinforcement Learning Cheat Sheet" I Replaced Chrome with Safari for AI Browser Automation. Here's What Broke (and What Finally Worked) How Python Borrows Other People's Work The $40 Architecture: Processing 1 Billion API Requests with 99.99% Uptime Vibe Coding: A Workflow Guide (From Zero to SaaS) Most webhook security guides protect the wrong side. The scary part is delivery. Headless CMS for TanStack Start: Build a Blog with Cosmic EU Age Verification App "Hacked in 2 Minutes" — What Actually Happened Comfy Cloud’s delete function does not actually remove files Running AI Models on GPU Cloud Servers: A Beginner Guide Event-driven media intelligence with AWS Step Functions and Bedrock I scored 500 AI prompts across 8 quality dimensions — here's what broke How to Call Google Gemini API from Next.js (Free Tier, No Backend Needed) The Portal Protocol: Reclaiming Human Connection in the Age of AI How to Fix Your Team's Scattered Knowledge Problem With a Self-Hosted Forum Intro to tc Cloud Functors: A Graph-First Mental Model for the Modern Cloud Designing Multi-Tenant Backends With Both Ownership and Team Access I Built a Neumorphic CSS Library with 77+ Components — Here's What I Learned PostgreSQL Performance Optimization: Why Connection Pooling Is Critical at Scale Cómo construí un SaaS multi-rubro para gestionar expensas en Argentina con FastAPI + Vue 3 🚀 I Built an Ethical Hacking Scanner Tool – Open Source Project I Replaced /usage and /context in Claude Code With a Single Statusline A Pythonic Way to Handle Emails (IMAP/SMTP) with Auto-Discovery and AI-Ready Design I Collected 8.9 Million Polymarket Price Points — Here's What I Found About How Markets Really Move EcoTrack AI — Carbon Footprint Tracker & Dashboard Everyone's Using AI. No One Agrees How. 5 self-hosted ebook managers worth trying in 2026 Building Your First AI Agent with LangChain: From Chatbot to Autonomous Assistant Common SOC 2 Failures (Real World) Stop Vibe-Checking Your AI App: A Practical Guide to Evals How to Use SonarQube and SonarScanner Locally to Level Up Your Code Quality Your Next To-Do App Is Dead — I Replaced Mine with an OpenClaw AI Sign a Nostr event in 60 lines of Python using coincurve — no nostr-sdk, no nbxplorer, no rust toolchain ITGC Audit Explained Like You’re in Big 4 Patch Tuesday abril 2026: Microsoft parcha 163 vulnerabilidades y un zero-day en SharePoint Stop scraping everything: a better way to track competitor price changes Listing on MCPize + the Official MCP Registry while routing payments OUTSIDE the marketplace — how I kept 100% of my x402 revenue Building an AI-Powered Risk Intelligence System Using Serverless Architecture Why We Ripped Function Overloading Out of Our AI Toolchain Testing AI-Generated Code: How to Actually Know If It Works SaaS Churn Is Killing Your Business. Here Is What to Do About It (Without a Support Team) The Speed of AI Is No Longer Linear - And Self-Improving Models Are Why How to Implement RBAC for MCP Tools: A Practical Guide for Engineering Teams From Standard Quote to Persuasive Proposal: AI Automation for Arborists I built a CLI that scaffolds complete multi-tenant SaaS apps Axios CVE-2025–62718: The Silent SSRF Bug That Could Be Hiding in Your Node.js App Right Now The dashboard that ended our friendship Data Pipelines Explained Simply (and How to Build Them with Python) The Hidden Cost of AI Systems Nobody Talks About. undefined vs undeclared, and how typeof behaves Switching from file-based jobs to NATS/Kafka in Rust without changing code io_uring Adventures: Rust Servers That Love Syscalls Why Agentic AI is Killing the Traditional Database The POUR principles of web accessibility for developers and designers Quantum Neural Network 3D — A Deep Dive into Interactive WebGL Visualization How To Install Caveman In Codex On macOS And Windows Automation Pipeline Reliability: Why Your Workflow Breaks When Nobody Is Watching I Built an 'Open World' AI Coding Agent — It Works From ANY Folder From Freelancing to Product: A Tech Service Company's SaaS Transformation China's AI Giants: Adding Tencent Hunyuan & ByteDance Doubao to AI University (74 Providers) On the Vibe Coders and Their Lies clerk: Auto-Summarize Your Claude Code Sessions AI Weekly — 2026/04/10–04/17 | The Model Lockdown Is Here, but the Toolchain Is the Real Battleground AI 週報 — 2026/04/10–2026/04/17 模型封鎖潮來了,但工具鏈才是真戰場 Maybe this is how Open-Source apps are born... 🚀 Fine-Tune LLMs with LoRA and QLoRA: 2026 Guide tRPC v11 + Next.js App Router: End-to-End Type Safety Without the Boilerplate ShadCN UI in 2026: Why I Stopped Installing Component Libraries and Started Owning My Components SaaS Billing in React Server Components: Stripe + Supabase Without a Single `useEffect` Join our DEV Weekend Challenge — $1,000 in Prizes Across TEN winners! Submissions Due April 20 at 6:59 AM UTC. Implementing FSRS Spaced Repetition in Flutter + Supabase — Adding Memory Science to an AI Learning App "I Texted My Localhost From the Train — Claude Code Fixed the Bug Before I Got Home" I Built a Sales Prep AI and It Went Deeper Than Expected Design to Code #2: One JSON, Eleven Outputs Solving the 100M-Row Problem: A Summary Table Pattern for High-Volume Push Notification Logs Flutter Web With Wasm: What Actually Changes For Developers I Built 50 Royalty-Free Soundtracks for My Side Project in a Weekend Using AI Music Generation The Vibe Coding Security Checklist: 7 Things to Check Before You Ship Stop Letting Googlebot Guess Fix Your React App's SEO Right Desconstruindo o Streaming do LinkedIn: Como Criar um Engine de Extração de Vídeo de Alta Performance com HLS e FFmpeg (EDA Part-1) EDA (Exploratory Data Analysis) Explained With Real Life — Why Looking at Your Data Is the Most Important Step in Machine Learning Brand Relationship Management at Scale: Our 4-Touch Outreach System for 200+ Brands Why String.fromEnvironment() Might Return an Empty String in Dart JGuardrails 1.0.0 — Hardening Java LLM Apps Against Jailbreaks, Toxicity, and Prompt Injection Plan and Schedule a Full Week of Threads Content From One Claude Conversation Coding Cat Oran Ep3, Five Tables Changed Everything Updated: BFF Pattern I'm done watching freelancers get buried by 200 proposals. So I'm building the alternative. This is my first post BFS Algorithm in Java Step by Step Tutorial with Examples Tracking LLM Pricing Monthly: An Open Dataset for 22 AI Models How We Measure Content ROI on a Comparison Site: Revenue Attribution Without Perfect Data Introducing Nova AI Ops: The AI-Native Operating System for SRE Teams I built a free desktop video downloader for Windows — Grabbit How Talkie OCR Helps Vision-Impaired & Dyslexic Users Read the World Around Them VRCFaceTracking安装和iPhone面捕配置教程,有bug Even CrowdStrike Can't See Your Agents The Automation Gold Rush: What n8n Workflows and Claude Are Opening Up for Developers Right Now
Data Consistency Patterns in Distributed Systems and .NET Core
Hossein Esmati · 2026-06-26 · via DEV Community

Hossein Esmati

This article is part of the Comprehensive Guide to Microservices Architecture in .NET Core, Cloud and Azure series.

Comparison of Patterns

Pattern Complexity Consistency Performance Use Case
Transactional Outbox Medium Strong Good General-purpose, reliable event publishing
CDC Low Strong Excellent Existing systems, minimal code changes
Event Sourcing High Strong Excellent Audit requirements, temporal queries
Saga Pattern High Eventual Good Complex distributed workflows

Best Practices and Decision Tree

For Transactional Outbox:

  • Use background services with proper error handling and retry logic
  • Implement idempotency in message consumers
  • Monitor outbox table size and clean up processed messages
  • Consider partitioning outbox table for high-volume scenarios

For CDC:

  • Regularly monitor change tracking overhead on the database
  • Configure appropriate retention periods
  • Handle schema changes carefully to avoid breaking CDC
  • Test CDC processor recovery after failures

For All Patterns:

  • Implement comprehensive logging and monitoring
  • Use distributed tracing to track operations across services
  • Design for idempotency at every level
  • Plan for failure scenarios and compensating actions
  • Consider using Azure Monitor and Application Insights for observability

Azure-Specific Considerations

When implementing these patterns on Azure:

  • Use Azure SQL Database with built-in CDC support
  • Leverage Azure Service Bus for reliable message delivery with features like dead-letter queues and scheduled messages
  • Implement Azure Functions as lightweight CDC processors for serverless scenarios
  • Use Azure Cosmos DB change feed as an alternative to CDC for NoSQL scenarios
  • Enable Application Insights for end-to-end transaction tracking
  • Consider Azure Durable Functions for implementing saga patterns with built-in state management

The Anti-Pattern

// ANTI-PATTERN: Dual write (not atomic)
public async Task CreateOrderAsync(Order order)
{
    await _database.SaveAsync(order); // Write 1
    await _serviceBus.PublishAsync(new OrderCreated(order.Id)); // Write 2

    // If the publish fails after the database save succeeds,
    // we have inconsistency - the order exists but no event was published!
}

This approach has several critical issues:

  • No atomicity between database write and message publication
  • Partial failures leave the system in an inconsistent state
  • Manual compensation logic is complex and error-prone
  • Difficult to recover from failures

Solution 1: Transactional Outbox Pattern

The Transactional Outbox pattern ensures atomicity by storing events in the same database transaction as the business data, then publishing them asynchronously. You can read more about distributed transactions in .NET core and Azure in here.

Implementation with EF Core 9 and Azure Service Bus

// Outbox message entity
public class OutboxMessage
{
    public Guid Id { get; set; }
    public string EventType { get; set; }
    public string Payload { get; set; }
    public DateTime CreatedAt { get; set; }
    public DateTime? ProcessedAt { get; set; }
    public int RetryCount { get; set; }
    public string? Error { get; set; }
}

// DbContext with outbox
public class AppDbContext : DbContext
{
    public DbSet<Order> Orders { get; set; }
    public DbSet<OutboxMessage> OutboxMessages { get; set; }

    public AppDbContext(DbContextOptions<AppDbContext> options) 
        : base(options)
    {
    }

    protected override void OnModelCreating(ModelBuilder modelBuilder)
    {
        modelBuilder.Entity<OutboxMessage>()
            .HasIndex(m => new { m.ProcessedAt, m.CreatedAt })
            .HasFilter("[ProcessedAt] IS NULL");
    }
}

// Service with transactional outbox
public class OrderService
{
    private readonly AppDbContext _context;
    private readonly ILogger<OrderService> _logger;

    public OrderService(AppDbContext context, ILogger<OrderService> logger)
    {
        _context = context;
        _logger = logger;
    }

    public async Task<Guid> CreateOrderAsync(CreateOrderRequest request)
    {
        // Use ExecuteInTransactionAsync for .NET 9
        await _context.Database.CreateExecutionStrategy().ExecuteInTransactionAsync(
            async () =>
            {
                var order = new Order
                {
                    Id = Guid.NewGuid(),
                    CustomerId = request.CustomerId,
                    TotalAmount = request.Items.Sum(i => i.Quantity * i.Price),
                    Status = OrderStatus.Created,
                    CreatedAt = DateTime.UtcNow
                };

                _context.Orders.Add(order);

                // Store event in outbox within the same transaction
                var outboxMessage = new OutboxMessage
                {
                    Id = Guid.NewGuid(),
                    EventType = nameof(OrderCreatedEvent),
                    Payload = JsonSerializer.Serialize(new OrderCreatedEvent
                    {
                        OrderId = order.Id,
                        CustomerId = order.CustomerId,
                        CreatedAt = order.CreatedAt
                    }),
                    CreatedAt = DateTime.UtcNow
                };

                _context.OutboxMessages.Add(outboxMessage);

                await _context.SaveChangesAsync();

                _logger.LogInformation(
                    "Order {OrderId} created with outbox message {MessageId}",
                    order.Id, outboxMessage.Id);

                return order.Id;
            },
            verifySucceeded: null);

        return await Task.FromResult(Guid.Empty); // This will be set by the transaction
    }
}

// Outbox processor with Azure Service Bus
public class OutboxProcessor : BackgroundService
{
    private readonly IServiceProvider _serviceProvider;
    private readonly ServiceBusSender _serviceBusSender;
    private readonly ILogger<OutboxProcessor> _logger;
    private const int BatchSize = 100;
    private const int MaxRetries = 5;

    public OutboxProcessor(
        IServiceProvider serviceProvider,
        ServiceBusClient serviceBusClient,
        ILogger<OutboxProcessor> logger)
    {
        _serviceProvider = serviceProvider;
        _serviceBusSender = serviceBusClient.CreateSender("order-events");
        _logger = logger;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("Outbox processor started");

        while (!stoppingToken.IsCancellationRequested)
        {
            try
            {
                await ProcessOutboxMessagesAsync(stoppingToken);
                await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Error processing outbox messages");
                await Task.Delay(TimeSpan.FromSeconds(30), stoppingToken);
            }
        }
    }

    private async Task ProcessOutboxMessagesAsync(CancellationToken cancellationToken)
    {
        using var scope = _serviceProvider.CreateScope();
        var context = scope.ServiceProvider.GetRequiredService<AppDbContext>();

        // Fetch unprocessed messages
        var messages = await context.OutboxMessages
            .Where(m => m.ProcessedAt == null && m.RetryCount < MaxRetries)
            .OrderBy(m => m.CreatedAt)
            .Take(BatchSize)
            .ToListAsync(cancellationToken);

        if (!messages.Any())
            return;

        _logger.LogInformation("Processing {Count} outbox messages", messages.Count);

        foreach (var message in messages)
        {
            try
            {
                // Publish to Azure Service Bus
                var serviceBusMessage = new ServiceBusMessage(message.Payload)
                {
                    MessageId = message.Id.ToString(),
                    Subject = message.EventType,
                    ContentType = "application/json"
                };

                await _serviceBusSender.SendMessageAsync(serviceBusMessage, cancellationToken);

                // Mark as processed
                message.ProcessedAt = DateTime.UtcNow;

                _logger.LogInformation(
                    "Outbox message {MessageId} published successfully",
                    message.Id);
            }
            catch (Exception ex)
            {
                message.RetryCount++;
                message.Error = ex.Message;

                _logger.LogWarning(ex,
                    "Failed to process outbox message {MessageId}. Retry count: {RetryCount}",
                    message.Id, message.RetryCount);
            }
        }

        await context.SaveChangesAsync(cancellationToken);
    }

    public override async Task StopAsync(CancellationToken cancellationToken)
    {
        _logger.LogInformation("Outbox processor stopping");
        await _serviceBusSender.CloseAsync();
        await base.StopAsync(cancellationToken);
    }
}

// Register services in Program.cs
builder.Services.AddDbContext<AppDbContext>(options =>
    options.UseSqlServer(builder.Configuration.GetConnectionString("DefaultConnection")));

builder.Services.AddSingleton(sp =>
    new ServiceBusClient(builder.Configuration.GetConnectionString("ServiceBus")));

builder.Services.AddHostedService<OutboxProcessor>();

Solution 2: Change Data Capture (CDC)

Change Data Capture automatically tracks changes in your database and publishes them as events, eliminating the need for application-level dual writes.

Implementation with SQL Server CDC and Azure

// Enable CDC on SQL Server
// Execute these SQL commands on your Azure SQL Database:
/*
ALTER DATABASE YourDatabase SET CHANGE_TRACKING = ON  
    (CHANGE_RETENTION = 2 DAYS, AUTO_CLEANUP = ON);

ALTER TABLE Orders ENABLE CHANGE_TRACKING  
    WITH (TRACK_COLUMNS_UPDATED = ON);
*/

// Change tracking model
public class OrderChange
{
    public Guid OrderId { get; set; }
    public Guid CustomerId { get; set; }
    public string Operation { get; set; } // I, U, D
    public long ChangeVersion { get; set; }
}

// CDC Processor with .NET 9
public class CdcProcessor : BackgroundService
{
    private readonly IServiceProvider _serviceProvider;
    private readonly ServiceBusSender _serviceBusSender;
    private readonly ILogger<CdcProcessor> _logger;
    private long _lastSyncVersion;

    public CdcProcessor(
        IServiceProvider serviceProvider,
        ServiceBusClient serviceBusClient,
        ILogger<CdcProcessor> logger)
    {
        _serviceProvider = serviceProvider;
        _serviceBusSender = serviceBusClient.CreateSender("order-events");
        _logger = logger;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("CDC processor started");

        // Initialize last sync version
        _lastSyncVersion = await GetCurrentChangeVersionAsync();

        while (!stoppingToken.IsCancellationRequested)
        {
            try
            {
                await ProcessChangesAsync(stoppingToken);
                await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Error processing CDC changes");
                await Task.Delay(TimeSpan.FromSeconds(30), stoppingToken);
            }
        }
    }

    private async Task ProcessChangesAsync(CancellationToken cancellationToken)
    {
        using var scope = _serviceProvider.CreateScope();
        var context = scope.ServiceProvider.GetRequiredService<AppDbContext>();

        // Query change tracking using FormattableString for safety
        var currentVersion = await GetCurrentChangeVersionAsync();

        var changes = await context.Database
            .SqlQuery<OrderChange>(
                $@"SELECT o.OrderId, o.CustomerId, 
                      CT.SYS_CHANGE_OPERATION as Operation,
                      CT.SYS_CHANGE_VERSION as ChangeVersion
                   FROM Orders o
                   RIGHT OUTER JOIN CHANGETABLE(CHANGES Orders, {_lastSyncVersion}) AS CT
                      ON o.OrderId = CT.OrderId
                   WHERE CT.SYS_CHANGE_VERSION <= {currentVersion}
                   ORDER BY CT.SYS_CHANGE_VERSION")
            .ToListAsync(cancellationToken);

        if (!changes.Any())
            return;

        _logger.LogInformation("Processing {Count} changes", changes.Count);

        foreach (var change in changes)
        {
            try
            {
                var eventMessage = change.Operation switch
                {
                    "I" => new ServiceBusMessage(JsonSerializer.Serialize(
                        new OrderCreatedEvent
                        {
                            OrderId = change.OrderId,
                            CustomerId = change.CustomerId
                        }))
                    {
                        Subject = nameof(OrderCreatedEvent)
                    },

                    "U" => new ServiceBusMessage(JsonSerializer.Serialize(
                        new OrderUpdatedEvent
                        {
                            OrderId = change.OrderId
                        }))
                    {
                        Subject = nameof(OrderUpdatedEvent)
                    },

                    "D" => new ServiceBusMessage(JsonSerializer.Serialize(
                        new OrderDeletedEvent
                        {
                            OrderId = change.OrderId
                        }))
                    {
                        Subject = nameof(OrderDeletedEvent)
                    },

                    _ => throw new InvalidOperationException(
                        $"Unknown operation: {change.Operation}")
                };

                eventMessage.MessageId = Guid.NewGuid().ToString();
                eventMessage.ContentType = "application/json";

                await _serviceBusSender.SendMessageAsync(eventMessage, cancellationToken);

                _logger.LogInformation(
                    "Published {EventType} for order {OrderId}",
                    eventMessage.Subject, change.OrderId);

                // Update last sync version
                _lastSyncVersion = Math.Max(_lastSyncVersion, change.ChangeVersion);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex,
                    "Failed to publish event for order {OrderId}",
                    change.OrderId);
            }
        }
    }

    private async Task<long> GetCurrentChangeVersionAsync()
    {
        using var scope = _serviceProvider.CreateScope();
        var context = scope.ServiceProvider.GetRequiredService<AppDbContext>();

        var version = await context.Database
            .SqlQuery<long>($"SELECT CHANGE_TRACKING_CURRENT_VERSION()")
            .FirstOrDefaultAsync();

        return version;
    }

    public override async Task StopAsync(CancellationToken cancellationToken)
    {
        _logger.LogInformation("CDC processor stopping");
        await _serviceBusSender.CloseAsync();
        await base.StopAsync(cancellationToken);
    }
}

Solution 3: Event Sourcing

Event Sourcing naturally solves the dual write problem by treating events as the single source of truth. See the dedicated Event Sourcing article for a complete implementation.

Solution 4: Saga Pattern for Distributed Transactions

For complex workflows spanning multiple services, the Saga pattern coordinates distributed transactions through compensating actions.

// Saga state machine using MassTransit
public class OrderSaga : MassTransitStateMachine<OrderSagaState>
{
    public State OrderCreated { get; private set; }
    public State PaymentProcessing { get; private set; }
    public State InventoryReserving { get; private set; }
    public State OrderCompleted { get; private set; }
    public State OrderFailed { get; private set; }

    public Event<OrderSubmitted> OrderSubmitted { get; private set; }
    public Event<PaymentSucceeded> PaymentSucceeded { get; private set; }
    public Event<PaymentFailed> PaymentFailed { get; private set; }
    public Event<InventoryReserved> InventoryReserved { get; private set; }
    public Event<InventoryReservationFailed> InventoryReservationFailed { get; private set; }

    public OrderSaga()
    {
        InstanceState(x => x.CurrentState);

        Event(() => OrderSubmitted, x => x.CorrelateById(m => m.Message.OrderId));
        Event(() => PaymentSucceeded, x => x.CorrelateById(m => m.Message.OrderId));
        Event(() => PaymentFailed, x => x.CorrelateById(m => m.Message.OrderId));
        Event(() => InventoryReserved, x => x.CorrelateById(m => m.Message.OrderId));
        Event(() => InventoryReservationFailed, x => x.CorrelateById(m => m.Message.OrderId));

        Initially(
            When(OrderSubmitted)
                .Then(context =>
                {
                    context.Saga.OrderId = context.Message.OrderId;
                    context.Saga.CustomerId = context.Message.CustomerId;
                })
                .TransitionTo(OrderCreated)
                .PublishAsync(context => context.Init<ProcessPayment>(new
                {
                    context.Message.OrderId,
                    context.Message.Amount
                })));

        During(OrderCreated,
            When(PaymentSucceeded)
                .TransitionTo(PaymentProcessing)
                .PublishAsync(context => context.Init<ReserveInventory>(new
                {
                    context.Message.OrderId,
                    context.Saga.CustomerId
                })),

            When(PaymentFailed)
                .TransitionTo(OrderFailed)
                .ThenAsync(async context =>
                {
                    // Compensate: Cancel order
                    await context.PublishAsync(new OrderCancelled
                    {
                        OrderId = context.Message.OrderId,
                        Reason = "Payment failed"
                    });
                }));

        During(PaymentProcessing,
            When(InventoryReserved)
                .TransitionTo(OrderCompleted)
                .Finalize(),

            When(InventoryReservationFailed)
                .TransitionTo(OrderFailed)
                .ThenAsync(async context =>
                {
                    // Compensate: Refund payment
                    await context.PublishAsync(new RefundPayment
                    {
                        OrderId = context.Message.OrderId
                    });
                })
                .ThenAsync(async context =>
                {
                    // Cancel order
                    await context.PublishAsync(new OrderCancelled
                    {
                        OrderId = context.Message.OrderId,
                        Reason = "Inventory unavailable"
                    });
                }));
    }
}

public class OrderSagaState : SagaStateMachineInstance
{
    public Guid CorrelationId { get; set; }
    public string CurrentState { get; set; }
    public Guid OrderId { get; set; }
    public Guid CustomerId { get; set; }
}

// Configure in Program.cs
builder.Services.AddMassTransit(x =>
{
    x.AddSagaStateMachine<OrderSaga, OrderSagaState>()
        .EntityFrameworkRepository(r =>
        {
            r.ExistingDbContext<AppDbContext>();
            r.UseSqlServer();
        });

    x.UsingAzureServiceBus((context, cfg) =>
    {
        cfg.Host(builder.Configuration.GetConnectionString("ServiceBus"));
        cfg.ConfigureEndpoints(context);
    });
});