How the asynchronous ingestion pipeline is wired for engineers: the entry point and every step, with the services it uses and the concrete resources it touches — database tables, SQL stored procedures, blob & PVC storage, Redis keys, SignalR, Hangfire and domain events.
The pipeline is asynchronous and choreographed: each step is a Hangfire background command that runs in its own transaction, records its outcome, pushes live progress over SignalR, and enqueues the next command. There is no central orchestrator — the chain advances itself, and a business failure simply stops it.
Every step inherits InboundDeliveryIngestionStepHandler<TCommand>, which runs the same lifecycle around the step's own ExecuteAsync:
GetWithStepsAsync (with steps + creator).HaltIfStopRequestedAsync reads the Redis stop flag; persists Stopped if set.Success / Skipped / BusinessFailure / NoOp (Success/Skipped may carry a WarningPayload).RecordStep* (persists the step row + any WarningPayload and the acting UserId) + SaveChangesAsync (dispatches domain events).Each command runs under TransactionResultCommandBehavior (ReadCommitted, 300 s). Transient exceptions
propagate and Hangfire retries (10×); a business failure records the step as failed and stops the chain (no next command).
Two views of the same machine: the control flow of commands and events, and the status the aggregate moves through.
→ command enqueued (Hangfire) · ⏱ delayed schedule (grace) · ⇢ domain / integration event. Each step is its own background transaction.
CreateInboundDeliveryIngestionCommand
→create
ingestion→ queued
⇢post-commit
…QueuedIntegrationEvent
→enqueue
StartPipelineWhenIngestionQueuedEventHandler
AwaitGracePeriodStepHandler
⏱schedule
CreateInboundDeliveryStepCommandjob id kept in Redis
CreateInboundDeliveryStepCommand
→handler
InboundDelivery.Create()
IInboundOrderSupplier.SupplyAsynccalled directly by the converge
→if overflow
inbound ordergenerated
InboundDeliveryCompletedInternalEvent
→
PickTower loadPT_Load_InboundDelivery
IInboundOrderSupplier.SupplyAsync directly (after the delivery's intermediate save) and records the returned order inline (→ IBO generated) — no follow-up command, no Succeeded→IboGenerated flash. InboundDeliveryCreatedInternalEvent is raised with SupplyInboundOrders=false, so SupplyInboundOrdersEventHandler skips ingestion-created deliveries (it still serves sync / acquisition / resell).InboundOrderCompletedIntegrationEvent
→
SucceedIngestion…Handler
→
ingestion→ succeeded
The aggregate's status and what drives each transition. Dashed blue = user re-entry.
resume continues the same attempt — from the paused step, or re-arms the grace window if paused there · restart begins a new attempt. Failed and Stopped are deletable. Internally the converge completes the ingestion and records any generated order inline in the same transaction — it lands directly in IBO generated when an order was produced.
Everything that can move an ingestion forward.
POST /DeliveriesManagement/Inbounds/IngestionsCreateInboundDeliveryIngestionCommand creates the aggregate in Queued.InboundDeliveryIngestionQueuedIntegrationEvent → StartPipelineWhenIngestionQueuedEventHandlerIngestionStepCommands.For(FromStepTypeCode ?? ParseRawPayload) — step 1 (ParseRawPayloadStepCommand) for create/restart, the paused step for a mid-pipeline resume. This is what actually starts the pipeline — nothing runs inside the HTTP request.EnqueueNext<Name>StepCommand via Hangfire. The chain is self-advancing.IBackgroundService.Schedule(CreateInboundDeliveryStepCommand, grace)ibd-ingestion-grace-job:{id}). A grace pause deletes the pending job (best-effort); Resume deletes any stale job and schedules a fresh one. The converge itself guards on status only (still InGracePeriod).POST {id}/stop · /resume · /restart · DELETE {id}InboundOrderCompletedIntegrationEvent → SucceedIngestionWhenGeneratedOrderCompletedEventHandlerIBO generated → Succeeded.PurgeStaleIngestionRawPayloadsCronJob (daily 03:00)Succeeded / IboGenerated ingestions older than the retention window.ReleaseStaleIngestionKeyReservationsCronJob (hourly · "release-stale-ingestion-reservations")Failed ingestions idle 24h+, and terminally fails Stopped / in-grace ingestions parked beyond the 30-day reservation horizon (their claims have expired, so converging would be unsafe).Referenced by every step — listed once here rather than repeated below.
IInboundDeliveryIngestionStorage at AppSettings.InboundDeliveryStoragePrefix.{prefix}/ingestions/{id}/attempt-{n}/ → {n}.raw, parsed.json, preparedkeys.json.{prefix}.…:ibd-ingestion:{id} — progress snapshot (24 h).…:ibd-ingestion-stop:{id} — stop flag (6 h); carries requestedBy so the persisting step can audit who requested the stop.…:ibd-key-reservation:{hash} + …-owned:{id} — key reservations (Lua, all-or-nothing, 30 days; explicit releases do the real cleanup).…:ibd-ingestion-grace-job:{id} — Hangfire job id of the scheduled converge (grace pause deletes the job).IngestionStepType.MarkerFloor = 1000 and above, which is the boundary the tester filters on.lock:ibd-ingestion-upload:{id} — upload lock (30 s).IIngestionProgressService broadcasts InboundDeliveryIngestionChanged (the ingestion's shape — status, attempt, current attempt's steps as bare codes; no error payloads) on WholeSaleHub (/notifications) via INotificationService.BroadcastToGroup.inbound-delivery-ingestions, joined via SubscribeToInboundDeliveryIngestions and gated on the same roles as the REST endpoints. InboundDeliveryCompleted rides the same group. INotificationService no longer exposes a broadcast-to-all at all — only BroadcastToGroup and SendNotification (one user) — since every logged-in app user holds an open connection.PublishStepProgressAsync → InboundDeliveryIngestionTick: counters only, no cache read, no snapshot rebuild.ErrorPayload, previous attempts) are never broadcast — clients read GET …/Ingestions/{id} on demand, and only when the summary shows the panel's content changed (new attempt, a step's HasError/HasWarning, a new IBD/IBO reference).ProjectTo<InboundDeliveryIngestionSummaryResponse> instead of loading the aggregate — ErrorPayload is only tested in SQL (HasError = a CASE collapsed to a bit) and the attempt filter is pushed into the join, so earlier attempts' rows are never read.RequestStopAsync now carries requestedBy; GetStopRequestedByAsync reads it back so a step persisting Stopped can audit who asked.IBackgroundService.Enqueue (next step) / Schedule (delayed converge — returns the Hangfire job id) / DeleteJob (drops a queued job, used by grace pause/resume), routed by MediatRHangfireBridge over Redis.HangfireFailedJobNotification email.WarningPayload (JSON list, catalogued in IngestionWarnings) via StepOutcome.Success(warningPayload) / Skipped(warningPayload). RecordStepSucceeded / RecordStepSkipped persist it onto the IngestionStepExecution.WarningPayload column; the run continues. Surfaced on InboundDeliveryIngestionProgressResponse / IngestionStepProgressDto and rendered as amber blocks in the tester.archive_entries_skipped (Unpack), duplicate_keys_accepted (Check keys), keys_near_limit (Generate keys), origin_keys_demoted (Converge).UserId. Marker rows are stamped with who acted: Stop/Pause (from the Redis stop flag's requestedBy, or set directly for a grace-pause), Resume, Restart; executed step rows carry the pipeline initiator. The tester shows a “by <user> · <time>” tooltip on the marker pills.IngestionStepWarningsAndAuditUser_Deliveries — nullable WarningPayload (nvarchar(max)) + UserId (uniqueidentifier).What a running step touches. The handler is a Hangfire worker; everything else is a backing system it reads from or writes to.
Per step: what it does, the services it depends on, the concrete things it calls, the artifacts it reads/writes, the events it raises, and how it can fail.
Stop / resume / restart / delete — and how each touches state vs. cache.
While running, sets the Redis stop flag only — no aggregate write, zero concurrency risk; the next step's start check catches it and persists Stopped. The flag carries requestedBy (read via GetStopRequestedByAsync) so the Stopped marker row is stamped with the requesting user. During the grace window (no step running) it persists Stopped synchronously via PauseGracePeriod (sole writer, user set directly) and deletes the pending Hangfire converge job outright (best-effort, job id read from Redis).
Clears the stop flag (after EnsureCanBeResumed; Failed → 400, use Restart). Paused mid-pipeline → Resume() (→ Running, same attempt) publishes the queued event with FromStepTypeCode = GetStepToResumeFrom(), so the paused step is enqueued directly and no completed step is re-run. Paused in grace → deletes any lingering converge job, calls ResumeGracePeriod (prepared keys reused) and Schedules a fresh converge, storing its job id in Redis; a surviving old converge no-ops on the status guard.
Clears the stop flag and calls RestartFromBeginning — bumps AttemptNumber and re-runs from step 1; the prior attempt is kept as history.
Removes the aggregate, clears the cache + broadcasts InboundDeliveryIngestionRemoved, releases any Redis key reservations, and sweeps the PVC attempt folder.
The converge records any generated IBO inline (returned synchronously by IInboundOrderSupplier.SupplyAsync, which it calls directly after the delivery's intermediate save) and moves the ingestion to IBO generated — no follow-up command. When the user completes that order, InboundOrderCompletedIntegrationEvent → SucceedIngestionWhenGeneratedOrderCompletedEventHandler → Succeeded.