Large pushes and bulk files
A push of up to 50 records is processed while you wait and answered 200 with every record’s outcome (Pushes and mandates). A larger push — up to 10,000 records — is accepted, then processed: MasterDB checks everything that could refuse the whole batch, answers 202 Accepted at once, and checks and publishes the records in the background, 200 at a time. For more than 10,000 records, upload a bulk file.
Every step below has a runnable sample, in TypeScript and Python, that runs against the sandbox (The sandbox). They share one small file for your key, the product record, sealing and polling; the samples use it, and you can copy it into your own system.
The shared helpers: the key, a product, sealing, sending and polling
// The pieces the large-push samples share: your integration key and client, a product record to send, sealing// a batch, sending it, and reading a push's status until it settles (/businesses/large-pushes/).//// MASTERDB_INTEGRATION_KEY_FILE the integration key's private half (PKCS #8 PEM), as in the push sample// MASTERDB_BUSINESS_UUID your business's public identifier// MASTERDB_AI_POLICY_VERSION your AI policy's version in force (GET /v1/seal-context: ai_policy_version)import { createPrivateKey } from 'node:crypto';import { readFileSync } from 'node:fs';import { type SealedBatch, SANDBOX, createBusinessClient, createEd25519Signer, createPublicClient, idempotencyKey, sealBatch } from '@masterdb/client';
export const baseUrl = process.env.MASTERDB_API_URL ?? SANDBOX.business;export const signer = createEd25519Signer(createPrivateKey(readFileSync(process.env.MASTERDB_INTEGRATION_KEY_FILE as string)));export const business = createBusinessClient({ baseUrl, signer });
// The seal names the certificate in force: read its cert_id from the public certificate endpoint.const cert = await createPublicClient({ baseUrl }).GET('/v1/certificates/{uuid}', { params: { path: { uuid: process.env.MASTERDB_BUSINESS_UUID as string } } });if (cert.error) throw new Error(cert.error.code);if (cert.data.cert_id === undefined) throw new Error('the business has no certificate in force');const certId = cert.data.cert_id;
/** * Product `n` of a series. The same series and number are always the same product (the same `business_product_id`), * so a sample that runs again publishes new versions of the same products instead of filling the sandbox. * Money is a decimal string, never a number. */export function product(series: string, n: number) { const code = `EMBERS-${series}${String(n).padStart(4, '0')}`; return { schema: 'masterdb/products/1', language: 'en', business_product_id: code, product_name: `Ember Spindle Roof Box ${series}${n}`, countries: ['US', 'GB'], vertical: 'automotive', category: 'parts_accessories', channel: 'online', brand: 'Ember Spindle', tags: ['sandbox', 'automotive'], short_description: 'A fictional roof box from the MasterDB sandbox documentation. It does not exist.', prices: [ { country: 'US', currency: 'USD', amount: '189.00' }, { country: 'GB', currency: 'GBP', amount: '149.00' }, ], on_sale: false, availability: 'available', product_url: `https://ember-spindle.sandbox.masterdb.ai/products/${code.toLowerCase()}`, };}
/** One seal over the batch. The bytes you seal are the bytes you send: each record is serialised once, here. */export function seal(records: object[]): Promise<SealedBatch> { return sealBatch({ records: records.map((r) => new TextEncoder().encode(JSON.stringify(r))), signer, certId, aiPolicyVersion: Number(process.env.MASTERDB_AI_POLICY_VERSION ?? 0), seq: Date.now(), // must rise with every batch this key sends: a replayed older batch is refused recordType: 'products', });}
/** * Sends a sealed batch (a new request signature each time). Up to 50 records are processed while you wait and * answered 200; more are accepted and answered 202. Either way the answer names the push: read it from there. */export async function send(body: SealedBatch): Promise<{ pushId: string; http: number }> { const answer = await business.POST('/v1/publish/{type}', { params: { path: { type: 'products' }, header: { 'Idempotency-Key': idempotencyKey() } }, body }); if (answer.error) throw new Error(`${answer.error.code}: ${answer.error.detail ?? ''}`); if (answer.data.push_id === undefined) throw new Error('the answer names no push'); return { pushId: answer.data.push_id, http: answer.response.status };}
/** The push's status and one page of its records' outcomes; `outcome` keeps only the records with that outcome. */export async function read(pushId: string, query: { outcome?: 'accepted' | 'updated' | 'rejected' | 'held' | 'pending'; cursor?: string } = {}) { const page = await business.GET('/v1/pushes/{push_id}', { params: { path: { push_id: pushId }, query } }); if (page.error) throw new Error(`${page.error.code}: ${page.error.detail ?? ''}`); return page.data;}
/** * Reads the status every `poll_after_seconds`, for as long as the answer carries that member; then the push has * settled (`complete`, `held` or `failed`). A push that `failed` or is `stalled` (no step for ten minutes) is * resumed by sending the exact batch again, which `resume` does. */export async function settle(pushId: string, resume?: () => Promise<unknown>) { let resumed = 0; for (;;) { const status = await read(pushId); console.log(` ${status.state}: ${status.progress.stage} ${status.progress.chunk}/${status.progress.chunks}, ${status.counts.pending} of ${status.records} records pending`); const stuck = status.state === 'failed' || status.stalled === true; if (stuck && resume !== undefined && resumed < 3) { resumed += 1; console.log(` ${status.state === 'failed' ? `stopped: ${status.error?.code ?? ''}` : 'stalled'}: sending the same batch again`); await resume(); } else if (stuck || status.poll_after_seconds === undefined) { return status; } await new Promise((r) => setTimeout(r, (status.poll_after_seconds ?? 5) * 1000)); }}
/** Every record's outcome, in leaf order, 100 a page. */export async function outcomes(pushId: string, outcome?: 'accepted' | 'updated' | 'rejected' | 'held' | 'pending') { const results = []; for (let cursor: string | undefined; ; ) { const page = await read(pushId, { ...(outcome === undefined ? {} : { outcome }), ...(cursor === undefined ? {} : { cursor }) }); results.push(...page.results); cursor = page.next_cursor; if (cursor === undefined) return results; }}"""The pieces the large-push samples share: your integration key and client, a product record to send, sealing abatch, sending it, and reading a push's status until it settles (/businesses/large-pushes/).
MASTERDB_INTEGRATION_KEY_FILE the integration key's private half (PKCS #8 PEM), as in the push sample MASTERDB_BUSINESS_UUID your business's public identifier MASTERDB_AI_POLICY_VERSION your AI policy's version in force (GET /v1/seal-context: ai_policy_version)"""
import jsonimport osimport timeimport uuid
import httpxfrom masterdb_signing import SANDBOX_API, MasterDBAuth, business_tag, load_key, seal_batch
base_url = os.environ.get("MASTERDB_API_URL", SANDBOX_API)key = load_key(os.environ["MASTERDB_INTEGRATION_KEY_FILE"])# A push of up to 500 records may be processed inside its request where background processing is not set up.business = httpx.Client(base_url=base_url, auth=MasterDBAuth(key, tag=business_tag), timeout=90)
# The seal names the certificate in force: read its cert_id from the public certificate endpoint._cert = httpx.get(f"{base_url}/v1/certificates/{os.environ['MASTERDB_BUSINESS_UUID']}")_cert.raise_for_status()cert_id = _cert.json()["cert_id"]
def product(series, n): """Product ``n`` of a series. The same series and number are always the same product (the same ``business_product_id``), so a sample that runs again publishes new versions of the same products instead of filling the sandbox. Money is a decimal string, never a number.""" code = f"EMBERS-{series}{n:04d}" return { "schema": "masterdb/products/1", "language": "en", "business_product_id": code, "product_name": f"Ember Spindle Roof Box {series}{n}", "countries": ["US", "GB"], "vertical": "automotive", "category": "parts_accessories", "channel": "online", "brand": "Ember Spindle", "tags": ["sandbox", "automotive"], "short_description": "A fictional roof box from the MasterDB sandbox documentation. It does not exist.", "prices": [{"country": "US", "currency": "USD", "amount": "189.00"}, {"country": "GB", "currency": "GBP", "amount": "149.00"}], "on_sale": False, "availability": "available", "product_url": f"https://ember-spindle.sandbox.masterdb.ai/products/{code.lower()}", }
def seal(records): """One seal over the batch. The bytes you seal are the bytes you send: each record is serialised once, here.""" return seal_batch( [json.dumps(r, separators=(",", ":"), ensure_ascii=False).encode("utf-8") for r in records], key, cert_id=cert_id, ai_policy_version=int(os.environ.get("MASTERDB_AI_POLICY_VERSION", "0")), seq=time.time_ns() // 1_000_000, # must rise with every batch this key sends: a replayed older batch is refused )
def send(body): """Sends a sealed batch (a new request signature each time). Up to 50 records are processed while you wait and answered 200; more are accepted and answered 202. Either way the answer names the push: read it from there.""" answer = business.post("/v1/publish/products", json=body, headers={"Idempotency-Key": str(uuid.uuid4())}) if answer.status_code not in (200, 202): raise SystemExit(f"{answer.json()['code']}: {answer.json().get('detail', '')}") return answer.json()["push_id"], answer.status_code
def read(push_id, **query): """The push's status and one page of its records' outcomes; ``outcome=`` keeps only the records with that outcome.""" page = business.get(f"/v1/pushes/{push_id}", params=query) page.raise_for_status() return page.json()
def settle(push_id, resume=None): """Reads the status every ``poll_after_seconds``, for as long as the answer carries that member; then the push has settled (``complete``, ``held`` or ``failed``). A push that ``failed`` or is ``stalled`` (no step for ten minutes) is resumed by sending the exact batch again, which ``resume`` does.""" resumed = 0 while True: status = read(push_id) c, p = status["counts"], status["progress"] print(f" {status['state']}: {p['stage']} {p['chunk']}/{p['chunks']}, {c['pending']} of {status['records']} records pending") stuck = status["state"] == "failed" or status.get("stalled") is True if stuck and resume is not None and resumed < 3: resumed += 1 print(f" {'stopped: ' + (status['error'] or {}).get('code', '') if status['state'] == 'failed' else 'stalled'}: sending the same batch again") resume() elif stuck or "poll_after_seconds" not in status: return status time.sleep(status.get("poll_after_seconds", 5))
def outcomes(push_id, outcome=None): """Every record's outcome, in leaf order, 100 a page.""" results, cursor = [], None while True: query = {k: v for k, v in (("outcome", outcome), ("cursor", cursor)) if v is not None} page = read(push_id, **query) results.extend(page["results"]) cursor = page.get("next_cursor") if cursor is None: return resultsWhat is checked before the answer
Section titled “What is checked before the answer”Everything that refuses a batch is checked before you get an answer, exactly as for a small push: the request signature, your key and its mandate (record type, countries, caps, source allow-list), that your business is verified, the batch seal, that the records are the ones the seal names (the Merkle root), your certificate and AI policy version, seq, and your fair use (below). A batch that fails any of these is refused with the reason, and nothing of it is kept.
What happens later is per record: each record is checked against every rule, then published — or rejected with its reasons. One bad record never fails the batch.
1. Push
Section titled “1. Push”Send the same POST /v1/publish/products body as for a small push. The answer is 202 with a Location header and the push’s status:
{ "push_id": "bat_7d3a4e8f9b6c1a2b3c4d5e6f", "state": "queued", "records": 8000, "counts": { "rejected": 0, "pending": 8000, "accepted": 0, "updated": 0, "held": 0 }, "progress": { "stage": "check", "chunk": 0, "chunks": 40 }, "status_url": "/v1/pushes/bat_7d3a4e8f9b6c1a2b3c4d5e6f", "poll_after_seconds": 5}The sample below sends 300 records, and polls to complete (section 2).
// A large push: 300 records are checked at the door and accepted at once (202), checked and published in the// background 200 at a time, and read from the push's status until it is complete (/businesses/large-pushes/).//// Uses the key and helpers of push-kit.ts. Where background processing is not set up, up to 500 records are processed// inside the request and answered 200: the rest of the sample reads the same push status either way.import assert from 'node:assert/strict';import { outcomes, product, seal, send, settle } from './push-kit.ts';
const records = Array.from({ length: 300 }, (_, i) => product('L', i + 1));const sent = await send(await seal(records));console.log(`HTTP ${sent.http}: push ${sent.pushId}`);assert.ok(sent.http === 200 || sent.http === 202, `the push was answered ${sent.http}`);
// Poll `GET /v1/pushes/{push_id}` every `poll_after_seconds`, until that member is gone.const done = await settle(sent.pushId);assert.equal(done.state, 'complete');assert.equal(done.records, 300);assert.equal(done.counts.pending, 0);assert.equal(done.counts.rejected, 0);// `accepted` for a record that was not live, `updated` for a new version of one that was (every run after the first).assert.equal(done.counts.accepted + done.counts.updated, 300);
// The per-record outcomes, in leaf order, 100 a page.const all = await outcomes(sent.pushId);assert.equal(all.length, 300);assert.ok(all.every((r, i) => r.leaf_index === i && (r.outcome === 'accepted' || r.outcome === 'updated')));console.log(`complete: ${done.counts.accepted} accepted, ${done.counts.updated} updated, ${done.counts.rejected} rejected`);"""A large push: 300 records are checked at the door and accepted at once (202), checked and published in thebackground 200 at a time, and read from the push's status until it is complete (/businesses/large-pushes/).
Uses the key and helpers of push_kit.py. Where background processing is not set up, up to 500 records are processedinside the request and answered 200: the rest of the sample reads the same push status either way."""
from push_kit import outcomes, product, seal, send, settle
records = [product("L", i + 1) for i in range(300)]push_id, http = send(seal(records))print(f"HTTP {http}: push {push_id}")assert http in (200, 202), f"the push was answered {http}"
# Poll GET /v1/pushes/{push_id} every poll_after_seconds, until that member is gone.done = settle(push_id)counts = done["counts"]assert done["state"] == "complete"assert done["records"] == 300assert counts["pending"] == 0 and counts["rejected"] == 0# `accepted` for a record that was not live, `updated` for a new version of one that was (every run after the first).assert counts["accepted"] + counts["updated"] == 300
# The per-record outcomes, in leaf order, 100 a page.results = outcomes(push_id)assert len(results) == 300assert all(r["leaf_index"] == i and r["outcome"] in ("accepted", "updated") for i, r in enumerate(results))print(f"complete: {counts['accepted']} accepted, {counts['updated']} updated, {counts['rejected']} rejected")2. Poll
Section titled “2. Poll”GET /v1/pushes/{push_id}, signed with tag="mdb-business-read" like every read, answers the push’s totals, its progress and a page of its records’ outcomes in leaf order (100 a page; follow next_cursor). Read it again every poll_after_seconds while that member is present.
statemovesqueued→processing→complete.counts.pendingreaches 0 and the other counts add up torecords.held: the push tripped the price-shock rule; nothing in it is live until an owner or admin confirms it in the Business Portal (Pushes and mandates). Confirming a large push is processed in the background too; the state then movesconfirming→confirmed.stalled: true(with aresumesentence): no step of the push has completed for ten minutes. MasterDB is alerted to a stalled push; you can resume it yourself at once with the next section. The flag is absent while the push is moving and once it has finished.failed: the push stopped on something it cannot get past by itself — your business was put on hold, say — anderrorsays what. What it published stays published.?outcome=rejectedlists only the rejected records, each with every reason, to fix and send again in a new push (below).
Your business’s developers also get an email when a large push stops, is held for confirmation, or has rejected records. A push made through the API never sends an email for each product or each push: your developers get one daily summary, and an email at once only when a push needs a person. Each developer chooses whether they get the summary, only the emails that need a person, or neither (Choosing which emails you get). If your business has no developer, the emails go to your owners.
Fix the records that were rejected
Section titled “Fix the records that were rejected”One bad record never fails the push: it completes with that record rejected. This sample sends 60 records, three of them with a price that is a number instead of a decimal string, lists the rejected ones with ?outcome=rejected — each with its code, the JSON pointer to the field and a detail — and sends the corrected three as a new push with a new seq.
// A large push with records that cannot be published: the push completes, the bad records are `rejected` with every// reason, and you fix just those and send them in a new push (/businesses/large-pushes/).//// Uses the key and helpers of push-kit.ts.import assert from 'node:assert/strict';import { outcomes, product, seal, send, settle } from './push-kit.ts';
// 60 records; three are wrong on purpose: the price is the number 189, not the string "189.00".const good = Array.from({ length: 60 }, (_, i) => product('R', i + 1));const wrong = new Set([6, 22, 40]);const records = good.map((r, i) => (wrong.has(i) ? { ...r, prices: [{ country: 'US', currency: 'USD', amount: 189 }] } : r));
const sent = await send(await seal(records));const done = await settle(sent.pushId);// One bad record never fails the push: it completes, with the other 57 published.assert.equal(done.state, 'complete');assert.equal(done.counts.rejected, 3);assert.equal(done.counts.accepted + done.counts.updated, 57);
// `?outcome=rejected` lists only the rejected records, each with every reason.const rejected = await outcomes(sent.pushId, 'rejected');assert.deepEqual(rejected.map((r) => r.leaf_index), [6, 22, 40]);for (const r of rejected) { console.log(r.leaf_index, r.business_product_id, JSON.stringify(r.errors?.map((e) => `${e.code} at ${e.pointer ?? '/'}`))); assert.equal(r.errors?.[0]?.code, 'money_not_string');}
// Fix those records (here: the correct versions we kept) and send them alone, as a new batch with a new `seq`.const fixed = rejected.map((r) => good[r.leaf_index] as object);const again = await send(await seal(fixed));const second = await settle(again.pushId);assert.equal(second.state, 'complete');assert.equal(second.counts.rejected, 0);assert.equal(second.counts.accepted + second.counts.updated, 3);console.log(`${fixed.length} fixed records published in push ${again.pushId}`);"""A large push with records that cannot be published: the push completes, the bad records are `rejected` with everyreason, and you fix just those and send them in a new push (/businesses/large-pushes/).
Uses the key and helpers of push_kit.py."""
from push_kit import outcomes, product, seal, send, settle
# 60 records; three are wrong on purpose: the price is the number 189, not the string "189.00".good = [product("R", i + 1) for i in range(60)]wrong = {6, 22, 40}records = [{**r, "prices": [{"country": "US", "currency": "USD", "amount": 189}]} if i in wrong else r for i, r in enumerate(good)]
push_id, _ = send(seal(records))done = settle(push_id)# One bad record never fails the push: it completes, with the other 57 published.assert done["state"] == "complete"assert done["counts"]["rejected"] == 3assert done["counts"]["accepted"] + done["counts"]["updated"] == 57
# ?outcome=rejected lists only the rejected records, each with every reason.rejected = outcomes(push_id, "rejected")assert [r["leaf_index"] for r in rejected] == [6, 22, 40]for r in rejected: print(r["leaf_index"], r["business_product_id"], [f"{e['code']} at {e.get('pointer', '/')}" for e in r["errors"]]) assert r["errors"][0]["code"] == "money_not_string"
# Fix those records (here: the correct versions we kept) and send them alone, as a new batch with a new seq.fixed = [good[r["leaf_index"]] for r in rejected]again_id, _ = send(seal(fixed))second = settle(again_id)assert second["state"] == "complete"assert second["counts"]["rejected"] == 0assert second["counts"]["accepted"] + second["counts"]["updated"] == 3print(f"{len(fixed)} fixed records published in push {again_id}")A small push’s answer also names its push_id and status_url: you can read any push’s status the same way.
3. Resume
Section titled “3. Resume”A push that stops (failed), or that has shown no progress for more than ten minutes (its status says stalled: true), is resumed by sending the exact batch again — the same batch_seal and the same records, with a new request signature. It is accepted even after the seal’s five-minute sealed_at window, because MasterDB accepted that batch before. The records already settled keep their outcome; the rest are completed; nothing is published twice. While a push is still moving, sending it again only answers its status. The answer to a large push you send again is 202 with the push’s status and its push_id, the same one as before.
This sample seals one batch, sends it twice as if the first answer was lost, and gets the same push both times; it then polls to the end, resending the same batch if the push failed or shows stalled: true; and a last send, after the push is complete, changes nothing.
// Resuming a push: when an answer is lost, or a push stops or stalls, send the exact batch again. The same seal and// the same records, a new request signature: you are answered the same push, and nothing is published twice// (/businesses/large-pushes/).//// Uses the key and helpers of push-kit.ts.import assert from 'node:assert/strict';import { outcomes, product, seal, send, settle } from './push-kit.ts';
const records = Array.from({ length: 300 }, (_, i) => product('S', i + 1));const body = await seal(records); // sealed once: this exact batch is what you send again
const first = await send(body);// Suppose that answer never reached you. Send the exact batch again: a push still moving is only reported (the// same `push_id`, never a second push), and one already finished is answered with its status.const again = await send(body);assert.equal(again.pushId, first.pushId);console.log(`the same batch again: HTTP ${again.http}, the same push ${again.pushId}`);
// Poll to the end. A push that fails or stalls (no step for ten minutes) is resumed by the same send.const done = await settle(first.pushId, () => send(body));assert.equal(done.state, 'complete');assert.equal(done.counts.pending, 0);assert.equal(done.counts.accepted + done.counts.updated, 300);
// Once it is complete, sending it again changes nothing: the same push, the same totals.const after = await send(body);assert.equal(after.pushId, first.pushId);const unchanged = await settle(after.pushId);assert.equal(unchanged.state, 'complete');assert.deepEqual(unchanged.counts, done.counts);assert.equal((await outcomes(first.pushId)).length, 300);console.log(`complete, and unchanged by sending it again: ${done.counts.accepted} accepted, ${done.counts.updated} updated`);"""Resuming a push: when an answer is lost, or a push stops or stalls, send the exact batch again. The same seal andthe same records, a new request signature: you are answered the same push, and nothing is published twice(/businesses/large-pushes/).
Uses the key and helpers of push_kit.py."""
from push_kit import outcomes, product, seal, send, settle
records = [product("S", i + 1) for i in range(300)]body = seal(records) # sealed once: this exact batch is what you send again
first_id, _ = send(body)# Suppose that answer never reached you. Send the exact batch again: a push still moving is only reported (the# same push_id, never a second push), and one already finished is answered with its status.again_id, http = send(body)assert again_id == first_idprint(f"the same batch again: HTTP {http}, the same push {again_id}")
# Poll to the end. A push that fails or stalls (no step for ten minutes) is resumed by the same send.done = settle(first_id, lambda: send(body))assert done["state"] == "complete"assert done["counts"]["pending"] == 0assert done["counts"]["accepted"] + done["counts"]["updated"] == 300
# Once it is complete, sending it again changes nothing: the same push, the same totals.after_id, _ = send(body)assert after_id == first_idunchanged = settle(after_id)assert unchanged["state"] == "complete"assert unchanged["counts"] == done["counts"]assert len(outcomes(first_id)) == 300print(f"complete, and unchanged by sending it again: {done['counts']['accepted']} accepted, {done['counts']['updated']} updated")Bulk files: up to 100,000 records
Section titled “Bulk files: up to 100,000 records”- Ask for an upload:
POST /v1/bulk-uploadswith{"record_type": "products"}, signed withtag="mdb-push". The answer is a signed upload: aurland itsfields, valid for an hour, for one file of at most 256 MiB. - Upload the file straight to that
urlasmultipart/form-data: every member offields, then the file asfile. The file is one record a line, each line the base64 of the record’s exact bytes (application/x-ndjson). Line 1 is leaf 0 of your batch seal’s Merkle tree. - Import it:
POST /v1/bulk-importswith{"upload_id": "…", "batch_seal": {…}}— the same batch seal as a push’s, over the file’s records (tree_sizethe number of lines). The answer is202with the import’s status; poll it as above.
This sample uploads 100 records and imports them (a bulk file may hold up to 100,000), then polls the import like any push; its status has source: "bulk".
// A bulk file: ask for a signed upload, upload one NDJSON file (one record a line, each the base64 of the record's// exact bytes), then import it under the batch seal over its records and poll it (/businesses/large-pushes/).//// Uses the key and helpers of push-kit.ts. The file here is 100 records; a bulk file may hold up to 100,000.import assert from 'node:assert/strict';import { idempotencyKey } from '@masterdb/client';import { business, outcomes, product, seal, settle } from './push-kit.ts';
const body = await seal(Array.from({ length: 100 }, (_, i) => product('B', i + 1)));
// 1. Ask for an upload: a signed POST (a `url` and its `fields`, valid for an hour) for one file.const slot = await business.POST('/v1/bulk-uploads', { params: { header: { 'Idempotency-Key': idempotencyKey() } }, body: { record_type: 'products' } });if (slot.error) throw new Error(`${slot.error.code}: ${slot.error.detail ?? ''}`);
// 2. Upload the file straight to that `url` as multipart/form-data: every member of `fields`, then the file as `file`.// Line i is leaf i of the seal's Merkle tree: `body.records` are already the base64 lines.const form = new FormData();for (const [name, value] of Object.entries(slot.data.fields)) form.append(name, value);form.append('file', new Blob([`${body.records.join('\n')}\n`], { type: slot.data.content_type }), 'records.ndjson');const uploaded = await fetch(slot.data.url, { method: 'POST', body: form });assert.ok(uploaded.ok, `the upload was answered ${uploaded.status}`);
// 3. Import it: the upload and the seal over its records. Answered 202 with the import's status.const imported = await business.POST('/v1/bulk-imports', { params: { header: { 'Idempotency-Key': idempotencyKey() } }, body: { upload_id: slot.data.upload_id, batch_seal: body.batch_seal },});if (imported.error) throw new Error(`${imported.error.code}: ${imported.error.detail ?? ''}`);assert.equal(imported.response.status, 202);console.log(`import ${imported.data.push_id} of upload ${slot.data.upload_id}`);
// Poll it as any push: the file is read once in the background, then its records are checked and published.const done = await settle(imported.data.push_id);assert.equal(done.state, 'complete');assert.equal(done.source, 'bulk');assert.equal(done.counts.accepted + done.counts.updated, 100);assert.equal((await outcomes(imported.data.push_id)).length, 100);console.log(`complete: ${done.counts.accepted} accepted, ${done.counts.updated} updated`);"""A bulk file: ask for a signed upload, upload one NDJSON file (one record a line, each the base64 of the record'sexact bytes), then import it under the batch seal over its records and poll it (/businesses/large-pushes/).
Uses the key and helpers of push_kit.py. The file here is 100 records; a bulk file may hold up to 100,000."""
import uuid
import httpxfrom push_kit import business, outcomes, product, seal, settle
body = seal([product("B", i + 1) for i in range(100)])
# 1. Ask for an upload: a signed POST (a url and its fields, valid for an hour) for one file.slot = business.post("/v1/bulk-uploads", json={"record_type": "products"}, headers={"Idempotency-Key": str(uuid.uuid4())})if slot.status_code != 201: raise SystemExit(f"{slot.json()['code']}: {slot.json().get('detail', '')}")slot = slot.json()
# 2. Upload the file straight to that url as multipart/form-data: every member of fields, then the file as `file`.# Line i is leaf i of the seal's Merkle tree: body["records"] are already the base64 lines.ndjson = ("\n".join(body["records"]) + "\n").encode("ascii")uploaded = httpx.post(slot["url"], data=slot["fields"], files={"file": ("records.ndjson", ndjson, slot["content_type"])}, timeout=60)assert uploaded.is_success, f"the upload was answered {uploaded.status_code}"
# 3. Import it: the upload and the seal over its records. Answered 202 with the import's status.imported = business.post( "/v1/bulk-imports", json={"upload_id": slot["upload_id"], "batch_seal": body["batch_seal"]}, headers={"Idempotency-Key": str(uuid.uuid4())},)if imported.status_code != 202: raise SystemExit(f"{imported.json()['code']}: {imported.json().get('detail', '')}")push_id = imported.json()["push_id"]print(f"import {push_id} of upload {slot['upload_id']}")
# Poll it as any push: the file is read once in the background, then its records are checked and published.done = settle(push_id)assert done["state"] == "complete"assert done["source"] == "bulk"assert done["counts"]["accepted"] + done["counts"]["updated"] == 100assert len(outcomes(push_id)) == 100print(f"complete: {done['counts']['accepted']} accepted, {done['counts']['updated']} updated")The file is read once in the background. If its records are not exactly the ones your seal names, the import is failed with seal_invalid before any record is checked. Uploaded files are deleted as soon as their import completes, and after 7 days in any case. Sending the same import request again answers its status, and resumes one that stopped.
Limits
Section titled “Limits”| Limit | Answer when over | |
|---|---|---|
| Records in one push | 10,000 (50 or fewer: answered inline) | 400 request_invalid |
| Request body | 32 MiB | 413 |
| One record | 153,600 bytes (204,800 base64 characters) | that record rejected |
| Records in one bulk file | 100,000 | 400 limit_exceeded |
| One bulk file | 256 MiB | 400 limit_exceeded |
| Your key’s records an hour and bytes a day | as your mandate states | 429 rate_limited, cap: records_per_hour or bytes_per_day |
| Large pushes and bulk imports in progress at once, per business | 3 | 429 rate_limited, cap: concurrent_pushes, Retry-After: 60 |
| Records an hour, per business, across all its keys | 250,000 | 429 rate_limited, cap: business_records_per_hour, retry at the next hour |
A 429 carries Retry-After and the problem’s retry_after_seconds, cap and limit.
Retrying a 503
Section titled “Retrying a 503”A 503 with Retry-After on a push or a bulk route is always safe to retry with the same batch: it is the same request, and a push that was already accepted is never processed twice.