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/collectorThe 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| Value | Where to get it |
|---|---|
Each *_CONNECTION_KEY | The matching integration's Collector setup tab; public and safe to keep in ordinary runtime configuration |
Each *_INGEST_TOKEN | The 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
- Start the backend with the Connection key and ingest token configured.
- Trigger the handler once with a safe test payload.
- Wait for the five-second automatic flush, or call
await webhookCollector.flush(). - 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
- Observe webhooks well: boundaries, route templates, and the outcome contract.
- Collector SDK reference: exact metadata, result types, and status projections for every wrapper.
- Observation envelope: how the four protocols appear on the wire.
- Redaction: exactly what leaves your process.
