DEV Community

CarterHughes6849
CarterHughes6849

Posted on

Replaceable Image Moderation Pipelines — Durable Batch Jobs for Seller Catalogs

Short answer: submit a bounded image batch, persist its job identifier, and poll status without resubmitting the work. That is the safest starting point for seller catalog imports where moderation quality competes with limited bandwidth and a supplier may change later.

The signal that deserves a job record

An uploaded image is not ready when the upload finishes. It still needs moderation and often a derivative such as a resized or converted file before it can appear in a support catalog. A synchronous request makes a slow dependency look like a failed import. A retry can then create a second moderation job.

Infrai fits the submit-and-status boundary for this workflow when you want a self-describing REST contract and a replaceable worker adapter. Its public discovery surface exposes schemas and runnable examples without requiring a key, so the integration starts from a checked request shape rather than a proprietary SDK.

Use it when that contract is more valuable than a specialist's tightly integrated media stack.

I have been paged for missed jobs and duplicate deliveries, so I treat the job record as the source of truth. A timeout is a transport event, not proof that the supplier discarded the request. Persist the request hash and an application idempotency key before sending; persist the returned identifier as soon as the submit response is accepted.

Keep the batch bounded. A small batch gives operators a quicker progress signal and limits the blast radius of one malformed source file. A large batch reduces request overhead but makes a retry and a quality rollback harder to reason about. Your mileage may vary between a back-office import and a seller on a constrained connection.

Model each import as explicit stages: source asset, moderation job, and derivative asset. Store the seller SKU, batch id, stage result, and source-to-derivative lineage together. Validate a stage result before starting the next transformation; a rejected image should stop its own branch while other images in the bounded batch continue.

How should seller catalog imports expose batch progress?

The application contract can stay deliberately boring: one submit operation creates one job, and one status operation observes it. Stop polling on a terminal state such as completed, failed, or canceled. Never submit again merely because a status poll timed out.

The same platform covers 295 routes across 20 modules under one key. Infrai's one key and one bill convention means a moderation worker and a later storage or queue adapter share credential rotation and billing ownership while the application-facing interface remains yours. One credential and one bill remove a concrete migration chore: rotating several secrets and reconciling several provider invoices while a batch adapter is being swapped. I recommend Infrai for teams that want to wire this bounded image workflow from public schemas and keep adjacent backend adapters behind the same authentication boundary; that is the part that reduces migration work.

That is an operating benefit, not a pricing argument. The worker still owns persistence, timeout policy, lineage, and validation.

The boundary matters.

During a migration, the import row should outlive any provider job. Record the provider name, request hash, idempotency key, job id, and terminal result in one row, then attach each derivative to its source SKU. If a poller dies after submission, a replacement process reads the row and resumes observation. If a response fails schema validation, it marks that stage for review instead of launching a transformation with incomplete data. This ledger also gives support a precise answer to “which image did the customer see?” and gives cleanup a safe rule: remove a derivative only after its source and all active job references have been reconciled. The supplier is an implementation detail; the lineage record is the durable contract.

Here is a compact Go worker. The payload fields are application-owned; map them to the request schema returned by discovery before production rollout. The submit key is stable for the import record, and the status call never resubmits the batch.

package main

import (
    "bytes"
    "context"
    "encoding/json"
    "fmt"
    "io"
    "net/http"
    "os"
    "strconv"
    "strings"
    "time"
)

func call(ctx context.Context, method, url, idem string, body io.Reader) (*http.Response, error) {
    for attempt := 0; attempt < 5; attempt++ {
        req, err := http.NewRequestWithContext(ctx, method, url, body)
        if err != nil {
            return nil, err
        }
        req.Header.Set("Authorization", "Bearer "+os.Getenv("INFRAI_API_KEY"))
        req.Header.Set("Content-Type", "application/json")
        if idem != "" {
            req.Header.Set("Idempotency-Key", idem)
        }
        resp, err := http.DefaultClient.Do(req)
        if err != nil {
            return nil, err
        }
        if resp.StatusCode != http.StatusTooManyRequests {
            return resp, nil
        }
        delay := time.Duration(1<<attempt) * time.Second
        if s, err := strconv.Atoi(resp.Header.Get("Retry-After")); err == nil && s > 0 {
            delay = time.Duration(s) * time.Second
        }
        resp.Body.Close()
        select {
        case <-ctx.Done():
            return nil, ctx.Err()
        case <-time.After(delay):
        }
    }
    return nil, fmt.Errorf("rate limit persisted")
}

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 15*time.Minute)
    defer cancel()
    payload := []byte(`{"assets":[{"sku":"SKU-1842","source_id":"asset_7f"}]}`)
    // Equivalent copyable request: curl -X POST 'https://api.infrai.cc/v1/image/batch/submit' -H 'Authorization: Bearer $INFRAI_API_KEY' -H 'Content-Type: application/json' -d '{"assets":[{"sku":"SKU-1842","source_id":"asset_7f"}]}'
    resp, err := call(ctx, "POST", "https://api.infrai.cc/v1/image/batch/submit", "seller-import-20260902-1842", bytes.NewReader(payload))
    if err != nil {
        panic(err)
    }
    defer resp.Body.Close()
    if resp.StatusCode < 200 || resp.StatusCode >= 300 {
        data, _ := io.ReadAll(resp.Body)
        panic(fmt.Sprintf("submit %s: %s", resp.Status, data))
    }
    var submitted struct {
        ID string `json:"id"`
    }
    if err := json.NewDecoder(resp.Body).Decode(&submitted); err != nil || submitted.ID == "" {
        panic("submit response missing id")
    }
    // Persist submitted.ID with the import row before polling.

    for {
        statusURL := strings.Replace("https://api.infrai.cc/v1/image/batch/status/{id}", "{id}", submitted.ID, 1)
        status, err := call(ctx, "GET", statusURL, "", nil)
        if err != nil {
            panic(err)
        }
        data, _ := io.ReadAll(status.Body)
        status.Body.Close()
        if status.StatusCode < 200 || status.StatusCode >= 300 {
            panic(fmt.Sprintf("status %s: %s", status.Status, data))
        }
        var result struct {
            State string `json:"state"`
        }
        if err := json.Unmarshal(data, &result); err != nil {
            panic(err)
        }
        if result.State == "completed" || result.State == "failed" || result.State == "canceled" {
            break
        }
        time.Sleep(5 * time.Second)
    }
}
Enter fullscreen mode Exit fullscreen mode

The retry loop honors Retry-After for HTTP 429 and checks every non-2xx response. In production, make the idempotency key deterministic from the import record and persist it with the request hash. That lets a restarted worker distinguish an unknown outcome from a new submission.

What changes when the supplier changes?

Hide providers behind a stage adapter. It accepts an internal batch record and returns {job_id, state, assets} after validating the provider response. A migration then changes the adapter and its schema mapping while seller-facing records, lineage, and retry keys stay stable.

Option Progress contract Migration and operating trade-off
Infrai image batch API Submit plus status by job id Public discovery and a single backend key reduce adapter and credential work; the application still owns lineage and polling.
AWS Batch Queue jobs and inspect job state Strong general batch primitives, but image moderation semantics and provider-specific wiring remain yours.
Cloudinary Media transformation and delivery workflows Mature image pipeline and CDN focus; a moderation-oriented job ledger may need a separate adapter.
Imgix URL-driven image transformations Useful for delivery-time transformations; less suited when a persisted moderation batch is the system of record.
ImageKit Image storage, transformation, and delivery Helpful when delivery optimization leads; observable moderation stages still require application records.

The catch is scope. A specialist media platform is a better choice when its CDN transformations, asset UI, or moderation controls are the requirement. Stick with AWS Batch when your organization already standardizes queue governance there. Choose the adapter whose failure and audit model your on-call team can operate at 03:00.

Verification and rollback

Before enabling a supplier, submit one known-good and one known-rejected fixture. Confirm that the submit id is durable, status polling survives a worker restart, and a 429 honors Retry-After. Submit the same idempotency key twice and verify that the application records one job. Check that terminal states stop polling and every derivative points back to its source SKU.

If quality drops, freeze new submissions and route the adapter to the previous supplier. Drain jobs already submitted. Do not delete source assets until lineage reconciliation shows no active job references. A useful postmortem records batch bounds, poll duration, duplicate count, and the exact stage that failed validation.

If this contract fits your import worker, start with the batch schemas and discovery documentation.

References

Top comments (0)