1226 lines
56 KiB
C#
1226 lines
56 KiB
C#
using System.Threading.Channels;
|
|
using Microsoft.Extensions.Logging;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Services;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Observability;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Types;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
|
|
using ZB.MOM.WW.ScadaBridge.SiteEventLogging;
|
|
|
|
namespace ZB.MOM.WW.ScadaBridge.StoreAndForward;
|
|
|
|
/// <summary>
|
|
/// Core store-and-forward service.
|
|
///
|
|
/// Lifecycle:
|
|
/// 1. Caller attempts immediate delivery via IDeliveryHandler
|
|
/// 2. On transient failure → buffer in SQLite → retry loop
|
|
/// 3. On success → remove from buffer
|
|
/// 4. On reaching MaxRetries → park (a MaxRetries of 0 means "no limit" — the
|
|
/// message is retried until delivered and is never parked for retry exhaustion)
|
|
/// 5. Permanent failures are returned to caller immediately (never buffered)
|
|
///
|
|
/// Fixed retry interval (not exponential). Per-source-entity retry settings.
|
|
/// Background timer-based retry sweep.
|
|
///
|
|
/// Parked messages queryable, retryable, and discardable.
|
|
///
|
|
/// Buffer depth reported as health metric. Activity logged to site event log.
|
|
///
|
|
/// CachedCall idempotency is the caller's responsibility.
|
|
/// This service does not deduplicate — if the same message is enqueued twice,
|
|
/// it will be delivered twice. Callers using ExternalSystem.CachedCall() must
|
|
/// design their payloads to be idempotent (e.g., include unique request IDs
|
|
/// and handle duplicate detection on the remote end).
|
|
/// </summary>
|
|
public class StoreAndForwardService
|
|
{
|
|
private readonly StoreAndForwardStorage _storage;
|
|
private readonly StoreAndForwardOptions _options;
|
|
private readonly ReplicationService? _replication;
|
|
private readonly ILogger<StoreAndForwardService> _logger;
|
|
/// <summary>
|
|
/// Site-side observer notified
|
|
/// after every cached-call delivery attempt. Optional — when null no
|
|
/// telemetry is emitted; the legacy retry loop behaviour is
|
|
/// preserved exactly.
|
|
/// </summary>
|
|
private readonly ICachedCallLifecycleObserver? _cachedCallObserver;
|
|
/// <summary>
|
|
/// Optional site operational-event log. When non-null the service maps
|
|
/// its own buffer/retry/park activity (the same activity that drives
|
|
/// <see cref="OnActivity"/>) onto site events — <c>store_and_forward</c> for the
|
|
/// cached-call categories and <c>notification</c> for the site's
|
|
/// forward-to-central notification path. Best-effort and fire-and-forget so a
|
|
/// failing logger never affects delivery bookkeeping.
|
|
/// </summary>
|
|
private readonly ISiteEventLogger? _siteEventLogger;
|
|
/// <summary>
|
|
/// Site id stamped onto the
|
|
/// cached-call attempt context so the audit bridge can build the
|
|
/// <see cref="SiteCallOperational"/> half of the telemetry packet.
|
|
/// </summary>
|
|
private readonly string _siteId;
|
|
|
|
/// <summary>
|
|
/// Distinctive marker stamped onto cached-call audit
|
|
/// telemetry when the host has not registered an
|
|
/// <see cref="IStoreAndForwardSiteContext"/>. Chosen with a leading <c>$</c>
|
|
/// so it cannot collide with a real site id (which is a configuration
|
|
/// identifier and never starts with <c>$</c>). Surfacing this in the
|
|
/// central audit log makes a missing site-context binding immediately
|
|
/// recognisable instead of an unattributable empty string.
|
|
/// </summary>
|
|
public const string UnknownSiteSentinel = "$unknown-site";
|
|
|
|
/// <summary>
|
|
/// The detail-string prefix written by <see cref="EnqueueAsync"/>
|
|
/// when an immediate forward attempt throws and the message is buffered for
|
|
/// the retry sweep. <see cref="EmitSiteEvent"/> matches on this same prefix
|
|
/// to distinguish a forward <i>failure</i> (logged) from a routine
|
|
/// no-handler enqueue (not logged), so both the construction site and the
|
|
/// check reference this single constant rather than duplicating the
|
|
/// literal — keeping the two ends from drifting apart.
|
|
/// </summary>
|
|
private const string BufferedForRetryDetailPrefix = "Buffered for retry";
|
|
|
|
private Timer? _retryTimer;
|
|
private int _retryInProgress;
|
|
|
|
/// <summary>
|
|
/// Active-node delivery gate. When set, the retry sweep runs ONLY when the
|
|
/// gate reports true — the standby node applies replicated buffer operations
|
|
/// but must never deliver (Component-StoreAndForward.md "standby-passive"
|
|
/// model). Re-evaluated on every sweep tick so a failover resumes delivery
|
|
/// within one RetryTimerInterval without re-wiring. Null (tests, central
|
|
/// hosts) preserves the ungated legacy behaviour. A throwing gate is treated
|
|
/// as standby — safe-by-default, mirroring SiteCommunicationActor's
|
|
/// DefaultIsActiveCheck fallback.
|
|
/// </summary>
|
|
private Func<bool>? _deliveryGate;
|
|
|
|
/// <summary>Installs the active-node delivery gate (see <see cref="_deliveryGate"/>).</summary>
|
|
/// <param name="gate">Predicate returning <c>true</c> when this node should run the retry sweep.</param>
|
|
public void SetDeliveryGate(Func<bool> gate) => _deliveryGate = gate;
|
|
|
|
/// <summary>
|
|
/// The in-flight retry sweep <see cref="Task"/>, or
|
|
/// <c>null</c> when no sweep is currently running. Captured when the timer
|
|
/// callback starts a sweep so <see cref="StopAsync"/> can wait for it to
|
|
/// finish before the host disposes downstream dependencies
|
|
/// (<see cref="_storage"/>, <see cref="_replication"/>) that the sweep is
|
|
/// still touching. Written from the timer thread and from
|
|
/// <see cref="StopAsync"/>, so reads are synchronised via the
|
|
/// <see cref="Volatile"/> APIs.
|
|
/// </summary>
|
|
private Task? _sweepTask;
|
|
|
|
/// <summary>
|
|
/// Ordered, single-reader queue of cached-call audit-observer notifications.
|
|
/// The retry sweep <b>posts</b> to this channel instead of awaiting the
|
|
/// observer inline, so a slow observer (a SQLite audit write) cannot stretch
|
|
/// the sweep; the single reader (<see cref="_observerPump"/>) preserves
|
|
/// per-operation event order. Re-created in <see cref="StartAsync"/> because
|
|
/// <see cref="StopAsync"/> completes it (a restarted instance needs a fresh
|
|
/// channel). Before <see cref="StartAsync"/> starts the pump, posts fall back
|
|
/// to inline processing (see <see cref="PostObserverNotification"/>).
|
|
/// </summary>
|
|
private Channel<Func<Task>> _observerQueue =
|
|
Channel.CreateUnbounded<Func<Task>>(new UnboundedChannelOptions { SingleReader = true });
|
|
|
|
/// <summary>
|
|
/// The single-reader pump draining <see cref="_observerQueue"/>, or
|
|
/// <c>null</c> when the service is not started. Started in
|
|
/// <see cref="StartAsync"/>, completed + awaited in <see cref="StopAsync"/>.
|
|
/// </summary>
|
|
private Task? _observerPump;
|
|
|
|
/// <summary>
|
|
/// How long <see cref="StopAsync"/> waits for an
|
|
/// in-flight retry sweep to finish before returning. The default — 10 s —
|
|
/// is generous enough to let a typical sweep over the buffered queue drain,
|
|
/// but bounded so a hung downstream call (a stuck SQLite write, a
|
|
/// long-running delivery handler) cannot block host shutdown indefinitely.
|
|
/// On timeout the wait is abandoned and the timer is still disposed; the
|
|
/// sweep keeps running but will throw on the next call into a disposed
|
|
/// dependency — preferred to blocking shutdown forever.
|
|
/// </summary>
|
|
private static readonly TimeSpan SweepShutdownWaitTimeout = TimeSpan.FromSeconds(10);
|
|
|
|
/// <summary>
|
|
/// Cached count of messages currently buffered for
|
|
/// forwarding — i.e. rows in <see cref="StoreAndForwardMessageStatus.Pending"/>,
|
|
/// the live store-and-forward queue waiting to be delivered. This backs the
|
|
/// <c>scadabridge.store_and_forward.queue.depth</c> observable gauge.
|
|
/// <para>
|
|
/// The gauge's collection callback is synchronous and is invoked frequently by
|
|
/// the OpenTelemetry/Prometheus collector, so it must never run an async SQLite
|
|
/// <c>COUNT(*)</c>. Instead this <see cref="long"/> is seeded once from storage
|
|
/// in <see cref="StartAsync"/> and then adjusted in-process on the existing
|
|
/// paths that change the Pending population: <see cref="BufferAsync"/> (+1),
|
|
/// successful-retry removal and Pending→Parked transitions in
|
|
/// <see cref="RetryMessageAsync"/> (-1), and operator requeue in
|
|
/// <see cref="RetryParkedMessageAsync"/> (+1). The provider registered with
|
|
/// <see cref="ScadaBridgeTelemetry.SetQueueDepthProvider"/> reads it via
|
|
/// <see cref="Interlocked.Read"/> — non-blocking and sync-safe. It is an
|
|
/// approximate, eventually-consistent gauge (concurrent failover replication
|
|
/// applies to the standby's own store, not this counter), which is exactly
|
|
/// what a queue-depth metric needs.
|
|
/// </para>
|
|
/// </summary>
|
|
private long _bufferedCount;
|
|
|
|
/// <summary>
|
|
/// Test seam: simulates a concurrent pre-seed
|
|
/// <see cref="BufferAsync"/> increment landing on <see cref="_bufferedCount"/>
|
|
/// before <see cref="StartAsync"/> seeds it, so a test can prove the seed uses
|
|
/// <see cref="Interlocked.Add"/> (additive) rather than Exchange (clobbering).
|
|
/// </summary>
|
|
internal void TestOnly_IncrementBufferedCount() =>
|
|
Interlocked.Increment(ref _bufferedCount);
|
|
|
|
/// <summary>
|
|
/// An instance field that guards against a single instance
|
|
/// registering the queue-depth provider (and re-seeding the counter) more than
|
|
/// once — e.g. a second <see cref="StartAsync"/> on the same instance. It does NOT
|
|
/// coordinate across instances: the gauge slot in <see cref="ScadaBridgeTelemetry"/>
|
|
/// is process-global, so in a multi-instance process the last <see cref="StartAsync"/>
|
|
/// wins the global slot. 0 = not yet registered, 1 = done.
|
|
/// </summary>
|
|
private int _queueDepthProviderRegistered;
|
|
|
|
/// <summary>
|
|
/// The exact provider delegate this instance registered with
|
|
/// the process-global <see cref="ScadaBridgeTelemetry"/> gauge slot, retained so
|
|
/// <see cref="StopAsync"/> can deregister it by identity (compare-and-clear). Holding
|
|
/// the same reference both registers and clears the slot; the identity check ensures a
|
|
/// late stop of this instance cannot clear a newer instance's provider. Null until
|
|
/// <see cref="StartAsync"/> registers.
|
|
/// </summary>
|
|
private Func<long>? _queueDepthProvider;
|
|
|
|
/// <summary>
|
|
/// Delivery handler delegate. The return value / exception is interpreted
|
|
/// the same way on both the immediate-delivery path (<see cref="EnqueueAsync"/>)
|
|
/// and the background retry path (<c>RetryMessageAsync</c>):
|
|
/// <list type="bullet">
|
|
/// <item><description><c>true</c> — delivered successfully. The message is
|
|
/// removed from the buffer (or, on the immediate path, never buffered).</description></item>
|
|
/// <item><description><c>false</c> — permanent failure. On the immediate path
|
|
/// the message is NOT buffered; on a retry the message is already buffered and
|
|
/// is parked immediately (no further retries).</description></item>
|
|
/// <item><description>throws — transient failure. On the immediate path the
|
|
/// message is buffered for retry; on a retry the retry count is incremented and
|
|
/// the message is parked once <see cref="StoreAndForwardMessage.MaxRetries"/> is
|
|
/// reached.</description></item>
|
|
/// </list>
|
|
/// </summary>
|
|
private readonly Dictionary<StoreAndForwardCategory, Func<StoreAndForwardMessage, Task<bool>>> _deliveryHandlers = new();
|
|
|
|
/// <summary>
|
|
/// Event callback for logging S&F activity to site event log.
|
|
/// </summary>
|
|
public event Action<string, StoreAndForwardCategory, string>? OnActivity;
|
|
|
|
/// <summary>
|
|
/// Initializes a new instance of the StoreAndForwardService.
|
|
/// </summary>
|
|
/// <param name="storage">The storage backend for buffered messages.</param>
|
|
/// <param name="options">Configuration options.</param>
|
|
/// <param name="logger">Logger instance.</param>
|
|
/// <param name="replication">Optional replication service for standby synchronization.</param>
|
|
/// <param name="cachedCallObserver">Optional observer for cached call lifecycle events.</param>
|
|
/// <param name="siteId">The site identifier this service belongs to.</param>
|
|
/// <param name="siteEventLogger">
|
|
/// Optional site operational-event log. When non-null, buffer/retry/park
|
|
/// activity is mirrored to site events (<c>store_and_forward</c> /
|
|
/// <c>notification</c> by category). Optional with a <c>null</c> default so the
|
|
/// many direct-construction tests still compile unchanged.
|
|
/// </param>
|
|
public StoreAndForwardService(
|
|
StoreAndForwardStorage storage,
|
|
StoreAndForwardOptions options,
|
|
ILogger<StoreAndForwardService> logger,
|
|
ReplicationService? replication = null,
|
|
ICachedCallLifecycleObserver? cachedCallObserver = null,
|
|
string siteId = "",
|
|
ISiteEventLogger? siteEventLogger = null)
|
|
{
|
|
_storage = storage;
|
|
_options = options;
|
|
_logger = logger;
|
|
_replication = replication;
|
|
_cachedCallObserver = cachedCallObserver;
|
|
_siteId = string.IsNullOrWhiteSpace(siteId) ? UnknownSiteSentinel : siteId;
|
|
_siteEventLogger = siteEventLogger;
|
|
|
|
// Ride the existing activity hook to emit site operational events.
|
|
// RaiseActivity already isolates a throwing subscriber, so a failing
|
|
// event log can never be misclassified as a transient delivery failure.
|
|
// Only subscribe when a logger is wired so the
|
|
// legacy (test/central) construction path stays a no-op.
|
|
if (_siteEventLogger != null)
|
|
{
|
|
OnActivity += EmitSiteEvent;
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Maps one store-and-forward activity to a site operational event,
|
|
/// following the Site Event Logging spec's per-category scope
|
|
/// (Component-SiteEventLogging.md §"Events Logged"):
|
|
/// <list type="bullet">
|
|
/// <item><description>Cached-call categories
|
|
/// (<see cref="StoreAndForwardCategory.ExternalSystem"/> /
|
|
/// <see cref="StoreAndForwardCategory.CachedDbWrite"/>) log under
|
|
/// <c>store_and_forward</c> for queued / retried / parked / retry-delivered
|
|
/// activity.</description></item>
|
|
/// <item><description>The site's notification forward-to-central path
|
|
/// (<see cref="StoreAndForwardCategory.Notification"/>) logs under
|
|
/// <c>notification</c> ONLY on a forward FAILURE (buffered after the
|
|
/// immediate forward threw) or a park (long-buffered / retries exhausted).
|
|
/// Routine enqueue and forward-success are deliberately NOT logged — central's
|
|
/// <c>Notifications</c> table is the record of audit; the site only fills the
|
|
/// in-transit blind spot when central is unreachable.</description></item>
|
|
/// </list>
|
|
/// A successful immediate cached-call <c>Delivered</c> is the normal hot path and
|
|
/// is not logged.
|
|
/// </summary>
|
|
private void EmitSiteEvent(string action, StoreAndForwardCategory category, string detail)
|
|
{
|
|
var logger = _siteEventLogger;
|
|
if (logger == null)
|
|
{
|
|
return;
|
|
}
|
|
|
|
// An immediate-delivery success is the normal hot path, not an
|
|
// operational event. A retry-loop success (detail "Delivered to … after
|
|
// N retries") IS logged for cached calls — it records a recovery.
|
|
if (action == "Delivered" && detail.StartsWith("Immediate", StringComparison.Ordinal))
|
|
{
|
|
return;
|
|
}
|
|
|
|
if (category == StoreAndForwardCategory.Notification)
|
|
{
|
|
// Spec: log only forward-failure (the immediate forward threw and the
|
|
// notification was buffered for retry — detail prefixed
|
|
// BufferedForRetryDetailPrefix) and park. A routine "No handler
|
|
// registered, buffered" enqueue and a forward-success "Delivered"
|
|
// are deliberately NOT logged.
|
|
var isForwardFailure = action == "Queued"
|
|
&& detail.StartsWith(BufferedForRetryDetailPrefix, StringComparison.Ordinal);
|
|
if (!isForwardFailure && action != "Parked")
|
|
{
|
|
return;
|
|
}
|
|
|
|
var notifSeverity = action == "Parked" ? "Error" : "Warning";
|
|
_ = logger.LogEventAsync(
|
|
"notification", notifSeverity, instanceId: null,
|
|
source: "StoreAndForwardService",
|
|
message: $"Notification {action.ToLowerInvariant()}: {detail}");
|
|
return;
|
|
}
|
|
|
|
// Cached-call categories: queued / retried / parked / retry-delivered.
|
|
// Severity: parking is an Error (delivery abandoned for retry purposes);
|
|
// queue/retry/requeue are Warning; a retry-loop Delivered is Info.
|
|
var severity = action switch
|
|
{
|
|
"Parked" => "Error",
|
|
"Delivered" => "Info",
|
|
_ => "Warning",
|
|
};
|
|
|
|
_ = logger.LogEventAsync(
|
|
"store_and_forward", severity, instanceId: null,
|
|
source: "StoreAndForwardService",
|
|
message: $"Operation {action.ToLowerInvariant()}: {detail}");
|
|
}
|
|
|
|
/// <summary>
|
|
/// Registers a delivery handler for a given message category. See the
|
|
/// <c>_deliveryHandlers</c> field documentation for the true/false/throws contract,
|
|
/// which applies identically on the immediate and retry paths.
|
|
/// </summary>
|
|
/// <param name="category">The message category to handle.</param>
|
|
/// <param name="handler">The delivery handler function.</param>
|
|
public void RegisterDeliveryHandler(
|
|
StoreAndForwardCategory category,
|
|
Func<StoreAndForwardMessage, Task<bool>> handler)
|
|
{
|
|
_deliveryHandlers[category] = handler;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Initializes storage and starts the background retry timer.
|
|
/// </summary>
|
|
/// <returns>A task representing the asynchronous start operation.</returns>
|
|
public async Task StartAsync()
|
|
{
|
|
await _storage.InitializeAsync();
|
|
|
|
var pending = await _storage.GetMessageCountByStatusAsync(
|
|
StoreAndForwardMessageStatus.Pending);
|
|
if (Interlocked.CompareExchange(ref _queueDepthProviderRegistered, 1, 0) == 0)
|
|
{
|
|
Interlocked.Add(ref _bufferedCount, pending);
|
|
var provider = (Func<long>)(() => Interlocked.Read(ref _bufferedCount));
|
|
_queueDepthProvider = provider;
|
|
ScadaBridgeTelemetry.SetQueueDepthProvider(provider);
|
|
}
|
|
|
|
// (Re)create the observer queue + start its single-reader pump. A prior
|
|
// StopAsync completes the channel, so a restarted instance needs a fresh
|
|
// one. The pump is best-effort: an observer that throws is logged and
|
|
// swallowed so a failing audit pipeline never corrupts retry bookkeeping.
|
|
_observerQueue = Channel.CreateUnbounded<Func<Task>>(
|
|
new UnboundedChannelOptions { SingleReader = true });
|
|
_observerPump = Task.Run(async () =>
|
|
{
|
|
await foreach (var work in _observerQueue.Reader.ReadAllAsync())
|
|
{
|
|
try
|
|
{
|
|
await work();
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogWarning(ex,
|
|
"Cached-call audit observer threw; ignored (best-effort per alog.md §7)");
|
|
}
|
|
}
|
|
});
|
|
|
|
_retryTimer = new Timer(
|
|
_ => KickSweep(),
|
|
null,
|
|
_options.RetryTimerInterval,
|
|
_options.RetryTimerInterval);
|
|
|
|
_logger.LogInformation(
|
|
"Store-and-forward service started. Retry interval: {Interval}s",
|
|
_options.DefaultRetryInterval.TotalSeconds);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Stops the background retry timer and waits (bounded) for any in-flight
|
|
/// retry sweep to finish before returning.
|
|
/// </summary>
|
|
/// <returns>A task representing the asynchronous stop operation.</returns>
|
|
public async Task StopAsync()
|
|
{
|
|
if (_retryTimer != null)
|
|
{
|
|
// Stop the periodic callback first so no new sweep starts while we
|
|
// are waiting for the in-flight one to drain.
|
|
await _retryTimer.DisposeAsync();
|
|
_retryTimer = null;
|
|
}
|
|
|
|
var inflight = Volatile.Read(ref _sweepTask);
|
|
if (inflight is not null && !inflight.IsCompleted)
|
|
{
|
|
try
|
|
{
|
|
// WaitAsync with a finite timeout: a hung delivery handler /
|
|
// storage call cannot block host shutdown indefinitely. On timeout
|
|
// the sweep keeps running but the host is free to proceed with
|
|
// disposal — preferred to never returning.
|
|
await inflight.WaitAsync(SweepShutdownWaitTimeout).ConfigureAwait(false);
|
|
}
|
|
catch (TimeoutException)
|
|
{
|
|
_logger.LogWarning(
|
|
"Store-and-forward retry sweep did not finish within {Timeout}; " +
|
|
"shutdown is proceeding while the sweep is still in-flight",
|
|
SweepShutdownWaitTimeout);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
// The sweep itself already logs at Error on failure (see
|
|
// RetryPendingMessagesAsync's catch); we only log here so a
|
|
// surprise fault during shutdown is still visible. Swallow so the
|
|
// host's shutdown sequence can continue regardless.
|
|
_logger.LogWarning(ex,
|
|
"Store-and-forward retry sweep faulted during shutdown wait");
|
|
}
|
|
}
|
|
|
|
// Complete the observer queue and wait (bounded) for the pump to drain any
|
|
// posted-but-not-yet-processed notifications, so shutdown does not silently
|
|
// drop audit events that the sweep already handed off.
|
|
_observerQueue.Writer.TryComplete();
|
|
var pump = _observerPump;
|
|
if (pump is not null)
|
|
{
|
|
try
|
|
{
|
|
await pump.WaitAsync(SweepShutdownWaitTimeout).ConfigureAwait(false);
|
|
}
|
|
catch (TimeoutException)
|
|
{
|
|
_logger.LogWarning(
|
|
"Audit-observer pump did not drain within {Timeout}", SweepShutdownWaitTimeout);
|
|
}
|
|
_observerPump = null;
|
|
}
|
|
|
|
var provider = _queueDepthProvider;
|
|
if (provider is not null)
|
|
{
|
|
ScadaBridgeTelemetry.ClearQueueDepthProvider(provider);
|
|
_queueDepthProvider = null;
|
|
}
|
|
Interlocked.Exchange(ref _bufferedCount, 0);
|
|
Interlocked.Exchange(ref _queueDepthProviderRegistered, 0);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Enqueues a message for store-and-forward delivery.
|
|
/// Attempts immediate delivery first. On transient failure, buffers for retry.
|
|
/// On permanent failure (handler returns false), returns false immediately.
|
|
///
|
|
/// Retry-count lifecycle — the immediate (or caller-made) delivery attempt
|
|
/// is attempt 0 and is not counted; the background retry sweep increments
|
|
/// <see cref="StoreAndForwardMessage.RetryCount"/> on each retry. A buffered
|
|
/// message is parked once <c>RetryCount</c> reaches <paramref name="maxRetries"/>
|
|
/// — <b>but only when <paramref name="maxRetries"/> is greater than 0</b>. A
|
|
/// <paramref name="maxRetries"/> of <c>0</c> means <b>no limit</b>: the message is
|
|
/// retried on every sweep until it is delivered and is <b>never parked</b> on a
|
|
/// retry-count basis. It is therefore <i>not</i> a "do not retry" value — callers
|
|
/// that want delivery abandoned after a bounded number of attempts must pass a
|
|
/// positive <paramref name="maxRetries"/>.
|
|
///
|
|
/// CachedCall idempotency note — this method does not deduplicate.
|
|
/// The caller (e.g., ExternalSystem.CachedCall()) is responsible for ensuring
|
|
/// that the remote system can handle duplicate deliveries safely.
|
|
/// </summary>
|
|
/// <param name="category">Message category (selects the delivery handler).</param>
|
|
/// <param name="target">Target system name (external system / notification list / DB connection).</param>
|
|
/// <param name="payloadJson">JSON-serialized call payload, treated opaquely.</param>
|
|
/// <param name="originInstanceName">Instance that originated the message (survives instance deletion).</param>
|
|
/// <param name="maxRetries">
|
|
/// Maximum background retry-sweep attempts before the message is parked.
|
|
/// <b><c>0</c> = no limit</b> — the message is retried on every sweep until
|
|
/// delivered and is never parked for exhausting retries; it is <b>not</b> a
|
|
/// "never retry" value. <c>null</c> uses <see cref="StoreAndForwardOptions.DefaultMaxRetries"/>.
|
|
/// Must be positive to bound delivery attempts. Mirrors the
|
|
/// <see cref="StoreAndForwardMessage.MaxRetries"/> contract.
|
|
/// </param>
|
|
/// <param name="retryInterval">Fixed interval between retry sweeps for this message; <c>null</c> uses the configured default.</param>
|
|
/// <param name="attemptImmediateDelivery">
|
|
/// When <c>false</c>, the caller has already made its own delivery attempt and the
|
|
/// message is buffered directly for the retry sweep (the handler is not invoked here).
|
|
/// </param>
|
|
/// <param name="deferToSweep">
|
|
/// When <c>true</c>, the handler is never invoked inline — the message is buffered
|
|
/// due-immediately (<see cref="StoreAndForwardMessage.LastAttemptAt"/> stays null)
|
|
/// and a background sweep is kicked. Use for callers on latency-sensitive threads
|
|
/// (e.g. script dispatchers) whose delivery involves a long Ask timeout. Unlike
|
|
/// <paramref name="attemptImmediateDelivery"/><c>: false</c>, no delivery attempt is
|
|
/// presumed to have already been made, so the row is due on the very next sweep.
|
|
/// </param>
|
|
/// <param name="messageId">
|
|
/// An explicit, caller-supplied message id. <c>null</c> (the default) makes the
|
|
/// service mint a fresh GUID. The Notification Outbox enqueue path supplies its own
|
|
/// id so the script-generated <c>NotificationId</c> is the single idempotency key —
|
|
/// it is the buffered row's <see cref="StoreAndForwardMessage.Id"/>, it is carried
|
|
/// inside the payload, and it is the id the forwarder submits to central.
|
|
/// </param>
|
|
/// <param name="executionId">
|
|
/// The originating script execution's
|
|
/// per-run correlation id. Threaded onto the buffered row so the retry-loop
|
|
/// cached-call audit rows carry it. <c>null</c> for callers (e.g. notifications)
|
|
/// that do not supply one.
|
|
/// </param>
|
|
/// <param name="sourceScript">
|
|
/// The originating script identifier,
|
|
/// threaded onto the buffered row alongside <paramref name="executionId"/>
|
|
/// so the retry-loop audit rows carry the same provenance the script-side
|
|
/// cached rows do. <c>null</c> when not known.
|
|
/// </param>
|
|
/// <param name="parentExecutionId">
|
|
/// The <c>ExecutionId</c> of the
|
|
/// inbound-API request that spawned the originating script execution.
|
|
/// Threaded onto the buffered row alongside <paramref name="executionId"/>
|
|
/// so the retry-loop cached-call audit rows carry it. <c>null</c> for a
|
|
/// non-routed run and for callers (e.g. notifications) that
|
|
/// do not supply one.
|
|
/// </param>
|
|
/// <returns>A task that resolves to a result indicating whether the message was delivered or buffered.</returns>
|
|
public async Task<StoreAndForwardResult> EnqueueAsync(
|
|
StoreAndForwardCategory category,
|
|
string target,
|
|
string payloadJson,
|
|
string? originInstanceName = null,
|
|
int? maxRetries = null,
|
|
TimeSpan? retryInterval = null,
|
|
bool attemptImmediateDelivery = true,
|
|
bool deferToSweep = false,
|
|
string? messageId = null,
|
|
Guid? executionId = null,
|
|
string? sourceScript = null,
|
|
Guid? parentExecutionId = null)
|
|
{
|
|
var message = new StoreAndForwardMessage
|
|
{
|
|
Id = messageId ?? Guid.NewGuid().ToString("N"),
|
|
Category = category,
|
|
Target = target,
|
|
PayloadJson = payloadJson,
|
|
RetryCount = 0,
|
|
MaxRetries = maxRetries ?? _options.DefaultMaxRetries,
|
|
RetryIntervalMs = (long)(retryInterval ?? _options.DefaultRetryInterval).TotalMilliseconds,
|
|
CreatedAt = DateTimeOffset.UtcNow,
|
|
Status = StoreAndForwardMessageStatus.Pending,
|
|
OriginInstanceName = originInstanceName,
|
|
ExecutionId = executionId,
|
|
SourceScript = sourceScript,
|
|
ParentExecutionId = parentExecutionId
|
|
};
|
|
|
|
// Latency-sensitive caller: never run the handler inline. Buffer
|
|
// due-immediately (LastAttemptAt intentionally left null so the row is due
|
|
// on the very next sweep, not held back by RetryInterval) and kick a
|
|
// background sweep so the healthy-path latency is milliseconds.
|
|
if (deferToSweep)
|
|
{
|
|
await BufferAsync(message);
|
|
RaiseActivity("Queued", category, $"Deferred to sweep: {target}");
|
|
TriggerSweep();
|
|
return new StoreAndForwardResult(true, message.Id, true);
|
|
}
|
|
|
|
// Attempt immediate delivery — unless the caller has already made a
|
|
// delivery attempt of its own (attemptImmediateDelivery: false). In that
|
|
// case re-invoking the handler here would dispatch the request twice.
|
|
if (attemptImmediateDelivery && _deliveryHandlers.TryGetValue(category, out var handler))
|
|
{
|
|
try
|
|
{
|
|
var success = await handler(message);
|
|
if (success)
|
|
{
|
|
RaiseActivity("Delivered", category, $"Immediate delivery to {target}");
|
|
return new StoreAndForwardResult(true, message.Id, false);
|
|
}
|
|
|
|
// Permanent failure — do not buffer
|
|
return new StoreAndForwardResult(false, message.Id, false);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
// Transient failure — buffer for retry. The immediate attempt is
|
|
// attempt 0; RetryCount tracks only sweep retries, so it stays 0
|
|
// here.
|
|
_logger.LogWarning(ex,
|
|
"Immediate delivery to {Target} failed (transient), buffering for retry",
|
|
target);
|
|
|
|
message.LastAttemptAt = DateTimeOffset.UtcNow;
|
|
message.LastError = ex.Message;
|
|
await BufferAsync(message);
|
|
|
|
RaiseActivity("Queued", category, $"{BufferedForRetryDetailPrefix}: {target} ({ex.Message})");
|
|
return new StoreAndForwardResult(true, message.Id, true);
|
|
}
|
|
}
|
|
|
|
// Either no handler is registered yet, or the caller already attempted
|
|
// delivery itself — buffer for the background retry sweep to deliver.
|
|
// The initial attempt (caller-made, or skipped because no handler is
|
|
// registered) is attempt 0; RetryCount tracks only sweep retries and
|
|
// therefore stays 0 here.
|
|
if (!attemptImmediateDelivery)
|
|
{
|
|
message.LastAttemptAt = DateTimeOffset.UtcNow;
|
|
}
|
|
await BufferAsync(message);
|
|
RaiseActivity("Queued", category, attemptImmediateDelivery
|
|
? $"No handler registered, buffered: {target}"
|
|
: $"{BufferedForRetryDetailPrefix}: {target}");
|
|
return new StoreAndForwardResult(true, message.Id, true);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Persists a message to the local SQLite buffer and replicates the
|
|
/// add to the standby node so a failover does not lose the buffered message.
|
|
/// </summary>
|
|
private async Task BufferAsync(StoreAndForwardMessage message)
|
|
{
|
|
await _storage.EnqueueAsync(message);
|
|
_replication?.ReplicateEnqueue(message);
|
|
Interlocked.Increment(ref _bufferedCount);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Kicks a background retry sweep now (fire-and-forget). No-op before
|
|
/// <see cref="StartAsync"/> (storage may be uninitialized and the timer, whose
|
|
/// presence gates this, is not yet set). Overlap-safe: publishes into
|
|
/// <see cref="_sweepTask"/> only when the CAS is won — see <see cref="KickSweep"/>.
|
|
/// </summary>
|
|
public void TriggerSweep()
|
|
{
|
|
if (_retryTimer == null) return;
|
|
KickSweep();
|
|
}
|
|
|
|
/// <summary>The current drain handle — test seam for the N3 clobber regression.</summary>
|
|
internal Task? CurrentSweepTaskForTest => Volatile.Read(ref _sweepTask);
|
|
|
|
/// <summary>
|
|
/// Starts a sweep IF none is in flight, publishing the new task into
|
|
/// <see cref="_sweepTask"/> only when this call wins the <see cref="_retryInProgress"/>
|
|
/// CAS. A kick that loses the CAS returns without touching <see cref="_sweepTask"/> —
|
|
/// pre-fix it overwrote the drain handle with an instantly-completed no-op, so
|
|
/// <see cref="StopAsync"/> proceeded with disposal under a still-running sweep
|
|
/// (review 02 round 2, N3) — defeating the drain exactly when a sweep outlives a tick.
|
|
/// </summary>
|
|
private void KickSweep()
|
|
{
|
|
if (Interlocked.CompareExchange(ref _retryInProgress, 1, 0) != 0)
|
|
return;
|
|
Volatile.Write(ref _sweepTask, RunSweepOwnedAsync());
|
|
}
|
|
|
|
/// <summary>
|
|
/// Background retry sweep. Processes all pending messages that are due for retry.
|
|
/// Self-CASes for direct callers (existing tests); the entry CAS is skipped when
|
|
/// ownership is already held by <see cref="KickSweep"/>.
|
|
/// </summary>
|
|
/// <returns>A task representing the asynchronous retry sweep.</returns>
|
|
internal async Task RetryPendingMessagesAsync()
|
|
{
|
|
// Prevent overlapping retry sweeps
|
|
if (Interlocked.CompareExchange(ref _retryInProgress, 1, 0) != 0)
|
|
return;
|
|
await RunSweepOwnedAsync();
|
|
}
|
|
|
|
/// <summary>
|
|
/// The actual retry sweep body. The caller MUST already own the
|
|
/// <see cref="_retryInProgress"/> flag (won the CAS); this method releases it in its
|
|
/// finally. Never call directly without holding ownership.
|
|
/// </summary>
|
|
private async Task RunSweepOwnedAsync()
|
|
{
|
|
try
|
|
{
|
|
var gate = _deliveryGate;
|
|
if (gate != null)
|
|
{
|
|
bool isActive;
|
|
try { isActive = gate(); }
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogWarning(ex,
|
|
"S&F delivery gate threw; treating this node as standby for this sweep");
|
|
isActive = false;
|
|
}
|
|
if (!isActive)
|
|
{
|
|
_logger.LogDebug("S&F retry sweep skipped: this node is not the active site node");
|
|
return;
|
|
}
|
|
}
|
|
|
|
var messages = await _storage.GetMessagesForRetryAsync(_options.SweepBatchLimit);
|
|
if (messages.Count == 0) return;
|
|
|
|
_logger.LogDebug("Retry sweep: {Count} messages due for retry", messages.Count);
|
|
|
|
// Group due rows into per-(category,target) lanes. GroupBy is stable, so
|
|
// each lane preserves created_at-ASC (oldest-first) order. Lanes run
|
|
// concurrently up to SweepTargetParallelism, so one slow/dead target no
|
|
// longer blocks delivery to healthy ones; within a lane delivery stays
|
|
// strictly sequential (per-target FIFO).
|
|
var lanes = messages
|
|
.GroupBy(m => (m.Category, m.Target))
|
|
.Select(g => g.ToList())
|
|
.ToList();
|
|
|
|
using var laneCap = new SemaphoreSlim(Math.Max(1, _options.SweepTargetParallelism));
|
|
var laneTasks = lanes.Select(async lane =>
|
|
{
|
|
await laneCap.WaitAsync();
|
|
try
|
|
{
|
|
foreach (var message in lane)
|
|
{
|
|
// Short-circuit, per lane: one transient failure means
|
|
// the target is down — skip the rest of this lane this sweep
|
|
// rather than burning a full timeout per remaining message.
|
|
// Skipped rows keep their RetryCount/LastAttemptAt untouched.
|
|
var outcome = await RetryMessageAsync(message);
|
|
if (outcome == RetryOutcome.TransientFailure)
|
|
return;
|
|
}
|
|
}
|
|
finally { laneCap.Release(); }
|
|
}).ToList();
|
|
await Task.WhenAll(laneTasks);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogError(ex, "Error during retry sweep");
|
|
}
|
|
finally
|
|
{
|
|
Interlocked.Exchange(ref _retryInProgress, 0);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Outcome of a single message's retry attempt, used by the sweep to
|
|
/// short-circuit a <c>(category, target)</c> lane after its first transient
|
|
/// failure. <see cref="Skipped"/> = no attempt was made (no handler, or the
|
|
/// row's status changed under a concurrent apply); its RetryCount is untouched.
|
|
/// </summary>
|
|
internal enum RetryOutcome { Delivered, Parked, TransientFailure, Skipped }
|
|
|
|
private async Task<RetryOutcome> RetryMessageAsync(StoreAndForwardMessage message)
|
|
{
|
|
if (!_deliveryHandlers.TryGetValue(message.Category, out var handler))
|
|
{
|
|
_logger.LogWarning("No delivery handler for category {Category}", message.Category);
|
|
return RetryOutcome.Skipped;
|
|
}
|
|
|
|
// Measure per-attempt
|
|
// duration so the audit row carries a meaningful DurationMs. Captured
|
|
// around the handler invocation only — storage / replication overhead
|
|
// is excluded.
|
|
var attemptStartUtc = DateTime.UtcNow;
|
|
var attemptStopwatch = System.Diagnostics.Stopwatch.StartNew();
|
|
|
|
try
|
|
{
|
|
var success = await handler(message);
|
|
attemptStopwatch.Stop();
|
|
if (success)
|
|
{
|
|
await _storage.RemoveMessageAsync(message.Id);
|
|
_replication?.ReplicateRemove(message.Id);
|
|
Interlocked.Decrement(ref _bufferedCount);
|
|
RaiseActivity("Delivered", message.Category,
|
|
$"Delivered to {message.Target} after {message.RetryCount} retries");
|
|
|
|
// Terminal Delivered observer notification — the audit
|
|
// bridge maps this to Attempted + CachedResolve(Delivered).
|
|
await PostObserverNotification(
|
|
message,
|
|
CachedCallAttemptOutcome.Delivered,
|
|
lastError: null,
|
|
httpStatus: null,
|
|
occurredAtUtc: attemptStartUtc,
|
|
durationMs: (int)attemptStopwatch.ElapsedMilliseconds);
|
|
return RetryOutcome.Delivered;
|
|
}
|
|
|
|
// Permanent failure on retry — park immediately.
|
|
message.Status = StoreAndForwardMessageStatus.Parked;
|
|
message.LastAttemptAt = DateTimeOffset.UtcNow;
|
|
message.LastError = "Permanent failure (handler returned false)";
|
|
var parked = await _storage.UpdateMessageIfStatusAsync(
|
|
message, StoreAndForwardMessageStatus.Pending);
|
|
if (!parked)
|
|
{
|
|
_logger.LogDebug(
|
|
"Message {MessageId} changed status during delivery; sweep park skipped",
|
|
message.Id);
|
|
return RetryOutcome.Skipped;
|
|
}
|
|
Interlocked.Decrement(ref _bufferedCount);
|
|
_replication?.ReplicatePark(message);
|
|
RaiseActivity("Parked", message.Category,
|
|
$"Permanent failure for {message.Target}: handler returned false");
|
|
|
|
// Terminal PermanentFailure observer notification — the
|
|
// audit bridge maps this to Attempted(Failed) + CachedResolve(Parked).
|
|
await PostObserverNotification(
|
|
message,
|
|
CachedCallAttemptOutcome.PermanentFailure,
|
|
lastError: message.LastError,
|
|
httpStatus: null,
|
|
occurredAtUtc: attemptStartUtc,
|
|
durationMs: (int)attemptStopwatch.ElapsedMilliseconds);
|
|
return RetryOutcome.Parked;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
attemptStopwatch.Stop();
|
|
// Transient failure — increment retry, check max
|
|
message.RetryCount++;
|
|
message.LastAttemptAt = DateTimeOffset.UtcNow;
|
|
message.LastError = ex.Message;
|
|
|
|
if (message.MaxRetries > 0 && message.RetryCount >= message.MaxRetries)
|
|
{
|
|
message.Status = StoreAndForwardMessageStatus.Parked;
|
|
var parked = await _storage.UpdateMessageIfStatusAsync(
|
|
message, StoreAndForwardMessageStatus.Pending);
|
|
if (!parked)
|
|
{
|
|
_logger.LogDebug(
|
|
"Message {MessageId} changed status during delivery; sweep park skipped",
|
|
message.Id);
|
|
return RetryOutcome.Skipped;
|
|
}
|
|
Interlocked.Decrement(ref _bufferedCount);
|
|
_replication?.ReplicatePark(message);
|
|
RaiseActivity("Parked", message.Category,
|
|
$"Max retries ({message.MaxRetries}) reached for {message.Target}");
|
|
_logger.LogWarning(
|
|
"Message {MessageId} parked after {MaxRetries} retries to {Target}",
|
|
message.Id, message.MaxRetries, message.Target);
|
|
|
|
// Terminal ParkedMaxRetries observer notification — the
|
|
// audit bridge maps this to Attempted(Failed) + CachedResolve(Parked).
|
|
await PostObserverNotification(
|
|
message,
|
|
CachedCallAttemptOutcome.ParkedMaxRetries,
|
|
lastError: ex.Message,
|
|
httpStatus: null,
|
|
occurredAtUtc: attemptStartUtc,
|
|
durationMs: (int)attemptStopwatch.ElapsedMilliseconds);
|
|
return RetryOutcome.Parked;
|
|
}
|
|
else
|
|
{
|
|
if (!await _storage.UpdateMessageIfStatusAsync(
|
|
message, StoreAndForwardMessageStatus.Pending))
|
|
{
|
|
_logger.LogDebug(
|
|
"Message {MessageId} changed status during delivery; sweep retry-count update skipped",
|
|
message.Id);
|
|
return RetryOutcome.Skipped;
|
|
}
|
|
RaiseActivity("Retried", message.Category,
|
|
$"Retry {message.RetryCount}/{message.MaxRetries} for {message.Target}: {ex.Message}");
|
|
|
|
// Per-attempt TransientFailure observer notification —
|
|
// the audit bridge maps this to Attempted(Failed).
|
|
await PostObserverNotification(
|
|
message,
|
|
CachedCallAttemptOutcome.TransientFailure,
|
|
lastError: ex.Message,
|
|
httpStatus: null,
|
|
occurredAtUtc: attemptStartUtc,
|
|
durationMs: (int)attemptStopwatch.ElapsedMilliseconds);
|
|
return RetryOutcome.TransientFailure;
|
|
}
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Hands a cached-call observer notification to the ordered single-reader pump
|
|
/// (see <see cref="_observerQueue"/>) and returns immediately, so a slow
|
|
/// observer cannot stretch the sweep. When the service is not started (no pump —
|
|
/// e.g. a unit test driving <see cref="RetryPendingMessagesAsync"/> directly),
|
|
/// the notification is processed inline so it is still delivered synchronously.
|
|
/// The buffered <paramref name="message"/> is not mutated again after this call
|
|
/// within a sweep (the next tick loads fresh row instances), so the deferred
|
|
/// read of its fields by the pump is safe.
|
|
/// </summary>
|
|
private ValueTask PostObserverNotification(
|
|
StoreAndForwardMessage message,
|
|
CachedCallAttemptOutcome outcome,
|
|
string? lastError,
|
|
int? httpStatus,
|
|
DateTime occurredAtUtc,
|
|
int? durationMs)
|
|
{
|
|
if (_cachedCallObserver == null)
|
|
{
|
|
return ValueTask.CompletedTask;
|
|
}
|
|
|
|
Func<Task> work = () => NotifyCachedCallObserverAsync(
|
|
message, outcome, lastError, httpStatus, occurredAtUtc, durationMs);
|
|
|
|
if (_observerPump is null)
|
|
{
|
|
// Not started: process inline (preserves prior test behavior).
|
|
return new ValueTask(work());
|
|
}
|
|
|
|
_observerQueue.Writer.TryWrite(work);
|
|
return ValueTask.CompletedTask;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Notify the registered
|
|
/// <see cref="ICachedCallLifecycleObserver"/> of the just-completed
|
|
/// attempt. Invoked from the single-reader observer pump (latency-isolated
|
|
/// from the sweep; order preserved per posting sequence) — or inline when the
|
|
/// service is not started. Only fires for cached-call categories
|
|
/// (<see cref="StoreAndForwardCategory.ExternalSystem"/> and
|
|
/// <see cref="StoreAndForwardCategory.CachedDbWrite"/>); the
|
|
/// <see cref="StoreAndForwardCategory.Notification"/> category has its
|
|
/// own central-side audit pipeline (Notification Outbox) and must
|
|
/// not surface on this hook.
|
|
/// </summary>
|
|
/// <remarks>
|
|
/// Best-effort: an observer that throws is logged and swallowed so a
|
|
/// failing audit pipeline cannot corrupt S&F retry bookkeeping
|
|
/// (alog.md §7 contract). Messages whose ids are not valid GUIDs (callers
|
|
/// that didn't thread a TrackedOperationId in) are silently
|
|
/// skipped — the observer requires a parseable id by contract.
|
|
/// </remarks>
|
|
private async Task NotifyCachedCallObserverAsync(
|
|
StoreAndForwardMessage message,
|
|
CachedCallAttemptOutcome outcome,
|
|
string? lastError,
|
|
int? httpStatus,
|
|
DateTime occurredAtUtc,
|
|
int? durationMs)
|
|
{
|
|
if (_cachedCallObserver == null)
|
|
{
|
|
return;
|
|
}
|
|
|
|
// Only cached-call categories generate audit telemetry on this hook —
|
|
// notifications have their own outbox-side audit pipeline.
|
|
var channel = message.Category switch
|
|
{
|
|
StoreAndForwardCategory.ExternalSystem => "ApiOutbound",
|
|
StoreAndForwardCategory.CachedDbWrite => "DbOutbound",
|
|
_ => null,
|
|
};
|
|
if (channel is null)
|
|
{
|
|
return;
|
|
}
|
|
|
|
if (!TrackedOperationId.TryParse(message.Id, out var trackedId))
|
|
{
|
|
_logger.LogWarning(
|
|
"Cached-call audit observer skipped: message id {MessageId} is not a parseable TrackedOperationId (category {Category}, outcome {Outcome}). " +
|
|
"Audit lifecycle for this operation will have no rows.",
|
|
message.Id, message.Category, outcome);
|
|
return;
|
|
}
|
|
|
|
CachedCallAttemptContext context;
|
|
try
|
|
{
|
|
context = new CachedCallAttemptContext(
|
|
TrackedOperationId: trackedId,
|
|
Channel: channel,
|
|
Target: message.Target,
|
|
SourceSite: _siteId,
|
|
Outcome: outcome,
|
|
RetryCount: message.RetryCount,
|
|
LastError: lastError,
|
|
HttpStatus: httpStatus,
|
|
CreatedAtUtc: message.CreatedAt.UtcDateTime,
|
|
OccurredAtUtc: DateTime.SpecifyKind(occurredAtUtc, DateTimeKind.Utc),
|
|
DurationMs: durationMs,
|
|
SourceInstanceId: message.OriginInstanceName,
|
|
// The buffered message carries the originating script
|
|
// execution's ExecutionId + SourceScript; surface them on the
|
|
// context so the bridge can stamp the retry-loop cached audit
|
|
// rows. Null on rows buffered before this field existed
|
|
// (back-compat).
|
|
ExecutionId: message.ExecutionId,
|
|
SourceScript: message.SourceScript,
|
|
// The buffered message also carries the spawning inbound-API
|
|
// request's ExecutionId; surface it so the bridge stamps it
|
|
// onto the retry-loop cached rows. Null for a non-routed run
|
|
// and on rows buffered before this field existed (back-compat).
|
|
ParentExecutionId: message.ParentExecutionId);
|
|
}
|
|
catch (Exception buildEx)
|
|
{
|
|
// Defensive — record construction shouldn't throw, but the alog.md
|
|
// §7 contract requires this path be exception-safe regardless.
|
|
_logger.LogWarning(buildEx,
|
|
"Failed to build cached-call attempt context for {MessageId}; observer skipped",
|
|
message.Id);
|
|
return;
|
|
}
|
|
|
|
try
|
|
{
|
|
await _cachedCallObserver.OnAttemptCompletedAsync(context, CancellationToken.None)
|
|
.ConfigureAwait(false);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
// alog.md §7 best-effort: an audit observer outage must NEVER be
|
|
// misclassified as a transient delivery failure or corrupt the
|
|
// S&F retry bookkeeping.
|
|
_logger.LogWarning(ex,
|
|
"ICachedCallLifecycleObserver threw for {MessageId} (Outcome {Outcome}); ignored",
|
|
message.Id, outcome);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gets parked messages for central query (Pattern 8).
|
|
/// </summary>
|
|
/// <param name="category">Optional category filter, or null for all categories.</param>
|
|
/// <param name="pageNumber">The page number (1-based).</param>
|
|
/// <param name="pageSize">The page size.</param>
|
|
/// <returns>A tuple of parked messages and the total count.</returns>
|
|
public async Task<(List<StoreAndForwardMessage> Messages, int TotalCount)> GetParkedMessagesAsync(
|
|
StoreAndForwardCategory? category = null,
|
|
int pageNumber = 1,
|
|
int pageSize = 50)
|
|
{
|
|
return await _storage.GetParkedMessagesAsync(category, pageNumber, pageSize);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Retries a parked message (moves back to pending queue).
|
|
///
|
|
/// An operator requeue is a buffer state change and is
|
|
/// replicated to the standby (as a <see cref="ReplicationOperationType.Requeue"/>)
|
|
/// so a failover preserves the operator's retry intent.
|
|
/// The activity-log entry carries the message's true
|
|
/// category rather than a hard-coded one.
|
|
/// The parked row is captured <i>before</i> the local
|
|
/// requeue write rather than re-read after it, so a concurrent
|
|
/// <c>RemoveMessageAsync</c> or <c>DiscardParkedMessageAsync</c> running
|
|
/// between the two storage calls cannot leave the standby in <c>Parked</c>
|
|
/// while the active node has already requeued — we always have the row in
|
|
/// hand for the <c>Requeue</c> replication.
|
|
/// </summary>
|
|
/// <param name="messageId">The identifier of the message to retry.</param>
|
|
/// <returns>True if successfully retried, false otherwise.</returns>
|
|
public async Task<bool> RetryParkedMessageAsync(string messageId)
|
|
{
|
|
var captured = await _storage.GetMessageByIdAsync(messageId);
|
|
if (captured is null || captured.Status != StoreAndForwardMessageStatus.Parked)
|
|
{
|
|
return false;
|
|
}
|
|
|
|
var success = await _storage.RetryParkedMessageAsync(messageId);
|
|
if (!success)
|
|
{
|
|
return false;
|
|
}
|
|
|
|
Interlocked.Increment(ref _bufferedCount);
|
|
|
|
// The active node just rewrote this row to Pending with retry_count = 0
|
|
// and cleared last_error / last_attempt_at (see
|
|
// StoreAndForwardStorage.RetryParkedMessageAsync). Reconstruct the
|
|
// post-requeue state on the captured POCO so the standby applies the
|
|
// same mutations even if a concurrent writer has already deleted the
|
|
// row underneath us.
|
|
captured.Status = StoreAndForwardMessageStatus.Pending;
|
|
captured.RetryCount = 0;
|
|
captured.LastError = null;
|
|
captured.LastAttemptAt = null;
|
|
_replication?.ReplicateRequeue(captured);
|
|
|
|
RaiseActivity("Retry", captured.Category,
|
|
$"Parked message {messageId} moved back to queue");
|
|
return true;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Permanently discards a parked message.
|
|
///
|
|
/// An operator discard is a buffer removal and is replicated
|
|
/// to the standby (as a <see cref="ReplicationOperationType.Remove"/>) so the
|
|
/// discarded message does not reappear after a failover.
|
|
/// The activity-log entry carries the message's true
|
|
/// category rather than a hard-coded one.
|
|
/// </summary>
|
|
/// <param name="messageId">The identifier of the message to discard.</param>
|
|
/// <returns>True if successfully discarded, false otherwise.</returns>
|
|
public async Task<bool> DiscardParkedMessageAsync(string messageId)
|
|
{
|
|
// Capture the category before the row is deleted so the activity log is
|
|
// labelled correctly.
|
|
var message = await _storage.GetMessageByIdAsync(messageId);
|
|
var success = await _storage.DiscardParkedMessageAsync(messageId);
|
|
if (success)
|
|
{
|
|
_replication?.ReplicateRemove(messageId);
|
|
RaiseActivity("Discard", message?.Category ?? StoreAndForwardCategory.ExternalSystem,
|
|
$"Parked message {messageId} discarded");
|
|
}
|
|
return success;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gets buffer depth by category for health reporting.
|
|
/// </summary>
|
|
/// <returns>A dictionary of buffer depths by category.</returns>
|
|
public async Task<Dictionary<StoreAndForwardCategory, int>> GetBufferDepthAsync()
|
|
{
|
|
return await _storage.GetBufferDepthByCategoryAsync();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gets count of S&F messages for a given instance (for verifying survival on deletion).
|
|
/// </summary>
|
|
/// <param name="instanceName">The instance name to query.</param>
|
|
/// <returns>The number of messages originating from the instance.</returns>
|
|
public async Task<int> GetMessageCountForInstanceAsync(string instanceName)
|
|
{
|
|
return await _storage.GetMessageCountByOriginInstanceAsync(instanceName);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Notification Outbox: looks up a buffered message by its id, or <c>null</c> if it
|
|
/// is not (or no longer) in the buffer. <c>Notify.Status</c> uses this to detect a
|
|
/// notification still in transit at the site — central reports it not-found while
|
|
/// the S&F buffer still holds it, which is the site-local <c>Forwarding</c> state.
|
|
/// </summary>
|
|
/// <param name="messageId">The message identifier.</param>
|
|
/// <returns>The message, or null if not found.</returns>
|
|
public async Task<StoreAndForwardMessage?> GetMessageByIdAsync(string messageId)
|
|
{
|
|
return await _storage.GetMessageByIdAsync(messageId);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Raises the S&F activity notification. The
|
|
/// delegate is snapshotted (so a concurrent unsubscribe cannot NRE) and every
|
|
/// subscriber invocation is wrapped so a slow/throwing subscriber (e.g. the site
|
|
/// event log) cannot abort the caller. Crucially, a subscriber exception raised
|
|
/// from <see cref="EnqueueAsync"/> or <c>RetryMessageAsync</c> must NOT be
|
|
/// misclassified as a transient delivery failure — pre-fix it escaped into the
|
|
/// delivery try/catch and caused a successfully delivered message to be buffered
|
|
/// (or its retry count to be bumped). Activity logging is best-effort.
|
|
/// </summary>
|
|
private void RaiseActivity(string action, StoreAndForwardCategory category, string detail)
|
|
{
|
|
var handlers = OnActivity;
|
|
if (handlers == null) return;
|
|
|
|
foreach (var handler in handlers.GetInvocationList().Cast<Action<string, StoreAndForwardCategory, string>>())
|
|
{
|
|
try
|
|
{
|
|
handler(action, category, detail);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_logger.LogWarning(ex,
|
|
"Store-and-forward activity subscriber threw for action {Action}; ignored",
|
|
action);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Result of an enqueue operation.
|
|
/// </summary>
|
|
public record StoreAndForwardResult(
|
|
/// <summary>True if the message was accepted (either delivered immediately or buffered).</summary>
|
|
bool Accepted,
|
|
/// <summary>Unique message ID for tracking.</summary>
|
|
string MessageId,
|
|
/// <summary>True if the message was buffered (not delivered immediately).</summary>
|
|
bool WasBuffered);
|