Repository navigation
Split batcher and ipfs publisher - #31
Conversation
There was a problem hiding this comment.
🟡 Changes recommended
Unresolved critical and moderate reliability issues remain in the batcher and IPFS publisher.
Get a fresh assessment by requesting another Copilot review.
Pull request overview
This PR separates telemetry batching from IPFS publication by adding a standalone @scp/batcher service and new Kafka/protobuf contracts.
Changes:
- Adds batcher flushing, offset handling, configuration, and tests.
- Adds
TelemetryBatchedPayload,telemetry.batched.v1, andipfs.published.v1. - Updates the IPFS publisher, documentation, Kafka initialization, and environment configuration.
File summaries
| File | Summary | Final review notes |
|---|---|---|
services/ipfs-publisher/tests/contracts.test.ts |
Updates publisher contract tests. | — |
services/ipfs-publisher/tests/config.test.ts |
Updates configuration tests. | — |
services/ipfs-publisher/src/index.ts |
Consumes batches and publishes to IPFS. | Critical (3): Failed batches can be skipped by later commits. Moderate (1): Upload success followed by result-send failure can republish the batch. |
services/ipfs-publisher/src/config.ts |
Removes obsolete batching configuration. | — |
services/ipfs-publisher/README.md |
Documents publisher responsibilities. | — |
services/batcher/tsdown.config.ts |
Adds batcher build configuration. | — |
services/batcher/tsconfig.json |
Adds TypeScript configuration. | — |
services/batcher/tests/service.shutdown.test.ts |
Tests shutdown flushing. | — |
services/batcher/tests/contracts.test.ts |
Tests batcher contracts. | — |
services/batcher/tests/batch-flusher.test.ts |
Tests serialized flush behavior. | — |
services/batcher/src/logger.ts |
Adds batcher logging. | — |
services/batcher/src/index.ts |
Implements batching and Kafka offset management. | Critical (3): Retry attempts generate new batch IDs. Moderate (3): Failed flushes may not re-arm the timer. Moderate (1): Joined flushes can leave threshold-sized batches pending. Moderate (1): Unbounded retries can stall partitions. Critical (1): Parse failures can be skipped by later commits. |
services/batcher/src/config.ts |
Adds batcher configuration loading. | — |
services/batcher/src/batch-flusher.ts |
Implements serialized flushing. | — |
services/batcher/README.md |
Documents batcher operation. | — |
services/batcher/package.json |
Defines the batcher package. | — |
README.md |
Updates the service overview. | Nit (2): Describes blockchain-anchor deduplication and result emission that are not currently implemented. |
proto/connectivity/v1/payload.proto |
Adds the batched telemetry payload schema. | — |
pnpm-lock.yaml |
Locks batcher dependencies. | — |
packages/core/tests/core.test.ts |
Tests the new topic constant. | — |
packages/core/src/topics.ts |
Adds the batched topic constant. | — |
docs/wp/wp-04-ipfs-publisher.md |
Notes the service split. | Nit (1): Remaining work-package guidance contradicts the new batcher/publisher interfaces. |
docs/architecture/project-architecture.md |
Updates pipeline architecture. | Nit (3): Blockchain Anchor still references the old input topic. |
docker/kafka-init-topics.sh |
Initializes the batched Kafka topic. | — |
.github/copilot-instructions.md |
Updates pipeline guidance. | — |
.env.example |
Adds batcher environment settings. | — |
Review details
Files not reviewed (1)
- pnpm-lock.yaml: Generated file
Suppressed comments (4)
docs/wp/wp-04-ipfs-publisher.md:6
- This update acknowledges the split but leaves the rest of the work-package contract unchanged: the following sections still require the publisher to consume
telemetry.authorized.v1, own deterministic batching, and emittelemetry.ipfs.result.v1. Since this document is implementation guidance, rewrite or clearly archive those sections so they no longer contradict the new batcher/publisher interfaces.
> **Update (batcher/publisher split):** batching and IPFS publication are now separate services. The `@scp/batcher` service consumes `telemetry.authorized.v1`, builds deterministic batches, and emits `telemetry.batched.v1`. The `@scp/ipfs-publisher` service consumes `telemetry.batched.v1`, publishes/pins the batch artifact to IPFS, and emits `ipfs.published.v1` (`cid`, `event_count`). The sections below describe the combined responsibility; the batching tasks now belong to the batcher and the publication tasks to the ipfs-publisher.
services/batcher/src/index.ts:383
- If a timer/lag flush is already in flight, this call joins that promise rather than starting another one. Messages arriving during that flush accumulate in a new batch; when this joined call returns, the code does not re-check the new batch, so a batch that reached
batchSizecan remain unflushed indefinitely because the timer was cleared by the first flush. After joining, schedule or execute another flush when pending items meet the threshold (and apply the same handling to the lag path).
await flushBatch();
services/batcher/src/index.ts:284
- All produce and commit failures are re-attached and retried without an attempt limit or a DLQ path. A permanent Kafka failure can therefore keep the earliest offsets uncommitted and stall that partition indefinitely while new items accumulate. Use the pipeline's bounded retry policy and route exhausted batches to
telemetry.dlq.v1before advancing the source offset.
} catch (error) {
const failedSensorIds = Array.from(
new Set(batchToPublish.map((b) => formatSensorId(b.sensorId)))
);
logWarn('batch produce failed, will retry on next flush', {
services/ipfs-publisher/src/index.ts:297
- The deduplication entry is added only after
publishBatchedcompletes both the IPFS upload and the result-event send. If the upload succeeds butproducer.sendfails, the offset stays uncommitted while this batch ID is still absent fromdedup, so the redelivery uploads/pins the same batch again. Persist completion or separate/idempotently retry the result emission so this failure window does not republish to IPFS.
dedup.add(batched.batchId);
- Files reviewed: 25/26 changed files
- Comments generated: 6
- Review effort level: Lite
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Updated descriptions for blockchain-anchor to clarify its function with the Robonomics blockchain and modified event consumption details in project architecture documentation. Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
Modify flusher error handling to reset batch timer if conditions are met. Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
- Derive batcher batch_id deterministically from partition/offset fingerprints instead of a fresh random UUID per produce attempt, so a retried flush (after a failed offset commit) reuses the same ID and downstream batch_id dedup in ipfs-publisher works correctly. - Route unparseable batcher records to the DLQ (telemetry.dlq.v1) as explicit "poison" batch items instead of silently skipping them; their offsets are only considered handled once the DLQ publish (and batch produce) succeeds, preventing a later valid record's offset commit from silently skipping past them. - In ipfs-publisher, retry failed batch processing in place with bounded attempts before giving up, then explicitly route the poisoned message to the DLQ and commit its offset, or stop processing the partition if DLQ routing itself fails - instead of continuing to the next message and letting its offset commit silently skip past the failed one. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
There was a problem hiding this comment.
🟡 Changes recommended
Unresolved moderate issues remain in batcher and publisher error handling, DLQ semantics, reliability, and integration startup.
Get a fresh assessment by requesting another Copilot review.
Review details
Files not reviewed (1)
- pnpm-lock.yaml: Generated file
Suppressed comments (15)
Previously missed (2) — in code that hasn't changed since the last review.
services/ipfs-publisher/src/index.ts:297
- Skipping a null-valued Kafka record without committing its offset can let a later valid batch commit past it, silently dropping the record; if it is the only record, it remains uncommitted. Route tombstones to the DLQ or explicitly commit the skipped offset before continuing.
This issue also appears on line 315 of the same file.
docs/wp/wp-04-ipfs-publisher.md:6
- The update acknowledges that the sections below are now stale, but they still specify
telemetry.authorized.v1as input andtelemetry.ipfs.result.v1as output; WP-05 likewise retains the old input topic. This leaves the work-package source documentation describing a different pipeline and can lead implementers to wire the wrong topics. Update or replace the remaining checklists/contracts for the split services.
README.md:98
- The integration workflow still starts
ipfs-publisherbut does not start the newbatcher; consequently CI sends authorized events while the publisher listens only totelemetry.batched.v1, so this new pipeline is never exercised and the publisher receives no input. Add the batcher to service startup and its health/log checks.
- **batcher**: Consumes `telemetry.authorized.v1`, groups events into deterministic batches (by size, lag, and a bounded flush timer), and emits `telemetry.batched.v1`. Serializes flushes (single active flush per instance) and flushes pending batches on graceful shutdown.
- **ipfs-publisher**: Consumes `telemetry.batched.v1`, publishes/pins each batch to IPFS (optionally XZ-compressed), deduplicates by `batch_id`, emits `ipfs.published.v1`.
services/batcher/src/index.ts:495
await maybeFlush()is inside the sametryas deserialization. A produce or offset-commit failure from a full batch therefore enters this handler, is logged as an envelope parse error, and adds the already-valid current record a second time as a poison item; the next flush can DLQ it and inflate the batch. Keep parsing/validation separate from the flush call, or distinguish these error sources before adding a poison item.
} catch (error) {
// The record could not be parsed. Rather than silently
// dropping it (which would let a later, valid record's
// offset commit skip past it undetected), keep its offset in
// the same batch as a poison item so it is explicitly routed
services/batcher/src/index.ts:367
- Every produce or commit failure is re-attached and retried by the timer indefinitely; this service has no maximum attempt count or DLQ path for a permanently failing batch. A persistent Kafka failure can therefore block the consumer partition forever, unlike the bounded retry/DLQ behavior expected for pipeline consumers. Add bounded attempts and route exhausted records to
telemetry.dlq.v1before committing/skipping them.
} catch (error) {
const failedSensorIds = Array.from(
new Set(validItems.map((b) => formatSensorId(b.sensorId)))
);
logWarn('batch produce failed, will retry on next flush', {
batch_size: batchToPublish.length,
poison_count: poisonItems.length,
unique_sensors: failedSensorIds.length,
sensor_ids: failedSensorIds,
error: error instanceof Error ? error.message : String(error),
});
// Re-throw so the flusher re-attaches the batch for retry.
throw error;
services/batcher/src/index.ts:443
- Skipping a null-valued Kafka record without tracking or committing its offset allows a later valid record's max-offset commit to advance past this record, silently dropping it; if no later record arrives, it remains uncommitted for redelivery. Treat tombstones as a DLQ/poison case or commit them explicitly before continuing.
if (!message.value) {
logWarn('received null message value; skipping');
continue;
}
services/batcher/src/index.ts:455
- An envelope with a non-authorized
eventTypeis ignored without entering the poison path or committing its offset. A later valid batch can then commit past it, silently losing it (or it can block forever at the head of the partition). Treat the mismatch as a poison record/DLQ case or explicitly commit the ignored record.
if (envelope.eventType !== TELEMETRY_TOPICS.AUTHORIZED) {
logDebug('non-authorized envelope ignored', {
eventType: envelope.eventType,
});
continue;
services/batcher/src/index.ts:383
- When a timer-triggered flush joins an already-active flush,
createBatchFlusher.flush()waits for that flush and then returns without flushing the items that arrived while it was in progress. Since this code clears the timer before joining, those items can remain buffered indefinitely when no further message arrives. Re-arm the timer or trigger another flush after the joined flush completes wheneverflusher.size() > 0.
await flusher.flush().catch(() => {
if (!shouldStop && flusher.size() > 0) {
resetBatchTimer();
}
});
services/batcher/src/index.ts:125
- The DLQ record puts the original authorized envelope bytes directly on
telemetry.dlq.v1, so the value still identifies itself as anAUTHORIZEDevent rather than a canonical DLQ envelope. Existing DLQ publishing wraps a payload inEnvelopeSchemawitheventType: TELEMETRY_TOPICS.DLQ(services/endpoint/src/producer.ts:112-128). Use that envelope convention while preserving the raw record and failure context for DLQ consumers.
await producer.send({
messages: poisonItems.map((item) => ({
topic: TELEMETRY_TOPICS.DLQ,
value: Buffer.from(item.raw ?? new Uint8Array()),
headers: {
source_topic: Buffer.from(TELEMETRY_TOPICS.AUTHORIZED),
source_service: Buffer.from(config.source),
source_partition: Buffer.from(String(item.partition)),
source_offset: Buffer.from(String(item.offset)),
reason: Buffer.from(item.parseError ?? 'unknown parse error'),
services/ipfs-publisher/README.md:11
- The failure description no longer matches the implementation: after three attempts,
src/index.tspublishes the raw message totelemetry.dlq.v1and commits the offset, so it is not always redelivered. Document the transient retry limit and exhausted-retry/DLQ behavior here.
- **Failure handling**: On failure the offset is not committed and the batch is redelivered; in-process `batch_id` deduplication prevents double publication of redelivered batches
services/ipfs-publisher/src/index.ts:319
- A non-batched envelope is skipped without committing its offset. If another valid batch follows, its per-message commit advances past the skipped record; if it is the last record, it will be replayed indefinitely. Route unexpected event types to the existing retry/DLQ path or commit the ignored offset explicitly.
if (envelope.eventType !== TELEMETRY_TOPICS.BATCHED) {
logDebug('non-batched envelope ignored', {
eventType: envelope.eventType,
});
continue messageLoop;
services/ipfs-publisher/src/index.ts:331
- Protobuf decoding does not enforce these fields, so a malformed payload with an empty
batchIdis accepted and the first one is added to the dedup set. Every subsequent malformed payload with the same empty ID will then be skipped as a duplicate after commit, and empty batch bytes can be published to IPFS. Validate the batch ID and payload contents before deduplication and send invalid records to the DLQ.
const batched = fromBinary(
TelemetryBatchedPayloadSchema,
envelope.payload
) as TelemetryBatchedPayload;
metrics.consumed += 1;
// Skip batches already published in this process; still commit
// the offset so the duplicate is not redelivered forever.
if (dedup.has(batched.batchId)) {
services/ipfs-publisher/src/index.ts:429
- When DLQ publication or the follow-up commit fails,
break messageLoopends the consumer task permanently, butstartedremains true and the health endpoint continues to report 200. A transient DLQ/Kafka outage therefore leaves the service alive but no longer consuming until a manual restart. Retry the DLQ/commit with bounded backoff or propagate the failure so the process supervisor can restart it.
} catch (dlqError) {
// Cannot safely commit past this message: stop consuming
// rather than let a later message's commit skip past it.
logError(
'failed to route batch to DLQ; stopping consumption to avoid silent offset skip',
dlqError,
{
partition: message.partition,
offset: message.offset.toString(),
}
);
break messageLoop;
services/ipfs-publisher/src/index.ts:288
- This deduplication state is reset on every process restart and evicts IDs after 10,000 batches. A redelivery after a restart or eviction is therefore published to IPFS again; if the previous result send was accepted but its acknowledgment was lost, this can also emit another result and cause duplicate downstream anchoring. Use durable idempotency state or a deterministic result key with downstream CID deduplication.
const dedup = createBoundedDedup(10000);
services/ipfs-publisher/src/index.ts:128
- The DLQ record puts the original batched envelope bytes directly on
telemetry.dlq.v1, so the value still identifies itself as aBATCHEDevent rather than a canonical DLQ envelope. Existing DLQ publishing wraps a payload inEnvelopeSchemawitheventType: TELEMETRY_TOPICS.DLQ(services/endpoint/src/producer.ts:112-128). Use that envelope convention while preserving the raw record and failure context for DLQ consumers.
await producer.send({
messages: [
{
topic: TELEMETRY_TOPICS.DLQ,
value: Buffer.from(raw),
headers: {
source_topic: Buffer.from(TELEMETRY_TOPICS.BATCHED),
source_service: Buffer.from(config.source),
source_partition: Buffer.from(String(partition)),
source_offset: Buffer.from(String(offset)),
reason: Buffer.from(reason),
- Files reviewed: 25/26 changed files
- Comments generated: 2
- Review effort level: Lite
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
This pull request introduces a new standalone batcher service to the telemetry pipeline, separating the responsibilities of batching and IPFS publishing. The batcher consumes authorized telemetry, groups it into deterministic batches, and emits them to a new Kafka topic for downstream publication. The documentation, configuration, and codebase have been updated to reflect this architectural change, and new schemas and dependencies have been added.
Pipeline & Service Architecture Updates
Introduced a new
@scp/batcherservice that consumestelemetry.authorized.v1, batches events, and emits them astelemetry.batched.v1for downstream consumers. Updated all relevant documentation and diagrams to reflect the new pipeline: endpoint → batcher → ipfs-publisher → blockchain-anchor. [1] [2] [3] [4] [5] [6] [7]Updated service and topic lists in documentation and configuration files to include the batcher and new topics, and clarified the responsibilities of batcher vs. ipfs-publisher. [1] [2] [3] [4]
Kafka Topics & Environment
telemetry.batched.v1throughout the codebase, including topic constants, Docker Kafka initialization, and environment variable examples. Updated the IPFS publisher's topic toipfs.published.v1. [1] [2] [3] [4]Protobuf Schema Changes
TelemetryBatchedPayloadmessage to the protobuf schema, defining the structure of batched telemetry events produced by the batcher and consumed by the ipfs-publisher.New Service Implementation
services/batcherpackage with its ownpackage.json, dependencies, README, and core batching logic (batch-flusher.ts) implementing serialized, race-safe batch flushing with retry and offset commit safety. [1] [2] [3]These changes decouple batching from publishing, improving reliability, scalability, and observability of the telemetry pipeline.