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!,});| Option | Type | Default | Purpose |
|---|---|---|---|
connectionKey | string | required | Public key from this integration's Collector setup tab (sw_conn_v1...) |
ingestToken | string | required | The backend-only secret token (sw_ing_...); also keys request signing |
endpoint | string | SEAMWARD_INGEST_URL env var, else https://api.seamward.com/ingest | Ingestion URL for custom or self-hosted deployments |
policy | object | seamward-default-v1 | Redaction policy: { version, dropFields, hashFields }; see Redaction |
deployment | object | auto-detected | { service?, release?, commitSha? }; explicit values win over environment detection |
maxBatchSize | number | 50 | Envelopes per delivery request |
flushIntervalMs | number | 5000 | Automatic flush interval |
maxQueueSize | number | 1000 | In-memory queue bound; overflow drops the oldest observations |
fetchFn | fetch | global fetch | Injectable 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:
| Wrapper | Direction | method | routeTemplate | payloadLocation | Success | Explicit rejection | Thrown host error |
|---|---|---|---|---|---|---|---|
observeQueuePublish | outbound | POST | queueName | message | 202 | 422 | 500 |
observeQueueConsumer | inbound | POST | queueName | message | 202 | 422 | 500 |
observeScheduledFeed | caller-selected | POST | feedName | message | 200 | 422 | 500 |
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, andpendingInspections(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.
| Situation | Behavior | Signal |
|---|---|---|
Endpoint unreachable, 409, 429, or 5xx | Batch stays queued; retried on the next flush | failedBatches increments |
Non-retryable 4xx or a rejected 207 item | Rejected envelopes leave the queue | rejected increments |
Queue exceeds maxQueueSize | Oldest observations are dropped; newest kept | dropped increments |
| An observation fails envelope validation | Silently discarded | buildErrors increments |
| Handler, fetch, publish, consume, or feed throws | Recorded, then rethrown unchanged | Your 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:
| Option | Default | Maximum | Purpose |
|---|---|---|---|
maxPendingInspections | 64 | 1,024 | Pending body inspections and deferred records |
maxPayloadBytes | 256,000 | 1,000,000 | Body byte cap and conservative JSON copy budget |
inspectionTimeoutMs | 1,000 | 10,000 | Total time to read one cloned body |
uploadTimeoutMs | 5,000 | 30,000 | Upload and acknowledgement deadline |
flushTimeoutMs | 10,000 | 60,000 | Deadline 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
- Node collector setup: install and record a first observation.
- Configuration: environment variables and option reference.
- Ingest observations: the wire contract the collector implements.
