package kafka-eio
sectionYPositions = computeSectionYPositions($el), 10)"
x-init="setTimeout(() => sectionYPositions = computeSectionYPositions($el), 10)"
>
On This Page
Eio-native Kafka client for OCaml 5, built on librdkafka
Install
dune-project
Dependency
Authors
Maintainers
Sources
v0.1.1.tar.gz
md5=88c2b9c6d263a1eb114cc5d1b570288a
sha512=44f5395fcc58b8d4303bc0b36113a17193da3a1925125cb0742fabeb96f32eb6b243c98d2fde283b6c95872ff36ddbf96de6e98b7612799e2aefeb2f85a1db95
doc/CHANGES.html
Changes
0.1.1
- Add
/usr/localand Homebrew C include/library paths for librdkafka on FreeBSD and macOS. - Skip the
/proc/self/fdfd-leak regression test on platforms without procfs.
0.1.0
- Initial standalone OPAM package:
Kafka_producerandKafka_consumeron top of librdkafka, extracted from the Sun platform after real in-tree usage. Idempotent/ transactional delivery, consumer-group streaming with explicit ack, partition-aware concurrency, and typed transport security (Kafka_security.t: plaintext/SSL/SASL) on every producer, consumer, and service config. - Delivery receipts flow through a Unix pipe from the C delivery callback (thread-safe, no OCaml runtime needed from C); a background Eio fiber reads the pipe and resolves pending promises.
- Blocking C calls (
ocaml_rd_kafka_flush,ocaml_rd_kafka_consumer_close,ocaml_rd_kafka_consumer_poll) release the OCaml domain lock at the FFI boundary around the blocking librdkafka call, so the GC can run and Eio's cancellation delivers cleanly once the call returns — not worked around withEio_unix.run_in_systhreador polling loops at the OCaml level.
Post-tag fix (#1): mutex-unlock skipped under Eio cancellation
Every Mutex.lock t.mutex; ...; Mutex.unlock t.mutex pair in Kafka_producer (delivery callback dispatch, pending-promise bookkeeping in produce_await, and cleanup in close) and the poll_exit_r resolution in Kafka_consumer's poll_fiber skipped the unlock/resolve step if an Eio.Cancel.Cancelled (or any other exception) was raised between the lock and unlock — leaving the mutex held forever, or a fiber awaiting poll_exit_r hanging forever after close. Fixed by wrapping every such critical section in Fun.protect ~finally, so the unlock/resolve always runs regardless of how the protected code exits.
sectionYPositions = computeSectionYPositions($el), 10)"
x-init="setTimeout(() => sectionYPositions = computeSectionYPositions($el), 10)"
>
On This Page