Skip to content
Open
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
129 changes: 96 additions & 33 deletions src/confluent_kafka/src/Producer.c
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,12 @@

#include "confluent_kafka.h"

#ifdef _WIN32
#include <windows.h>
#else
#include <unistd.h>
#endif


/**
* @brief KNOWN ISSUES
Expand Down Expand Up @@ -296,12 +302,11 @@
if (!dr_cb || dr_cb == Py_None)
dr_cb = self->u.Producer.default_dr_cb;

if (!self->rk) {
if (!Handle_enter_rk_use(self)) {
#ifdef RD_KAFKA_V_HEADERS
if (rd_headers)
rd_kafka_headers_destroy(rd_headers);
#endif
PyErr_SetString(PyExc_RuntimeError, ERR_MSG_PRODUCER_CLOSED);
return NULL;
}

Expand All @@ -323,6 +328,8 @@
key_len, msgstate);
#endif

Handle_exit_rk_use(self);

if (err) {
if (msgstate)
Producer_msgstate_destroy(msgstate);
Expand Down Expand Up @@ -421,12 +428,13 @@
if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|d", kws, &tmout))
return NULL;

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

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

Handle_exit_rk_use(self);

if (r == -1)
return NULL;

Expand Down Expand Up @@ -469,10 +477,8 @@
if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|d", kws, &tmout))
return NULL;

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

total_timeout_ms = cfl_timeout_ms(tmout);
CallState_begin(self, &cs);
Expand Down Expand Up @@ -507,6 +513,7 @@
* interruptibility) */
chunk_count++;
if (check_signals_between_chunks(self, &cs)) {
Handle_exit_rk_use(self);
return NULL; /* Signal detected */
}

Expand All @@ -524,12 +531,16 @@
}
}

if (!CallState_end(self, &cs))
if (!CallState_end(self, &cs)) {
Handle_exit_rk_use(self);
return NULL;
}

if (err) /* Get the queue length on error (timeout) */
qlen = rd_kafka_outq_len(self->rk);

Handle_exit_rk_use(self);

return cfl_PyInt_FromInt(qlen);
}

Expand All @@ -542,6 +553,37 @@
if (!self->rk)
Py_RETURN_TRUE;

/* Only one concurrent close() can destroy rk, otherwise,
* two threads could both reach rd_kafka_destroy() on the
* same handle (a double-free). The losing thread(s) wait for the
* winner to finish and then return True, same as a normal close(),
* rather than racing it. */
if (!atomic_int_cas(&self->closing, 0, 1)) {
while (self->rk) {
CallState_begin(self, &cs);
#ifdef _WIN32
Sleep(100);
#else
usleep(100000);

Check warning on line 567 in src/confluent_kafka/src/Producer.c

View check run for this annotation

SonarQube-Confluent / SonarQube Code Analysis

Remove use of this obsolete "usleep" function. Replace it by a call to "nanosleep" or "setitimer".

[S1911] Obsolete POSIX functions should not be used See more on https://sonarqube.confluent.io/project/issues?id=confluent-kafka-python&pullRequest=2313&issues=8d11b50b-d243-4507-8ff4-d90aba771682&open=8d11b50b-d243-4507-8ff4-d90aba771682
#endif
CallState_end(self, &cs);
}
Py_RETURN_TRUE;
}

/* Signal in-flight calls to stop, and wait for them to finish
* using self->rk before destroying it -- see Handle_enter_rk_use().
* New calls will see `closing` and fail with ERR_MSG_PRODUCER_CLOSED. */
while (atomic_int_get(&self->active_calls) > 0) {
CallState_begin(self, &cs);
#ifdef _WIN32
Sleep(100);
#else
usleep(100000);

Check warning on line 582 in src/confluent_kafka/src/Producer.c

View check run for this annotation

SonarQube-Confluent / SonarQube Code Analysis

Remove use of this obsolete "usleep" function. Replace it by a call to "nanosleep" or "setitimer".

[S1911] Obsolete POSIX functions should not be used See more on https://sonarqube.confluent.io/project/issues?id=confluent-kafka-python&pullRequest=2313&issues=50f38792-4a85-49bc-bde2-26a440c113fa&open=50f38792-4a85-49bc-bde2-26a440c113fa
#endif
CallState_end(self, &cs);
}

CallState_begin(self, &cs);

/* Flush any pending messages (wait indefinitely to ensure delivery) */
Expand Down Expand Up @@ -817,10 +859,8 @@
return cfl_PyInt_FromInt(0);
}

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

/* Allocate arrays for librdkafka messages and msgstates */
rkmessages = calloc(message_cnt, sizeof(*rkmessages));
Expand Down Expand Up @@ -849,6 +889,8 @@
messages_list, rkt, partition, rkmessages, msgstates, message_cnt);

cleanup:
Handle_exit_rk_use(self);

/* Cleanup resources */
if (rkt)
rd_kafka_topic_destroy(rkt);
Expand All @@ -871,21 +913,22 @@
if (!PyArg_ParseTuple(args, "|d", &tmout))
return NULL;

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

CallState_begin(self, &cs);

error = rd_kafka_init_transactions(self->rk, cfl_timeout_ms(tmout));

if (!CallState_end(self, &cs)) {
Handle_exit_rk_use(self);
if (error) /* Ignore error in favour of callstate exception */
rd_kafka_error_destroy(error);
return NULL;
}

Handle_exit_rk_use(self);

if (error) {
cfl_PyErr_from_error_destroy(error);
return NULL;
Expand All @@ -897,13 +940,13 @@
static PyObject *Producer_begin_transaction(Handle *self) {
rd_kafka_error_t *error;

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

error = rd_kafka_begin_transaction(self->rk);

Handle_exit_rk_use(self);

if (error) {
cfl_PyErr_from_error_destroy(error);
return NULL;
Expand All @@ -924,16 +967,17 @@
if (!PyArg_ParseTuple(args, "OO|d", &offsets, &metadata, &tmout))
return NULL;

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

if (!(c_offsets = py_to_c_parts(offsets)))
if (!(c_offsets = py_to_c_parts(offsets))) {
Handle_exit_rk_use(self);
return NULL;
}

if (!(cgmd = py_to_c_cgmd(metadata))) {
rd_kafka_topic_partition_list_destroy(c_offsets);
Handle_exit_rk_use(self);
return NULL;
}

Expand All @@ -946,11 +990,14 @@
rd_kafka_topic_partition_list_destroy(c_offsets);

if (!CallState_end(self, &cs)) {
Handle_exit_rk_use(self);
if (error) /* Ignore error in favour of callstate exception */
rd_kafka_error_destroy(error);
return NULL;
}

Handle_exit_rk_use(self);

if (error) {
cfl_PyErr_from_error_destroy(error);
return NULL;
Expand All @@ -967,21 +1014,22 @@
if (!PyArg_ParseTuple(args, "|d", &tmout))
return NULL;

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

CallState_begin(self, &cs);

error = rd_kafka_commit_transaction(self->rk, cfl_timeout_ms(tmout));

if (!CallState_end(self, &cs)) {
Handle_exit_rk_use(self);
if (error) /* Ignore error in favour of callstate exception */
rd_kafka_error_destroy(error);
return NULL;
}

Handle_exit_rk_use(self);

if (error) {
cfl_PyErr_from_error_destroy(error);
return NULL;
Expand All @@ -998,21 +1046,22 @@
if (!PyArg_ParseTuple(args, "|d", &tmout))
return NULL;

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

CallState_begin(self, &cs);

error = rd_kafka_abort_transaction(self->rk, cfl_timeout_ms(tmout));

if (!CallState_end(self, &cs)) {
Handle_exit_rk_use(self);
if (error) /* Ignore error in favour of callstate exception */
rd_kafka_error_destroy(error);
return NULL;
}

Handle_exit_rk_use(self);

if (error) {
cfl_PyErr_from_error_destroy(error);
return NULL;
Expand All @@ -1034,10 +1083,8 @@
&in_flight, &blocking))
return NULL;

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

if (in_queue)
purge_strategy = RD_KAFKA_PURGE_F_QUEUE;
Expand All @@ -1048,6 +1095,8 @@

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

Handle_exit_rk_use(self);

if (err) {
cfl_PyErr_Format(err, "Purge failed: %s",
rd_kafka_err2str(err));
Expand Down Expand Up @@ -1403,9 +1452,23 @@


static Py_ssize_t Producer__len__(Handle *self) {
if (!self->rk)
Py_ssize_t len;

/* __len__ must never raise, so we can't use Handle_enter_rk_use()
* (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)
return 0;
atomic_int_inc(&self->active_calls);
if (atomic_int_get(&self->closing) || !self->rk) {
atomic_int_dec(&self->active_calls);
return 0;
return rd_kafka_outq_len(self->rk);
}

len = rd_kafka_outq_len(self->rk);

atomic_int_dec(&self->active_calls);
return len;
}


Expand Down
Loading
Loading