スキル一覧に戻る
erymuzuan

messaging-events

by erymuzuan

Motorbike rental system for Thailand tourist areas (Phuket, Krabi, etc.) - Blazor Server + WASM PWA with MudBlazor

0🍴 0📅 2026年1月24日
GitHubで見るManusで実行

SKILL.md


name: messaging-events description: RabbitMQ pub/sub patterns for asynchronous processing of entity changes and system events.

Messaging & Events (RabbitMQ)

RabbitMQ pub/sub patterns from rx-erp for async processing.

Overview

ComponentDescription
ExchangeTopic exchange for routing
QueueMessage destination
Routing Key{Entity}.{Crud}.{Operation}
MessageBrokeredMessage with Entity payload

Message Flow

┌─────────────────┐     ┌─────────────────┐     ┌─────────────────┐
│  SubmitChanges  │────>│  RabbitMQ       │────>│  Subscriber     │
│  ("CheckIn")    │     │  Exchange       │     │  (Handler)      │
└─────────────────┘     └─────────────────┘     └─────────────────┘
        │                       │                       │
        ▼                       ▼                       ▼
  Rental.Changed.CheckIn   Route by key      Update inventory,
                                             Send notification

BrokeredMessage

public class BrokeredMessage
{
    public string RoutingKey => $"{Entity}.{Crud}.{Operation}";

    public Entity? Item { get; set; }
    public string? Operation { get; set; }
    public CrudOperation Crud { get; set; }
    public int? TryCount { get; set; }
    public string? Username { get; set; }
    public Dictionary<string, object> Headers { get; } = new();
    public TimeSpan RetryDelay { get; set; }

    public void Accept() => m_acknowledge(this, MessageReceiveStatus.Accepted);
    public void Reject() => m_acknowledge(this, MessageReceiveStatus.Rejected);
    public void Delay(TimeSpan ttl) { /* retry with delay */ }
}

public enum CrudOperation
{
    Added,
    Changed,
    Deleted
}

Publishing Messages

Messages are automatically published when calling SubmitChanges:

using var session = context.OpenSession();
session.Attach(rental);
await session.SubmitChanges("CheckIn");
// Publishes: Rental.Changed.CheckIn

Routing Key Format

ExampleEntityCrudOperation
Rental.Changed.CheckInRentalChangedCheckIn
Rental.Changed.CheckOutRentalChangedCheckOut
Motorbike.Changed.StatusUpdateMotorbikeChangedStatusUpdate
Payment.Added.RentalPaymentAddedRental
Renter.Added.RegistrationRenterAddedRegistration

Subscriber Base Pattern

public abstract class Subscriber<T> : Subscriber where T : Entity
{
    public abstract override string QueueName { get; }
    public abstract override string[] RoutingKeys { get; }

    protected abstract Task ProcessMessage(T item, BrokeredMessage message);

    public override void Run(IMessageBroker broker)
    {
        broker.OnMessageDeliveredAsync(async message =>
        {
            try
            {
                var item = message.Item as T;
                await ProcessMessage(item!, message);
                return MessageReceiveStatus.Accepted;
            }
            catch (Exception ex)
            {
                Logger.LogError(ex, "Failed to process message");
                return MessageReceiveStatus.Rejected;
            }
        }, new SubscriberOption
        {
            QueueName = QueueName,
            RoutingKeys = RoutingKeys
        });
    }
}

Example Subscribers

Rental Check-Out Handler

public class RentalCheckOutSubscriber : Subscriber<Rental>
{
    public override string QueueName => nameof(RentalCheckOutSubscriber);
    public override string[] RoutingKeys => [$"{nameof(Rental)}.{CrudOperation.Changed}.CheckOut"];

    protected override async Task ProcessMessage(Rental rental, BrokeredMessage message)
    {
        var context = new RentalDataContext();

        // Update motorbike status back to Available
        var motorbike = await context.LoadOneAsync<Motorbike>(
            m => m.MotorbikeId == rental.MotorbikeId);

        motorbike!.Status = "Available";

        using var session = context.OpenSession();
        session.Attach(motorbike);
        await session.SubmitChanges("StatusUpdate");

        // Send thank you notification
        await SendThankYouEmail(rental);

        message.Accept();
    }

    private async Task SendThankYouEmail(Rental rental)
    {
        // Email service call
    }
}

Expiry Warning Handler

public class RentalExpirySubscriber : Subscriber<Rental>
{
    public override string QueueName => nameof(RentalExpirySubscriber);
    public override string[] RoutingKeys => [$"{nameof(Rental)}.{CrudOperation.Changed}.ExpiryCheck"];

    protected override async Task ProcessMessage(Rental rental, BrokeredMessage message)
    {
        if (rental.Status != "Active")
        {
            message.Accept();
            return;
        }

        var daysRemaining = (rental.ExpectedEndDate - DateTimeOffset.Now).TotalDays;

        if (daysRemaining <= 1)
        {
            await SendExpiryWarning(rental);
        }

        message.Accept();
    }
}

Damage Report Handler

public class DamageReportSubscriber : Subscriber<DamageReport>
{
    public override string QueueName => nameof(DamageReportSubscriber);
    public override string[] RoutingKeys => [$"{nameof(DamageReport)}.{CrudOperation.Added}.*"];

    protected override async Task ProcessMessage(DamageReport damage, BrokeredMessage message)
    {
        // Notify shop owner
        await NotifyShopOwner(damage);

        // If major damage, flag motorbike for maintenance
        if (damage.Severity == "Major")
        {
            var context = new RentalDataContext();
            var motorbike = await context.LoadOneAsync<Motorbike>(
                m => m.MotorbikeId == damage.MotorbikeId);

            motorbike!.Status = "Maintenance";
            motorbike.Notes = $"Major damage reported: {damage.Description}";

            using var session = context.OpenSession();
            session.Attach(motorbike);
            await session.SubmitChanges("MaintenanceFlag");
        }

        message.Accept();
    }
}

Subscriber Registration

// Program.cs or HostedService
public class SubscriberHostedService : BackgroundService
{
    private readonly IMessageBroker m_broker;
    private readonly IServiceProvider m_services;

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        var subscribers = new Subscriber[]
        {
            new RentalCheckOutSubscriber(),
            new RentalExpirySubscriber(),
            new DamageReportSubscriber()
        };

        foreach (var subscriber in subscribers)
        {
            subscriber.Run(m_broker);
        }

        await Task.Delay(Timeout.Infinite, stoppingToken);
    }
}

Error Handling & Retry

protected override async Task ProcessMessage(Rental rental, BrokeredMessage message)
{
    try
    {
        await ProcessRentalAsync(rental);
        message.Accept();
    }
    catch (TransientException ex)
    {
        // Retry with delay
        if (message.TryCount < 3)
        {
            message.Delay(TimeSpan.FromMinutes(5));
        }
        else
        {
            Logger.LogError(ex, "Max retries exceeded");
            message.Reject();  // Move to dead letter queue
        }
    }
    catch (Exception ex)
    {
        Logger.LogError(ex, "Unrecoverable error");
        message.Reject();
    }
}

RabbitMQ Configuration

// appsettings.json
{
  "RabbitMQ": {
    "Host": "localhost",
    "Port": 5672,
    "Username": "guest",
    "Password": "guest",
    "VirtualHost": "/",
    "Exchange": "motorent.events"
  }
}

Common Use Cases

EventSubscriber Action
Rental.CheckInSend confirmation SMS
Rental.CheckOutSend receipt email
Rental.ExpiringSend reminder notification
DamageReport.AddedNotify shop owner
Payment.CompletedGenerate invoice
Motorbike.MaintenanceUpdate availability

Source

  • From: D:\project\work\rx-erp messaging patterns

スコア

総合スコア

50/100

リポジトリの品質指標に基づく評価

SKILL.md

SKILL.mdファイルが含まれている

+20
LICENSE

ライセンスが設定されている

0/10
説明文

100文字以上の説明がある

+10
人気

GitHub Stars 100以上

0/15
最近の活動

3ヶ月以内に更新がある

0/10
フォーク

10回以上フォークされている

0/5
Issue管理

オープンIssueが50未満

+5
言語

プログラミング言語が設定されている

+5
タグ

1つ以上のタグが設定されている

0/5

レビュー

💬

レビュー機能は近日公開予定です