Kafka implementation plan¶
Reference: Powertools TypeScript v2.35.0, commit 7bcc27b1574493f9452688673658f52b80c53847. Keep the standalone consumer distinct from Parser's existing Kafka schemas and envelope.
- KAF-SOURCE: Inspect the published consumer, primitive/JSON/Avro/Protobuf decoders, error exports and type declarations. Confirm lazy getters, repeated parsing, key/value differences, JSON fallback diagnostics, light event guard and schema metadata handling. Record the CommonJS reference entry and ESM loader boundary.
- KAF-CORE: Implement the independent consumer module without third-party dependencies, native Lambda wrapper, context propagation, ordered flattening, lazy reads, primitive/JSON decoding, headers, original metadata and typed errors. Reuse Commons Base64, UTF-8, numeric conversion and object ordering. Verify 143 actual TypeScript scenarios, 64 concurrent contexts, repeated reads, cancellation, original-field capture and error identity.
- KAF-LOCAL: Verify all 26 packaged modules and 23 standalone consumers, both CGO-disabled Linux builds, and native Lambda Parser/Idempotency/Logger/OTel composition. Passed 808/808 RIE assertions, 95/95 streaming checks and 14/14 saved Batch artifact checks (2026-09-23). Docker executed amd64; arm64 was cross-compiled. The ten Kafka archive files match current source. No AWS resources were used and disposable runtime resources were cleaned.
- KAF-AVRO-CORE: Implement the independent Avro adapter with hamba/avro schema parsing and reads, tagged unions, ignored logical annotations, reference floating-point zigzag behavior and lazy schema failure. Verify 370 actual TypeScript scenarios and native concurrency/isolation/cancellation checks.
- KAF-AVRO: Complete remaining invalid-schema diagnostics, overlong/oversized wire values, filesystem/prototype/native serialization boundaries and performance gates.
- KAF-PROTOBUF-CORE: Implement the independent callback/native-descriptor adapter, original buffer/position/length contract, Glue prefix, Confluent int32/sint32 byte skips, atomic adaptive preference and first-error retention. Verify 123 actual prefix scenarios and nine native proto2 descriptor/message scenarios, plus concurrent and reentrant callbacks.
- KAF-PROTOBUF: Complete native defaults/presence/unknown-wire/error/description serialization and arbitrary reentrant scheduling boundaries; establish performance/resource budgets.
- KAF-BINARY-LOCAL: Verify all 28 packaged modules and 25 standalone consumers with GOWORK=off, tests/vet/tidy and dependency isolation; build normal and streaming Lambda handlers for Linux amd64/arm64 with CGO_ENABLED=0. Passed 826/826 Docker RIE assertions, 95/95 streaming checks and 14/14 saved Batch checks (2026-09-23). Both adapter archives match all six current files. Docker executed amd64; arm64 was cross-compiled. Runtime resources were cleaned, and no AWS resources were used.
- KAF-MODES-JSON: Verify 22 additional pinned-consumer scenarios for MSK/self-managed JSON delivery, selected key/value attributes, original schema metadata, no-registry events and explicit codec selection. The 165-case core passes packaged tests/vet/tidy and an independent consumer build with no external modules. These are constructed events, not live captures; full SOURCE/service acceptance remains open. See KAFKA_MODES.md.
- KAF-MODES-SOURCE: Verify 172 constructed mixed SOURCE/JSON events against the pinned consumer, directly and through the native Go Lambda SDK. Cover both sources, Glue/Confluent metadata, every text/JSON/Avro/Protobuf key/value pair, parsers, original fields and null/empty/missing values. Record the SDK KafkaRecord metadata/presence loss and verify the RawMessage wrapper. See KAFKA_MODES.md.
- KAF-PROTOBUF-LENGTH: Match the upstream schemaId length comparison for JSON object/array IDs by reusing Commons number parsing. Verify 97 additional cases (220 prefix cases total), nine native messages and the 172 mixed events; retain cyclic/prototype boundaries explicitly.
- KAF-MODES-LOCAL: Verify the changed Protobuf and integration packaged modules with GOWORK=off, tests/vet/tidy and the independent Protobuf consumer. Session checkpoint-05 matches all six Protobuf and 45 integration archive files. Rebuilt Linux amd64/arm64 with CGO_ENABLED=0 and passed 829/829 RIE, 95/95 streaming and 14/14 Batch checks (2026-09-23). The runner reused completed module checks; Docker executed amd64 and cleanup succeeded. No AWS resources were used.
- KAF-MODES: Verify real SOURCE/JSON event examples with and without registry integration. The consumer must not invent a schema registry client: upstream consumes the provided schema and Lambda metadata. Include native self-managed Kafka events even though the reference declarations use MSK names.
- KAF-PARITY: Complete public type/API mappings, native/serialization/error boundaries, malformed header coercion, configuration mutation and async mapping, large integers and schema compatibility. Retain known differences in KAFKA.md.
- KAF-SERVICE: Validate actual event source retry/failure behavior when cloud testing is requested, then establish allocation/latency/payload budgets and release requirements.
The core's Decoder callback is now implemented by separate optional Avro and Protobuf modules. See KAFKA_BINARY.md for usage, native mappings and remaining boundaries. Current acceptance is recorded in MODULE_ACCEPTANCE.json, LOCAL_ACCEPTANCE.json, STREAMING_ACCEPTANCE.json and BATCH_ACCEPTANCE.json. The binary-adapter acceptance run executed all module checks and builds, then passed 826 RIE and 95 streaming checks. The subsequent core-only check in checkpoint-09 added 22 JSON-mode reference scenarios without changing production or integration Go code. The later Protobuf length-coercion fix and integration event corpus passed the scoped check in checkpoint-05; see KAF-MODES-LOCAL. Scoped report files track the most recent selected-module check, while complete workspace reports include subsequent Kafka regression coverage. The earlier primitive/JSON milestone required a runtime-only retry after correcting obsolete aggregate Idempotency counters; that historical run is documented in LOCAL_VALIDATION.md. All Go commands and subprocesses use CGO_ENABLED=0.
Binary adapter findings¶
The installed v2.35.0 package dynamically loads avro-js and protobufjs; neither is a mandatory dependency in its published manifest. Keep the corresponding Go libraries outside the core module. The reference fixture lock currently uses avro-js 1.12.1 and protobufjs 7.5.4, and uses the published CommonJS entry to avoid the observed ESM loader failure.
Avro source inspection shows that fromBuffer rejects trailing data and truncated input. Its long reader rejects potential precision loss rather than silently rounding every int64 to a JavaScript number. Union values preserve branch wrappers; bytes and fixed values are Node Buffers. These are acceptance requirements when evaluating a Go decoder, not optional output normalization. The Avro adapter now isolates hamba/avro v2.31.0 in its own module, uses a fresh schema cache and primitive reader, and owns the traversal needed to preserve reference union/logical/number behavior. The Kafka core retains no third-party dependencies.
The Protobuf source takes a caller-supplied message decoder and the event's schema metadata. An object without schemaId uses the whole buffer; absent metadata itself currently fails inside the decoder. A schemaId longer than ten characters selects a Glue path that consumes one uint32 before decoding. Other schema IDs select the Confluent path: read the index count as int32 or sint32 and skip that number of bytes, then decode. If the initial convention fails but the second succeeds, upstream reverses its module-level preferred order for later messages. A Go adapter must account for that shared state under concurrent invocations. The Kafka core passes metadata unchanged; the optional Protobuf module now implements these transformations with direct reference coverage. Native Go Protobuf object and error behavior remains explicit.