⚙ Ingestion — engineering reference ↗ UX guide
Wholesale · Deliveries · internals

Inbound Delivery Ingestion — engineering reference

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.

CQRS · MediatRHangfire (Redis)EF Core · SQL Server Redis (StackExchange)SignalRAzure BlobPVC file storageMassTransit · RabbitMQ

1 · High-level overview

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:

  1. Load the aggregate — GetWithStepsAsync (with steps + creator).
  2. Stale guard — drop if superseded attempt / paused / terminal.
  3. Stop check — HaltIfStopRequestedAsync reads the Redis stop flag; persists Stopped if set.
  4. Idempotency — skip if this step+attempt already recorded.
  5. Publish running — Redis snapshot + SignalR broadcast.
  6. ExecuteAsync — the step's real work (below).
  7. Classify outcome — Success / Skipped / BusinessFailure / NoOp (Success/Skipped may carry a WarningPayload).
  8. Record + save — RecordStep* (persists the step row + any WarningPayload and the acting UserId) + SaveChangesAsync (dispatches domain events).
  9. Publish the new snapshot, then EnqueueNext (unless failed / terminal).

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

2 · Flow & lifecycle

Two views of the same machine: the control flow of commands and events, and the status the aggregate moves through.

Command & event choreography

→ command enqueued (Hangfire) · ⏱ delayed schedule (grace) · ⇢ domain / integration event. Each step is its own background transaction.

command event handler state
1 · create & queue one transaction
clientPOST submit →send CreateInboundDeliveryIngestionCommand →create ingestion→ queued ⇢post-commit …QueuedIntegrationEvent →enqueue StartPipelineWhenIngestionQueuedEventHandler
↓
2 · pipeline Hangfire · each step its own transaction
1parse› 2unpack› 3partner› 4products› 5generate keys› 6reserve keys› 7check keys› 8upload keys› 9await grace ⏱after grace 10create
Every step's handler records the step → pushes progress (SignalR + Redis) → enqueues the next command. A business failure stops the chain.
↓
3 · grace window pause-able · nothing created yet
AwaitGracePeriodStepHandler ⏱schedule CreateInboundDeliveryStepCommandjob id kept in Redis
stop → Stopped + the pending Hangfire job is deleted (best-effort) · resume → deletes any stale job, schedules a fresh converge and stores its new job id · the converge runs only if status is still in-grace (status-only guard).
↓
4 · converge step 10 · one transaction
CreateInboundDeliveryStepCommand →handler InboundDelivery.Create()
IInboundOrderSupplier.SupplyAsynccalled directly by the converge →if overflow inbound ordergenerated
InboundDeliveryCompletedInternalEvent → PickTower loadPT_Load_InboundDelivery
ingestion→ succeeded / IBO generated
The converge calls 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).
↓
5 · generated order → success
ingestionIBO generated →user fills it InboundOrderCompletedIntegrationEvent → SucceedIngestion…Handler → ingestion→ succeeded

Ingestion status lifecycle

The aggregate's status and what drives each transition. Dashed blue = user re-entry.

Ingestion status state machine Queued advances to Running, then In grace, then Succeeded; an overflow order routes through IBO generated; Running can branch to Failed or Stopped; re-entry commands return Stopped or Failed to Running or In grace. first step grace · step 9 converge · no order converge · order order completed business fail restart stop resume pause resume Queued Running In grace Succeeded IBO generated Failed Stopped

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.

3 · Triggers & entry points

Everything that can move an ingestion forward.

4 · Shared infrastructure

Referenced by every step — listed once here rather than repeated below.

PVC Ingestion storage

  • IInboundDeliveryIngestionStorage at AppSettings.InboundDeliveryStoragePrefix.
  • Layout: {prefix}/ingestions/{id}/attempt-{n}/ → {n}.raw, parsed.json, preparedkeys.json.
  • Uploaded / extracted key files: flat temp files under {prefix}.

Redis db 1 keys

  • …: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).
  • Step codes are the running order (10, 20 … 100, spaced by 10 to leave room for insertions); markers sit at IngestionStepType.MarkerFloor = 1000 and above, which is the boundary the tester filters on.
  • lock:ibd-ingestion-upload:{id} — upload lock (30 s).

SignalR Live progress

  • 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.
  • Group 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.
  • Steps that report own progress call PublishStepProgressAsync → InboundDeliveryIngestionTick: counters only, no cache read, no snapshot rebuild.
  • The unbounded parts (per-key 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).
  • No sequence numbers: a change carries the ingestion's whole current shape, not a diff, so a lost message is superseded by the next one and a reconnect re-reads the list.
  • The list endpoints 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.

Hangfire Background

  • 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.
  • Retry policy: 10×; permanent failure → HangfireFailedJobNotification email.

DB Warnings & audit (step row)

  • Warnings — a step can attach a 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.
  • Four shipped: archive_entries_skipped (Unpack), duplicate_keys_accepted (Check keys), keys_near_limit (Generate keys), origin_keys_demoted (Converge).
  • Audit — step & marker rows carry a 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.
  • Both columns added by migration IngestionStepWarningsAndAuditUser_Deliveries — nullable WarningPayload (nvarchar(max)) + UserId (uniqueidentifier).

Systems map

What a running step touches. The handler is a Hangfire worker; everything else is a backing system it reads from or writes to.

Ingestion systems map The ingestion step handler reads and writes SQL Server, Redis, Azure Blob, PVC file storage, publishes to SignalR and RabbitMQ, and is driven by Hangfire. EF Core · tables + SPs cache · reservations · locks enqueue / schedule key files integration events live progress artifacts + temp SQL Server · EF Core deliveries.InboundDeliveryIngestion · .Key contactBook.Partner · products.Products SP PT_CheckDuplicates_Keys · PT_Load_… SeedWork.Counter (sequence) Redis · db 1 ibd-ingestion (snapshot, 24h) ibd-ingestion-stop (6h) ibd-key-reservation (Lua, 30d) ibd-ingestion-grace-job · upload lock Hangfire · Redis enqueue next · schedule converge Azure Blob container 'keys' · file keys + MD5 MassTransit · RabbitMQ integration events (post-commit) …QueuedIntegrationEvent InboundOrderCompletedIntegrationEvent SignalR → clients WholeSaleHub /notifications …IngestionUpdated PVC · file storage ingestions/{id}/attempt-{n}/ 1.raw · parsed.json · preparedkeys.json + temp files · keys mirror Ingestion step handler Hangfire worker · 1 txn / step

5 · Entry point & the ten steps

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.

DBSQL Server table SQLstored proc / sequence Redis BlobAzure PVCfile storage SignalR Hangfire Eventdomain / integration

6 · Re-entry commands

Stop / resume / restart / delete — and how each touches state vs. cache.

Stop / pause

POST {id}/stop

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

Resume

POST {id}/resume · stopped only

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.

Restart

POST {id}/restart · any non-succeeded

Clears the stop flag and calls RestartFromBeginning — bumps AttemptNumber and re-runs from step 1; the prior attempt is kept as history.

Delete

DELETE {id} · stopped / failed

Removes the aggregate, clears the cache + broadcasts InboundDeliveryIngestionRemoved, releases any Redis key reservations, and sweeps the PVC attempt folder.

Finalize (generated order)

inline + integration event

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.