Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
00325bd
Add a serializing Consumer reentrancy gate for free-threading support
ojasvajain Aug 18, 2026
32ff49e
Trigger CLA check
ojasvajain Aug 22, 2026
b43b8e5
Add Integration test coverage for AIO Producer
ojasvajain Aug 20, 2026
cc89579
Add Tx related test cases + Add abort_tx call in Producer close
ojasvajain Aug 21, 2026
f51339b
Fix styling
ojasvajain Aug 24, 2026
5b634e3
Serialize AIOConsumer calls that outlive their callback invocation in…
ojasvajain Aug 27, 2026
dae4611
Serialize AIOConsumer calls that outlive their callback invocation in…
ojasvajain Aug 27, 2026
08971b0
Trigger CLA check
ojasvajain Aug 22, 2026
b75f6fa
Trigger CLA check
ojasvajain Aug 22, 2026
14431f3
Add a serializing Consumer reentrancy gate for free-threading support
ojasvajain Aug 18, 2026
117cb5e
Style fixes
ojasvajain Aug 18, 2026
1617b76
Addressed comments
ojasvajain Aug 21, 2026
0bf4a1c
Trigger CLA check
ojasvajain Aug 22, 2026
31632cb
Address comments
ojasvajain Aug 24, 2026
61451ee
Add Tx related test cases + Add abort_tx call in Producer close
ojasvajain Aug 21, 2026
2adba42
Fix for serializing concurrent calls from a Async callback
ojasvajain Aug 24, 2026
e6bf843
Fix one test case
ojasvajain Aug 24, 2026
0a2acd5
Fix styling
ojasvajain Aug 24, 2026
9d3a6d0
[NOGIL] Guard AdminClient against concurrent close()-vs-method-call r…
ojasvajain Aug 24, 2026
e271be1
Add handle check in common APIs and refactor overall handle logic
ojasvajain Aug 25, 2026
32db4ff
Remove conflict marker
ojasvajain Aug 25, 2026
97c1c75
Remove another conflict marker
ojasvajain Aug 25, 2026
9f9b8c2
Rename handle functions
ojasvajain Aug 25, 2026
a2bd2db
Add more tests
ojasvajain Aug 25, 2026
cc91c79
Fix flaky test
ojasvajain Aug 26, 2026
3405f86
Address comments
ojasvajain Aug 27, 2026
069f6f3
Fix corrupted _common.py
ojasvajain Aug 28, 2026
2bec006
Fix other corrupted files
ojasvajain Aug 28, 2026
2dd1397
Fix flaky test case
ojasvajain Aug 29, 2026
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
614 changes: 358 additions & 256 deletions src/confluent_kafka/src/Admin.c

Large diffs are not rendered by default.

207 changes: 61 additions & 146 deletions src/confluent_kafka/src/Consumer.c

Large diffs are not rendered by default.

25 changes: 13 additions & 12 deletions src/confluent_kafka/src/Metadata.c
Original file line number Diff line number Diff line change
Expand Up @@ -366,18 +366,17 @@ PyObject *list_topics(Handle *self, PyObject *args, PyObject *kwargs) {
&tmout))
return NULL;

if (!self->rk) {
PyErr_SetString(PyExc_RuntimeError, ERR_MSG_HANDLE_CLOSED);
if (!Handle_common_enter(self))
return NULL;
}

if (topic != NULL) {
if (!(only_rkt = rd_kafka_topic_new(self->rk, topic, NULL))) {
return PyErr_Format(
PyExc_RuntimeError,
"Unable to create topic object "
"for \"%s\": %s",
topic, rd_kafka_err2str(rd_kafka_last_error()));
PyErr_Format(PyExc_RuntimeError,
"Unable to create topic object "
"for \"%s\": %s",
topic,
rd_kafka_err2str(rd_kafka_last_error()));
goto end; /* result and only_rkt are NULL */
}
}

Expand Down Expand Up @@ -409,6 +408,8 @@ PyObject *list_topics(Handle *self, PyObject *args, PyObject *kwargs) {
rd_kafka_topic_destroy(only_rkt);
}

Handle_common_exit(self);

return result;
}

Expand Down Expand Up @@ -608,7 +609,8 @@ PyObject *list_groups(Handle *self, PyObject *args, PyObject *kwargs) {
const struct rd_kafka_group_list *group_list = NULL;
const char *group = NULL;
double tmout = -1.0f;
static char *kws[] = {"group", "timeout", NULL};
static char *kws[] = {"group", "timeout",
NULL};

PyErr_WarnEx(PyExc_DeprecationWarning,
"list_groups() is deprecated, use list_consumer_groups() "
Expand All @@ -619,10 +621,8 @@ PyObject *list_groups(Handle *self, PyObject *args, PyObject *kwargs) {
&tmout))
return NULL;

if (!self->rk) {
PyErr_SetString(PyExc_RuntimeError, ERR_MSG_HANDLE_CLOSED);
if (!Handle_common_enter(self))
return NULL;
}

CallState_begin(self, &cs);

Expand All @@ -645,6 +645,7 @@ PyObject *list_groups(Handle *self, PyObject *args, PyObject *kwargs) {
if (group_list != NULL) {
rd_kafka_group_list_destroy(group_list);
}
Handle_common_exit(self);
return result;
}

Expand Down
67 changes: 30 additions & 37 deletions src/confluent_kafka/src/Producer.c
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,11 @@
*
****************************************************************************/

static int Producer_rk_use_begin(Handle *self) {
return Handle_rk_use_begin(self, ERR_MSG_PRODUCER_CLOSED);
}


/**
* Per-message state.
*/
Expand Down Expand Up @@ -302,7 +307,7 @@ Producer_produce(Handle *self, PyObject *args, PyObject *kwargs) {
if (!dr_cb || dr_cb == Py_None)
dr_cb = self->u.Producer.default_dr_cb;

if (!Handle_enter_rk_use(self)) {
if (!Producer_rk_use_begin(self)) {
#ifdef RD_KAFKA_V_HEADERS
if (rd_headers)
rd_kafka_headers_destroy(rd_headers);
Expand All @@ -328,7 +333,7 @@ Producer_produce(Handle *self, PyObject *args, PyObject *kwargs) {
key_len, msgstate);
#endif

Handle_exit_rk_use(self);
Handle_rk_use_end(self);

if (err) {
if (msgstate)
Expand Down Expand Up @@ -433,12 +438,12 @@ static PyObject *Producer_poll(Handle *self, PyObject *args, PyObject *kwargs) {
if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|d", kws, &tmout))
return NULL;

if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

r = Producer_poll0(self, cfl_timeout_ms(tmout));

Handle_exit_rk_use(self);
Handle_rk_use_end(self);

if (r == -1)
return NULL;
Expand Down Expand Up @@ -483,7 +488,7 @@ Producer_flush(Handle *self, PyObject *args, PyObject *kwargs) {
if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|d", kws, &tmout))
return NULL;

if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

total_timeout_ms = cfl_timeout_ms(tmout);
Expand Down Expand Up @@ -552,7 +557,7 @@ Producer_flush(Handle *self, PyObject *args, PyObject *kwargs) {
result = cfl_PyInt_FromInt(qlen);

exit:
Handle_exit_rk_use(self);
Handle_rk_use_end(self);
return result;
}

Expand All @@ -579,13 +584,7 @@ Producer_close(Handle *self, PyObject *args, PyObject *kwargs) {
*/
if (!atomic_int_cas(&self->closing, 0, 1)) {
while (self->rk && atomic_int_get(&self->closing)) {
CallState_begin(self, &cs);
#ifdef _WIN32
Sleep(100);
#else
usleep(100000);
#endif
if (!CallState_end(self, &cs))
if (!Handle_sleep(self, 100))
return NULL;
}
if (!self->rk)
Expand All @@ -598,25 +597,19 @@ Producer_close(Handle *self, PyObject *args, PyObject *kwargs) {
return NULL;
}

/* Record which thread won, so Handle_enter_rk_use() can let a
/* Record which thread won, so Handle_rk_use_begin() can let a
* reentrant call through if (and only if) it's this same thread --
* i.e. close()'s own delivery callback, fired synchronously from
* inside the flush() call below, calling back into the Producer. */
atomic_ulong_set(&self->closing_thread, PyThread_get_thread_ident());

/* Signal in-flight calls to stop, and wait for them to finish
* using self->rk before destroying it -- see Handle_enter_rk_use().
* using self->rk before destroying it -- see Handle_rk_use_begin().
* New calls will see `closing` and fail with ERR_MSG_PRODUCER_CLOSED. */
/* TODO NOGIL: replace this poll loop with a mutex/condvar wait so
* close() unblocks immediately instead of up to 100ms late. */
while (atomic_int_get(&self->active_calls) > 0) {
CallState_begin(self, &cs);
#ifdef _WIN32
Sleep(100);
#else
usleep(100000);
#endif
if (!CallState_end(self, &cs)) {
if (!Handle_sleep(self, 100)) {
/* Abort the attempt: rk was never touched, so
* reopen the gate for a future close() attempt.
*/
Expand Down Expand Up @@ -923,7 +916,7 @@ Producer_produce_batch(Handle *self, PyObject *args, PyObject *kwargs) {
return cfl_PyInt_FromInt(0);
}

if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

/* Allocate arrays for librdkafka messages and msgstates */
Expand Down Expand Up @@ -956,7 +949,7 @@ Producer_produce_batch(Handle *self, PyObject *args, PyObject *kwargs) {
if (rkt)
rd_kafka_topic_destroy(rkt);

Handle_exit_rk_use(self);
Handle_rk_use_end(self);

/* Cleanup resources not tied to self->rk */
if (rkmessages)
Expand All @@ -979,7 +972,7 @@ static PyObject *Producer_init_transactions(Handle *self, PyObject *args) {
if (!PyArg_ParseTuple(args, "|d", &tmout))
return NULL;

if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

CallState_begin(self, &cs);
Expand All @@ -1002,19 +995,19 @@ static PyObject *Producer_init_transactions(Handle *self, PyObject *args) {
Py_INCREF(result);

exit:
Handle_exit_rk_use(self);
Handle_rk_use_end(self);
return result;
}

static PyObject *Producer_begin_transaction(Handle *self) {
rd_kafka_error_t *error;

if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

error = rd_kafka_begin_transaction(self->rk);

Handle_exit_rk_use(self);
Handle_rk_use_end(self);

if (error) {
cfl_PyErr_from_error_destroy(error);
Expand All @@ -1037,7 +1030,7 @@ static PyObject *Producer_send_offsets_to_transaction(Handle *self,
if (!PyArg_ParseTuple(args, "OO|d", &offsets, &metadata, &tmout))
return NULL;

if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

if (!(c_offsets = py_to_c_parts(offsets)))
Expand Down Expand Up @@ -1071,7 +1064,7 @@ static PyObject *Producer_send_offsets_to_transaction(Handle *self,
rd_kafka_consumer_group_metadata_destroy(cgmd);
if (c_offsets)
rd_kafka_topic_partition_list_destroy(c_offsets);
Handle_exit_rk_use(self);
Handle_rk_use_end(self);
return result;
}

Expand All @@ -1084,7 +1077,7 @@ static PyObject *Producer_commit_transaction(Handle *self, PyObject *args) {
if (!PyArg_ParseTuple(args, "|d", &tmout))
return NULL;

if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

CallState_begin(self, &cs);
Expand All @@ -1107,7 +1100,7 @@ static PyObject *Producer_commit_transaction(Handle *self, PyObject *args) {
Py_INCREF(result);

exit:
Handle_exit_rk_use(self);
Handle_rk_use_end(self);
return result;
}

Expand All @@ -1122,7 +1115,7 @@ static PyObject *Producer_abort_transaction(Handle *self, PyObject *args) {

/* abort_transaction is called as part of close() so even if this call
* is rejected here (because a close is in progress), it is safe.*/
if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

CallState_begin(self, &cs);
Expand All @@ -1145,7 +1138,7 @@ static PyObject *Producer_abort_transaction(Handle *self, PyObject *args) {
Py_INCREF(result);

exit:
Handle_exit_rk_use(self);
Handle_rk_use_end(self);
return result;
}

Expand All @@ -1162,7 +1155,7 @@ static void *Producer_purge(Handle *self, PyObject *args, PyObject *kwargs) {
&in_flight, &blocking))
return NULL;

if (!Handle_enter_rk_use(self))
if (!Producer_rk_use_begin(self))
return NULL;

if (in_queue)
Expand All @@ -1174,7 +1167,7 @@ static void *Producer_purge(Handle *self, PyObject *args, PyObject *kwargs) {

err = rd_kafka_purge(self->rk, purge_strategy);

Handle_exit_rk_use(self);
Handle_rk_use_end(self);

if (err) {
cfl_PyErr_Format(err, "Purge failed: %s",
Expand Down Expand Up @@ -1533,7 +1526,7 @@ static PyMethodDef Producer_methods[] = {
static Py_ssize_t Producer__len__(Handle *self) {
Py_ssize_t len;

/* __len__ must never raise, so we can't use Handle_enter_rk_use()
/* __len__ must never raise, so we can't use Handle_rk_use_begin()
* (which sets an exception on failure) -- fall back to returning 0,
* , if the Handle is closed/closing. */
if (atomic_int_get(&self->closing) || !self->rk)
Expand Down
Loading
Loading