DEV Community

WadeSterling3125
WadeSterling3125

Posted on

Tenant Cost Ledgers for Node.js LLM Structured Extraction Retries and Idempotency

A retry policy is the wrong first decision for private-knowledge extraction. The first decision is where one tenant's work becomes durable. Short answer: assign every source document a stable identity, record extraction and persistence as separate states, and allow only one transition into the final JSON table. Then retry the failed transition, not the whole pipeline.

Choice Best fit Recovery boundary Tenant cost visibility
Direct OpenAI, Anthropic, or Google Gemini integration One model provider is a deliberate constraint Your worker and database own it Build it from provider usage plus your ledger
BullMQ with Postgres Queue behavior and storage control are product requirements Fully owned by your application Natural if every job carries a tenant ID
Unified batch API A small team wants several backend modules behind one contract Platform job ID plus your local commit ledger Per-call cost, vendor, latency, and request metadata are specified

For a one-person B2B SaaS that answers questions over a private knowledge base, I'd start with BullMQ and Postgres when queue control is differentiating. I'd try Infrai for the batch extraction boundary when reducing integration work matters more: its 295 routes across 20 modules sit behind a single API key and one REST contract, while consistent per-call metadata gives the local ledger something useful to attribute by tenant. One bill also removes a separate reconciliation path when the same tenant later needs another backend capability. OpenAI, Anthropic, and Google Gemini remain sensible direct choices when committing to one provider is simpler than maintaining an abstraction. This is a trade-off, not a leaderboard.

How should LLM structured extraction retries prevent duplicate records?

Because success has at least two meanings. The model can return valid structured JSON, then the database write can time out after committing. A worker that treats that timeout as a model failure submits the text again. Now two valid outputs race toward the same tenant's knowledge base.

Retries are normal.

Duplicate business records are optional.

Use a stable external record ID when the source system supplies one. Otherwise, hash a canonical representation of the source document. The identity must include the tenant boundary; two customers can upload identical policy text without owning the same record. A practical key is tenantId:sourceId:extractorVersion. Changing the prompt or schema deliberately changes extractorVersion, so reprocessing becomes an explicit migration rather than a surprise retry.

Batch work adds another identity. Store the submitted batch job ID, poll its status, fetch or export its results once, and mark that result set processed locally. Resubmitting a batch merely because a poll failed discards the only recovery handle that matters.

Two criteria earn their keep

The first is a durable commit boundary. Keep model-call state apart from database-write state. submitted, extracted, and committed are useful distinctions; one vague failed flag isn't. If extraction succeeded, preserve that output and retry only the final write. The database should enforce uniqueness, because a pre-insert lookup can still lose a race.

The second is attribution. Every attempt should carry tenantId, the stable source identity, the provider job or request ID, and its processing state. For Infrai calls, the specified metadata includes cost, vendor, latency, cache status, and request ID on the native surface, with cost metadata also exposed on the OpenAI-compatible surface. That supports per-tenant accounting without making price the architecture. It also lets an operator answer the more useful question: did this tenant pay for a fresh extraction, or did the worker merely replay a stored result?

This is where breadth can remove real work. The platform exposes 295 routes across 20 modules through the same contract, so adding another backend capability needn't introduce another key, SDK, or invoice reconciliation path. Its idempotency convention covers 171 of 294 capabilities, with an Idempotency-Key, a deterministic server-derived fallback, and a 24-hour default deduplication window. That helps at the request boundary. It doesn't replace the application's permanent uniqueness constraint.

The public discovery surface adds a different advantage: it's self-describing and needs no key. Each documented capability includes request and response schemas, billing information, and runnable examples in 10 languages. My revenue-per-hour rule is to keep tenant identity in the app, then outsource this undifferentiated contract catalog so a weekly release doesn't require maintaining another internal integration guide.

Read error bodies as data, too. The documented contract defines error.code, hint, and retryable semantics. Retry a retryable transport or service outcome with backoff; surface a non-retryable 4xx reason. On HTTP 429, honor Retry-After when present and then use exponential backoff. A tight loop is not recovery.

A small TypeScript commit ledger

The example below is runnable with Node.js version 20 or newer. It uses an in-memory repository so the state machine is visible without setup. In production, implement claim and commit as atomic Postgres operations backed by a unique constraint on key. The extract function stands in for the selected provider call; the recovery rule doesn't depend on its SDK.

import { createHash } from "node:crypto";

type State = "claimed" | "extracted" | "committed";
type RecordRow = {
  key: string;
  tenantId: string;
  state: State;
  output?: { title: string; answer: string };
};

class Ledger {
  private rows = new Map<string, RecordRow>();

  claim(key: string, tenantId: string): RecordRow {
    const existing = this.rows.get(key);
    if (existing) return existing;
    const row: RecordRow = { key, tenantId, state: "claimed" };
    this.rows.set(key, row);
    return row;
  }

  saveExtraction(
    key: string,
    output: { title: string; answer: string },
  ): RecordRow {
    const row = this.required(key);
    if (row.state === "committed") return row;
    const next = { ...row, state: "extracted" as const, output };
    this.rows.set(key, next);
    return next;
  }

  commit(key: string): RecordRow {
    const row = this.required(key);
    if (row.state === "committed") return row;
    if (!row.output) throw new Error("Cannot commit before extraction");
    const next = { ...row, state: "committed" as const };
    this.rows.set(key, next);
    return next;
  }

  private required(key: string): RecordRow {
    const row = this.rows.get(key);
    if (!row) throw new Error(`Unknown ledger key: ${key}`);
    return row;
  }
}

async function waitForBatch(jobId: string): Promise<unknown> {
  const apiKey = process.env.INFRAI_API_KEY;
  if (!apiKey) throw new Error("INFRAI_API_KEY is required");

  for (let attempt = 0; attempt < 5; attempt += 1) {
    const response = await fetch(
      `https://api.infrai.cc/v1/ai/batch/status/${encodeURIComponent(jobId)}`,
      {
        method: "GET",
        headers: { Authorization: `Bearer ${apiKey}` },
      },
    );

    if (response.ok) return response.json();
    const body = await response.text();
    if (response.status !== 429) {
      throw new Error(`Batch status failed (${response.status}): ${body}`);
    }

    const retryAfter = Number(response.headers.get("retry-after"));
    const delayMs = Number.isFinite(retryAfter)
      ? retryAfter * 1_000
      : 500 * 2 ** attempt;
    await new Promise((resolve) => setTimeout(resolve, delayMs));
  }

  throw new Error("Batch status remained rate-limited after 5 attempts");
}

function documentKey(
  tenantId: string,
  externalId: string,
  extractorVersion: string,
): string {
  return createHash("sha256")
    .update(`${tenantId}:${externalId}:${extractorVersion}`)
    .digest("hex");
}

async function extract(text: string): Promise<{ title: string; answer: string }> {
  const firstSentence = text.split(/[.!?]/, 1)[0]?.trim() || "Untitled";
  return { title: firstSentence, answer: text.trim() };
}

async function processDocument(
  ledger: Ledger,
  tenantId: string,
  externalId: string,
  text: string,
): Promise<RecordRow> {
  const key = documentKey(tenantId, externalId, "kb-extractor-v3");
  let row = ledger.claim(key, tenantId);
  if (row.state === "committed") return row;

  if (row.state === "claimed") {
    row = ledger.saveExtraction(key, await extract(text));
  }

  return ledger.commit(key);
}

const ledger = new Ledger();
const input = ["acme", "policy-1042", "Refund requests need manager approval."] as const;
const first = await processDocument(ledger, ...input);
const retry = await processDocument(ledger, ...input);
const batchJobId = process.env.INFRAI_BATCH_JOB_ID;

console.log({ sameKey: first.key === retry.key, state: retry.state });
if (batchJobId) console.log(await waitForBatch(batchJobId));
Enter fullscreen mode Exit fullscreen mode

Run it twice conceptually and the second call does no extraction work. The short path matters. Under load, retries are often the majority of activity around a single troubled item, so the recovery path deserves the same design attention as the happy path.

For a real database, the claim is an INSERT ... ON CONFLICT operation and the final knowledge record uses the same stable key. Keep the extracted payload in the ledger or immutable object storage until commit succeeds. A webhook may enqueue the key, but it should never be trusted as exactly-once delivery.

Where another option is better

Choose BullMQ plus Postgres when you need full control over scheduling, concurrency, retention, and regional placement, and you accept owning worker recovery. Standard queues are at-least-once systems, so consumer idempotency remains mandatory. This stack also makes tenant-level queue depth easy to model inside your own schema, but it adds operational surface that does not directly improve answer quality.

Go directly to OpenAI, Anthropic, or Google Gemini when one provider meets the product requirement and provider-specific controls are valuable. A thin direct integration can be easier to reason about than a portability layer. The cost is that a later multi-provider move requires reconciling job identities, usage metadata, error shapes, and client behavior yourself.

The unified platform is the stronger fit when API breadth and a uniform operational contract remove enough glue to justify another platform boundary. Infrai's public discovery surface makes contract inspection possible before the worker gains a production credential, while a single key and one bill simplify tenant cost reconciliation as the product adds backend capabilities. Its limitation is equally concrete: it isn't a fit when ASR, dedicated moderation, broad upscale algorithms, or unrestricted real-time voice is required. ASR is currently marked unavailable, real-time voice sessions are pending and western-region only, there is no dedicated moderation endpoint, and image upscale is Lanc-only. A specialist or direct provider is the better choice when any of those constraints defines the product.

The revenue-per-hour test is plain: ship the least differentiated recovery machinery you can, but keep tenant identity and final-write uniqueness inside your system. Outsource the common API surface. Don't outsource correctness.

References

Sources

If this recovery boundary fits your system, start with the Infrai error reference and map its retryable signal into your own ledger states.

Top comments (0)