Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 6 additions & 1 deletion src/confluent_kafka/src/Producer.c
Original file line number Diff line number Diff line change
Expand Up @@ -1175,7 +1175,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. */
Comment on lines +1178 to +1182
if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|ppp", kws, &in_queue,
&in_flight, &blocking))
return NULL;

Expand Down
51 changes: 51 additions & 0 deletions tests/test_Producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down