Node.js collector SDK

@seamward/collector observes third-party HTTP APIs, webhooks, queue operations, and scheduled feeds from a Node.js backend. It batches redacted evidence asynchronously and never sits on your critical request path. Collector failures do not throw into your application, while errors from your own handlers and clients are recorded and rethrown unchanged.

Requires Node.js 22 or later; the package is ESM-only. See Node collector setup for installation and a guided first integration.

createSeamwardCollector

Code example

import { createSeamwardCollector } from "@seamward/collector";
const seamward = createSeamwardCollector({  connectionKey: process.env.SEAMWARD_CONNECTION_KEY!,  ingestToken: process.env.SEAMWARD_INGEST_TOKEN!,});
OptionTypeDefaultPurpose
connectionKeystringrequiredPublic key from this integration's Collector setup tab (sw_conn_v1...)
ingestTokenstringrequiredThe backend-only secret token (sw_ing_...); also keys request signing
endpointstringSEAMWARD_INGEST_URL env var, else https://api.seamward.com/ingestIngestion URL for custom or self-hosted deployments
policyobjectseamward-default-v1Redaction policy: { version, dropFields, hashFields }; see Redaction
deploymentobjectauto-detected{ service?, release?, commitSha? }; explicit values win over environment detection
maxBatchSizenumber50Envelopes per delivery request
flushIntervalMsnumber5000Automatic flush interval
maxQueueSizenumber1000In-memory queue bound; overflow drops the oldest observations
fetchFnfetchglobal fetchInjectable transport, mainly for tests

Construction fails fast when the Connection key or ingest token is invalid. Key values are never echoed in errors.

Deployment context is resolved once at construction. Without explicit values, the collector reads SEAMWARD_SERVICE, SEAMWARD_RELEASE, and the first commit SHA found among SEAMWARD_COMMIT_SHA, VERCEL_GIT_COMMIT_SHA, RENDER_GIT_COMMIT, RAILWAY_GIT_COMMIT_SHA, and GITHUB_SHA. When no release is set, it falls back to the first 12 characters of the commit SHA. Invalid values are ignored rather than breaking your application.

Observe operations

The Connection key binds this collector to one Integration, so observe entry points need only operation metadata. An Integration has one direction, one protocol, and its own ingest token. A process that observes multiple boundaries creates one collector for each Integration using that Integration's Connection key and ingest token.

observeWebhook

Wraps an inbound handler. Your handler's return value or thrown error stays authoritative; observation happens after the fact.

Code example

const handleCandidateWebhook = seamward.observeWebhook(  {    routeTemplate: "/webhooks/candidates",  },  async (payload) => {    const candidate = await candidateService.accept(payload);    return {      statusCode: 202,      eventType: "candidate.create",      outcome: {        accepted: true,        businessObjectType: "candidate",        businessObjectId: candidate.id,      },      correlation: { sourceEventId: candidate.providerEventId },    };  },);

Meta fields include routeTemplate (a stable template, never a concrete id), plus optional method (default POST), attempt, and correlation. Webhook observations use payloadLocation: "message". Set attempt to the provider delivery attempt and pass an idempotency key in correlation when one exists; both feed retry analysis without exposing the raw identifier. The handler may return a WebhookResult with optional statusCode, eventType, outcome, and correlation. Identifier fields (sourceEventId, idempotencyKey, businessObjectId) are converted to keyed hashes inside your process; Seamward never receives the original values. A thrown handler error is recorded as statusCode: 500, accepted: false, and then rethrown unchanged.

observeFetch

Wraps an outbound call with the standard fetch signature:

Code example

const providerFetch = seamward.observeFetch({  routeTemplate: "/v1/candidates",  eventType: "candidate.create",});
const response = await providerFetch(  "https://provider.example/v1/candidates",  init,);

The original response is returned as soon as the provider resolves; the collector clones it and inspects the JSON asynchronously, so response timing is preserved and non-JSON bodies still produce transport evidence. durationMs measures the provider call, not the inspection. Request headers are never read. A rejected fetch is recorded as statusCode: 0, accepted: false, and rethrown.

Fetch metadata also accepts optional method, attempt, correlation, payloadLocation, and requestPayload. Outbound fetch observations default to payloadLocation: "response".

Choose payloadLocation: "request" to describe the request contract instead. The collector inspects a JSON string in init.body or clones and parses a Request without consuming the original. For non-JSON bodies, provide a JSON representation through requestPayload; the collector derives its structural shape locally and never adds the values to the envelope. Without a parseable body or requestPayload, the operation still records transport evidence with a null payload shape.

observeQueuePublish

Wraps a broker-neutral outbound publish operation:

Code example

const publishCandidate = queuePublisherCollector.observeQueuePublish(  {    queueName: "candidate-events",    eventType: "candidate.requested",    attempt: 1,  },  async (message) => {    await queueClient.publish("candidate-events", message);    return { disposition: "acknowledged" as const };  },);

QueueMeta requires queueName, a stable queue, topic, or subscription label. It also accepts optional eventType, attempt, and correlation. A normal return defaults to disposition: "acknowledged" and accepted: true; an explicit disposition: "rejected" defaults to accepted: false. QueueResult.outcome can override that business result, and eventType or correlation returned by the publisher takes precedence over the metadata.

observeQueueConsumer

Wraps a broker-neutral inbound consumer with the same metadata and result types:

Code example

const consumeCandidate = queueConsumerCollector.observeQueueConsumer(  {    queueName: "candidate-events",    eventType: "candidate.received",  },  async (message) => {    await candidateService.accept(message);    return {      disposition: "acknowledged" as const,      outcome: { accepted: true },    };  },);

The publish and consumer wrappers do not acknowledge, retry, or dead-letter a broker message for you. Seamward does not currently ship broker-specific adapters. Your existing Kafka, RabbitMQ, SQS, Pub/Sub, or other client remains authoritative.

observeScheduledFeed

Wraps one scheduled import or export run:

Code example

const exportCandidates = feedExportCollector.observeScheduledFeed(  {    feedName: "nightly-candidate-export",    direction: "outbound",    eventType: "candidate.exported",  },  async (batch) => {    await feedClient.upload(batch);    return { disposition: "completed" as const };  },);

ScheduledFeedMeta requires a stable feedName and direction (inbound for imports, outbound for exports). It also accepts optional eventType, attempt, and correlation. A normal return defaults to disposition: "completed" and accepted: true; disposition: "rejected" defaults to accepted: false. ScheduledFeedResult can supply an explicit outcome, eventType, or correlation.

Queue and feed compatibility projection

The queue and scheduled-feed entry points are protocol-native at the SDK boundary. Envelope 0.2 currently projects them into the shared transport object so existing APIs and workers remain compatible:

WrapperDirectionmethodrouteTemplatepayloadLocationSuccessExplicit rejectionThrown host error
observeQueuePublishoutboundPOSTqueueNamemessage202422500
observeQueueConsumerinboundPOSTqueueNamemessage202422500
observeScheduledFeedcaller-selectedPOSTfeedNamemessage200422500

These status codes describe the collector's compatibility projection, not a broker or scheduler response. Use outcome for the business result and attempt for delivery or run retries.

record

Use record(input) for an operation that cannot be wrapped. It accepts direction, protocol, method, routeTemplate, statusCode, and durationMs, plus optional payloadLocation, attempt, eventType, payload, correlation, and outcome. The default payload location is message for webhooks, queues, and scheduled feeds; outbound HTTP API calls default to response, and inbound HTTP API calls default to request.

statusCode must be 0 or an integer from 100 through 599. Zero means that no protocol response was received. When outcome is omitted, accepted defaults to false for status 0 and to statusCode < 400 otherwise.

Lifecycle and stats

Code example

process.once("SIGTERM", async () => {  await seamward.stop();  process.exit(0);});
  • flush() waits for the current bounded inspections, then starts or joins one bounded shipping cycle. Never rejects.
  • stop() clears the flush timer and performs one final flush. Never rejects.
  • stats() returns these counters: enqueued, shipped, rejected, dropped, failedBatches, queueLength (live), buildErrors, skippedInspections, and pendingInspections (live).

The flush timer is unreferenced, so the collector never keeps your process alive on its own.

Failure behavior

The design rule is fail-open: telemetry must never break the host application.

SituationBehaviorSignal
Endpoint unreachable, 409, 429, or 5xxBatch stays queued; retried on the next flushfailedBatches increments
Non-retryable 4xx or a rejected 207 itemRejected envelopes leave the queuerejected increments
Queue exceeds maxQueueSizeOldest observations are dropped; newest keptdropped increments
An observation fails envelope validationSilently discardedbuildErrors increments
Handler, fetch, publish, consume, or feed throwsRecorded, then rethrown unchangedYour normal error handling

Raw payload value size does not determine envelope size because values never ship. Deep or unusually wide structures can still produce large envelopes, so normal ingestion request limits apply. A persistently failing endpoint retries on every flush until the queue overflows into drop-oldest, so watch failedBatches in your health checks.

Advanced exports

Most applications never need these; they exist for self-hosted setups and custom pipelines.

  • createCollector(config): the unauthenticated core with the same surface, for callers managing raw scope ids and a complete policy themselves.
  • createShipper(config): batching, signing, and retry transport alone, for shipping pre-built envelopes.
  • buildEnvelope(input, policy, deps?): the pure envelope builder; deterministic with injected clock and id.
  • redactPayload(value, policy) / redactHeaders(headers, policy) / keyedHash(value, key): standalone redaction helpers for your own logging pipelines.
  • resolveDeploymentContext(explicit?, env?): the deployment detection used at construction.

Upload queue accounting

maxQueueSize bounds pending envelopes (default 1,000). One additional batch of up to maxBatchSize envelopes (default 50) can be in flight. Overflow drops the oldest pending envelopes and increments dropped; it cannot alter an active batch. A failed batch is returned to the front of the pending queue, where the same overflow rule applies. queueLength includes pending and in-flight envelopes. Acknowledgements update only the batch actually sent.

Node collector resource limits

These additive options apply to the Node collector:

OptionDefaultMaximumPurpose
maxPendingInspections641,024Pending body inspections and deferred records
maxPayloadBytes256,0001,000,000Body byte cap and conservative JSON copy budget
inspectionTimeoutMs1,00010,000Total time to read one cloned body
uploadTimeoutMs5,00030,000Upload and acknowledgement deadline
flushTimeoutMs10,00060,000Deadline for a shipping cycle

Values must be positive integers within the maximum; invalid values use the default. Inspection also limits nesting to 32 levels and visits to 8,192 JSON nodes. Only JSON content types (application/json and +json) are read from HTTP bodies. Oversized, unsupported and timed-out bodies produce transport-only evidence where capacity allows. A full inspection queue skips the observation. skippedInspections counts these skips and pendingInspections reports active work. No raw payload is transmitted.

Envelope 0.3 carries payload.schemaInspected: false for a skipped body. Its null placeholder is excluded from structural drift detection; status and timing still contribute to behavioral analysis. Older observations without the marker retain their existing interpretation. Deploy the updated API and worker, then update external strict readers, before upgrading Node collectors to envelope 0.3. Older parsers reject the new wire version. Existing 0.1 and 0.2 data remains readable. See the wire contract for the pending release boundary.

Concurrent flush callers share a shipping cycle. Each cycle processes at most the queue size it started with; new observations wait for the next cycle. stop() prevents new inspection work, waits for the bounded current inspections, and attempts a final shipping cycle. With defaults, it finishes within roughly 11 seconds plus event-loop scheduling time. Timed-out uploads remain queued for a future flush; after process exit, unshipped in-memory observations are lost.

Next steps