From 370a8311d8f4150ec51080132b70ba0712e46dde Mon Sep 17 00:00:00 2001 From: Devarsh Patel Date: Thu, 3 Sep 2026 13:07:37 +0530 Subject: [PATCH 1/3] Fix Producer.purge() ignoring boolean flags on big-endian platforms --- CHANGELOG.md | 5 +++++ src/confluent_kafka/src/Producer.c | 7 ++++++- 2 files changed, 11 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7e1c2a9a1..b069d6a9d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,11 @@ - Fix segmentation fault after calling `AdminClient.delete_records()` followed by another Admin API call (e.g. `list_topics()`) on Python 3.14. - Use `asyncio.get_running_loop()` instead of `asyncio.get_event_loop()` to avoid creating a new event loop and raise an error in case a loop isn't available (@AlexCai26, #2339). +- Fix `Producer.purge()` ignoring its `in_queue`, `in_flight` and `blocking` + arguments on big-endian platforms (e.g. s390x). They were parsed with the + 1-byte `"b"` format into 4-byte `int` targets, so on big-endian the value + landed on the high byte and the flags could not be cleared; for example + `purge(in_queue=False)` purged the queue anyway. ## v2.15.0 diff --git a/src/confluent_kafka/src/Producer.c b/src/confluent_kafka/src/Producer.c index 726262fb4..a981faafd 100644 --- a/src/confluent_kafka/src/Producer.c +++ b/src/confluent_kafka/src/Producer.c @@ -1034,7 +1034,12 @@ static void *Producer_purge(Handle *self, PyObject *args, PyObject *kwargs) { rd_kafka_resp_err_t err; static char *kws[] = {"in_queue", "in_flight", "blocking", NULL}; - if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|bbb", kws, &in_queue, + /* Use "p" (bool predicate -> int), not "b" (one byte): the targets are + * 4-byte ints, so "b" stores a single byte, which lands on the low byte + * on little-endian but the high byte on big-endian. There the flags, + * pre-initialised to 1, could never be cleared, so e.g. in_queue=False + * was ignored and the queue was purged anyway. */ + if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|ppp", kws, &in_queue, &in_flight, &blocking)) return NULL; From 114e05f6746941d1e4121c3fd1d2c751d274816b Mon Sep 17 00:00:00 2001 From: Devarsh Patel Date: Tue, 29 Sep 2026 11:18:02 +0530 Subject: [PATCH 2/3] Move the purge fix CHANGELOG entry to 2.16.0 --- CHANGELOG.md | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3a737c5e3..bb77a1d66 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,7 @@ v2.16.0 is a feature release with the following features, fixes and enhancements - Prefer httpx2 over httpx for Schema Registry to avoid Authlib deprecation warnings (#2351) - Fix KafkaError error strings raising/garbling on non-UTF-8 locales (#2331) - Fix crash on nullable array of $ref items in JSON Schema CSFLE (#2370) +- Fix `Producer.purge()` ignoring `in_queue`, `in_flight` and `blocking` set to `False` on big-endian platforms such as s390x (#2345) ## v2.15.1 @@ -44,11 +45,6 @@ v2.16.0 is a feature release with the following features, fixes and enhancements - Fix segmentation fault after calling `AdminClient.delete_records()` followed by another Admin API call (e.g. `list_topics()`) on Python 3.14. - Use `asyncio.get_running_loop()` instead of `asyncio.get_event_loop()` to avoid creating a new event loop and raise an error in case a loop isn't available (@AlexCai26, #2339). -- Fix `Producer.purge()` ignoring its `in_queue`, `in_flight` and `blocking` - arguments on big-endian platforms (e.g. s390x). They were parsed with the - 1-byte `"b"` format into 4-byte `int` targets, so on big-endian the value - landed on the high byte and the flags could not be cleared; for example - `purge(in_queue=False)` purged the queue anyway. confluent-kafka-python 2.15.1 is based on librdkafka 2.15.1, see the From c2b8f6a7f2dc1b8cb03d021fa14fbbfe51b6e820 Mon Sep 17 00:00:00 2001 From: Devarsh Patel Date: Tue, 29 Sep 2026 11:18:26 +0530 Subject: [PATCH 3/3] Test that purge() honours flags set to False --- tests/test_Producer.py | 51 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 51 insertions(+) diff --git a/tests/test_Producer.py b/tests/test_Producer.py index be345d4c0..497ae86ad 100644 --- a/tests/test_Producer.py +++ b/tests/test_Producer.py @@ -266,6 +266,57 @@ def on_delivery(err, msg): assert p.close(), "The producer was not closed" +@pytest.mark.parametrize( + "args, kwargs", + [ + ((), {"in_queue": False}), + ((), {"in_queue": False, "in_flight": False}), + ((), {"in_queue": False, "in_flight": False, "blocking": False}), + ((False, False, False), {}), + ], +) +def test_purge_keeps_queued_messages_when_in_queue_is_false(args, kwargs): + """ + purge() must honour flags set to False, however they are passed: with + in_queue=False a queued message stays queued, with no delivery report. + The flags used to be parsed with the one-byte "b" format, which on + big-endian platforms (e.g. s390x) left them stuck at True, so the queue + was purged anyway. + """ + p = Producer({"socket.timeout.ms": 10, "error_cb": error_cb, "message.timeout.ms": 30000}) + errors = [] + p.produce(topic="some_topic", value="testing", partition=9, callback=lambda err, msg: errors.append(err)) + + p.purge(*args, **kwargs) + p.flush(0.002) + assert errors == [] + assert len(p) == 1 + + p.purge() + p.flush(0.002) + assert [err.code() for err in errors] == [KafkaError._PURGE_QUEUE] + assert p.close(), "The producer was not closed" + + +@pytest.mark.parametrize("falsy", [0, None]) +def test_purge_flags_accept_any_falsy_value(falsy): + """ + 0 and None turn a purge() flag off, the same as False. + """ + p = Producer({"socket.timeout.ms": 10, "error_cb": error_cb, "message.timeout.ms": 30000}) + errors = [] + p.produce(topic="some_topic", value="testing", partition=9, callback=lambda err, msg: errors.append(err)) + + p.purge(in_queue=falsy) + p.flush(0.002) + assert errors == [] + + p.purge() + p.flush(0.002) + assert [err.code() for err in errors] == [KafkaError._PURGE_QUEUE] + assert p.close(), "The producer was not closed" + + def test_producer_bool_value(): """ Make sure producer has a truth-y bool value