Node collector

Set up @seamward/collector in a Node.js backend and record observations from one integration. The collector is observe-only: it batches redacted evidence asynchronously and never sits on your critical request path.

If you are not sure where to add the collector, start with assisted codebase setup. The setup agent discovers likely HTTP, webhook, queue, and scheduled-feed boundaries locally, creates a reviewable plan, previews every source change, and applies only the changes you approve.

Backend only

Create the collector in a trusted backend process. Never expose the secret ingest token in browser code, Git, logs, or support messages.

Before you begin

  • Node.js 22 or later; the package is ESM-only.
  • A public Connection key and secret ingest token from the integration's Collector setup tab.

Install

Code example

npm install @seamward/collector

The collector is a prerelease (0.1.0-alpha); pin the version through your lockfile as with any dependency.

Configure

Set the connection values in your backend environment. This example observes five Integration boundaries in one process:

Code example

SEAMWARD_WEBHOOK_CONNECTION_KEY=sw_conn_v1.replace_webhookSEAMWARD_HTTP_API_CONNECTION_KEY=sw_conn_v1.replace_http_apiSEAMWARD_QUEUE_PUBLISH_CONNECTION_KEY=sw_conn_v1.replace_queue_publishSEAMWARD_QUEUE_CONSUME_CONNECTION_KEY=sw_conn_v1.replace_queue_consumeSEAMWARD_FEED_EXPORT_CONNECTION_KEY=sw_conn_v1.replace_feed_exportSEAMWARD_WEBHOOK_INGEST_TOKEN=sw_ing_replace_webhookSEAMWARD_HTTP_API_INGEST_TOKEN=sw_ing_replace_http_apiSEAMWARD_QUEUE_PUBLISH_INGEST_TOKEN=sw_ing_replace_queue_publishSEAMWARD_QUEUE_CONSUME_INGEST_TOKEN=sw_ing_replace_queue_consumeSEAMWARD_FEED_EXPORT_INGEST_TOKEN=sw_ing_replace_feed_export
ValueWhere to get it
Each *_CONNECTION_KEYThe matching integration's Collector setup tab; public and safe to keep in ordinary runtime configuration
Each *_INGEST_TOKENThe matching integration's Collector setup tab; shown once when issued or rotated

Create one collector per Integration boundary and reuse that instance across operations belonging to the boundary:

Code example

import { createSeamwardCollector } from "@seamward/collector";
const connected = (connectionKey: string, ingestToken: string) =>  createSeamwardCollector({    connectionKey,    ingestToken,    deployment: {      service: "candidate-api",      release: process.env.RELEASE_VERSION,      commitSha: process.env.GITHUB_SHA,    },  });
export const webhookCollector = connected(  process.env.SEAMWARD_WEBHOOK_CONNECTION_KEY!,  process.env.SEAMWARD_WEBHOOK_INGEST_TOKEN!,);export const httpApiCollector = connected(  process.env.SEAMWARD_HTTP_API_CONNECTION_KEY!,  process.env.SEAMWARD_HTTP_API_INGEST_TOKEN!,);export const queuePublisherCollector = connected(  process.env.SEAMWARD_QUEUE_PUBLISH_CONNECTION_KEY!,  process.env.SEAMWARD_QUEUE_PUBLISH_INGEST_TOKEN!,);export const queueConsumerCollector = connected(  process.env.SEAMWARD_QUEUE_CONSUME_CONNECTION_KEY!,  process.env.SEAMWARD_QUEUE_CONSUME_INGEST_TOKEN!,);export const feedExportCollector = connected(  process.env.SEAMWARD_FEED_EXPORT_CONNECTION_KEY!,  process.env.SEAMWARD_FEED_EXPORT_INGEST_TOKEN!,);

One Integration boundary per collector

An Integration currently has one direction and one protocol. Use a separate Connection key and collector instance for inbound and outbound traffic, and for HTTP, webhook, queue, and scheduled-feed boundaries. Instances in the same workspace environment use their own Integration-owned ingest tokens.

deployment is optional; the collector also detects common platform variables (GITHUB_SHA, VERCEL_GIT_COMMIT_SHA, and others listed in configuration) automatically, and ignores invalid values instead of breaking your application.

Instrument a webhook

Wrap the existing business handler. Your handler's result or thrown error stays authoritative; observation happens after the operation.

Code example

const handleCandidateWebhook = webhookCollector.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 },    };  },);
fastify.post("/webhooks/candidates", async (request, reply) => {  const result = await handleCandidateWebhook(    request.body as Record<string, unknown>,  );  return reply.code(result.statusCode ?? 200).send();});

The identifier fields (sourceEventId, businessObjectId) are converted to keyed hashes inside your process before anything is queued; Seamward never receives the original values. For outbound calls, observeFetch wraps the standard fetch signature the same way.

Observe an outbound HTTP API

Wrap the provider call with observeFetch. Response-side evidence is the default, so the collector clones the response and inspects JSON asynchronously while returning the original Response as soon as the provider resolves.

Code example

const providerFetch = httpApiCollector.observeFetch({  routeTemplate: "/v1/candidates",  eventType: "candidate.create",});
const response = await providerFetch("https://provider.example/v1/candidates", {  method: "POST",  headers: { "content-type": "application/json" },  body: JSON.stringify({ externalId: "candidate_123" }),});

Use payloadLocation: "request" when the request contract is the useful side of the operation. JSON string bodies and cloned Request objects are inspected automatically. For form data, streams, protobuf, or another non-JSON body, provide a local JSON value through requestPayload:

Code example

const providerFetch = httpApiCollector.observeFetch({  routeTemplate: "/v1/documents",  eventType: "document.upload",  payloadLocation: "request",  requestPayload: {    documentId: "shape-only-local-value",    contentType: "application/pdf",  },});

requestPayload is used inside your process to calculate the structural shape. Its values are not added to the envelope. If no JSON representation is available, the operation still records transport evidence with a null payload shape. A network failure that receives no HTTP response records statusCode: 0, accepted: false, then rethrows the original fetch error.

Observe a queue

Use the publish wrapper for messages your service sends and the consumer wrapper for messages it receives. Both preserve the return value and rethrow host errors after recording them.

Code example

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

queueName must be a stable queue, topic, or subscription label, never a message id. Return disposition: "rejected" to record an explicit rejection. Set attempt from broker delivery metadata when it is available.

The wrappers are broker-neutral. Seamward does not currently ship adapters for Kafka, RabbitMQ, SQS, Pub/Sub, or another broker. Keep acknowledgement, visibility timeout, retry, and dead-letter behavior in your existing client.

Observe a scheduled feed

Wrap one import or export run and state its direction. An inbound feed brings provider data into your application; an outbound feed sends application data to the provider.

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 };  },);

Use a stable feedName, never a run id or timestamp. Return disposition: "rejected" for a completed run whose business result was rejected. Thrown job errors are recorded and rethrown unchanged, so your scheduler remains responsible for retry and failure handling.

Queue and scheduled-feed wrappers expose protocol-native options, but envelope 0.2 still stores them through a compatibility projection: the queue or feed label occupies transport.routeTemplate, method is POST, and payloadLocation is message. This keeps current APIs and workers compatible. Treat these transport fields as storage details, not broker commands.

Verify

  1. Start the backend with the Connection key and ingest token configured.
  2. Trigger the handler once with a safe test payload.
  3. Wait for the five-second automatic flush, or call await webhookCollector.flush().
  4. Open the integration in Seamward and confirm the setup status reports the first observation instead of No observation received yet.

If nothing arrives, seamward.stats() tells you where it stopped; the troubleshooting guide maps each counter to its fix.

Shutdown

Give queued observations a final flush during controlled shutdown:

Code example

process.once("SIGTERM", async () => {  await Promise.all([    webhookCollector.stop(),    httpApiCollector.stop(),    queuePublisherCollector.stop(),    queueConsumerCollector.stop(),    feedExportCollector.stop(),  ]);  process.exit(0);});

flush() and stop() never reject, and the flush timer never keeps your process alive on its own.

Bounded inspection and shutdown

The collector caps pending inspections at 64, reads at most 256,000 bytes per JSON body and stops a stalled inspection after one second by default. Check skippedInspections and pendingInspections alongside delivery counters. stop() finishes its current inspection and shipping cycle within roughly 11 seconds with defaults, plus event-loop scheduling time. Allow that grace period during shutdown. See the resource limits for configuration and the API-before-collector upgrade order.

Next steps