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
7 changes: 4 additions & 3 deletions include/pulsar/BatchReceivePolicy.h
Original file line number Diff line number Diff line change
Expand Up @@ -58,9 +58,10 @@ class PULSAR_PUBLIC BatchReceivePolicy {

/**
*
* @param maxNumMessage Max num message, if less than 0, it means no limit.
* @param maxNumBytes Max num bytes, if less than 0, it means no limit.
* @param timeoutMs If less than 0, it means no limit.
* @param maxNumMessage Max num message, a non-positive value means no limit.
* @param maxNumBytes Max num bytes, a non-positive value means no limit.
* @param timeoutMs The receive timeout, a non-positive value means no limit.
* @throws std::invalid_argument if all arguments are non-positive
*/
BatchReceivePolicy(int maxNumMessage, long maxNumBytes, long timeoutMs);

Expand Down
41 changes: 29 additions & 12 deletions include/pulsar/c/consumer_configuration.h
Original file line number Diff line number Diff line change
Expand Up @@ -89,12 +89,14 @@ typedef enum
pulsar_consumer_regex_sub_mode_AllTopics = 2
} pulsar_consumer_regex_subscription_mode;

// Though any field could be non-positive, if all of them are non-positive, this policy will be treated as
// invalid
typedef struct {
// Max num message, if less than 0, it means no limit.
int maxNumMessage;
// Max num bytes, if less than 0, it means no limit.
// Max num messages, a non-positive value means no limit.
int maxNumMessages;
Comment thread
shibd marked this conversation as resolved.
// Max num bytes, a non-positive value means no limit.
long maxNumBytes;
// If less than 0, it means no limit.
// The receive timeout, a non-positive value means no limit.
long timeoutMs;
} pulsar_consumer_batch_receive_policy_t;

Expand Down Expand Up @@ -340,18 +342,33 @@ pulsar_consumer_configuration_get_regex_subscription_mode(
/**
* Set batch receive policy.
*
* The default value: {maxNumMessage: -1, maxNumBytes: 10 * 1024 * 1024, timeoutMs: 100}
* @param consumer_configuration the consumer conf object.
* @param maxNumMessage default: Max num message, if less than 0, it means no limit.
* @param maxNumBytes Max num bytes, if less than 0, it means no limit.
* @param timeoutMs If less than 0, it means no limit.
* @param [in] consumer_configuration a non-null pointer of the consumer configuration
* @param [in] batch_receive_policy
* @return 0 on success and -1 on failure
*
* The possible failed reasons are:
* - batch_receive_policy is null
* - batch_receive_policy points to an invalid policy
*/
PULSAR_PUBLIC void pulsar_consumer_configuration_set_batch_receive_policy(
PULSAR_PUBLIC int pulsar_consumer_configuration_set_batch_receive_policy(
pulsar_consumer_configuration_t *consumer_configuration,
const pulsar_consumer_batch_receive_policy_t *batch_receive_policy);

PULSAR_PUBLIC pulsar_consumer_batch_receive_policy_t pulsar_consumer_configuration_get_batch_receive_policy(
pulsar_consumer_configuration_t *consumer_configuration);
/**
* Get the batch receive policy.
*
* @param [in] consumer_configuration a non-null pointer of the consumer configuration
* @param [out] batch_receive_policy
*
* If batch_receive_policy is not null, the instance that it points to will be updated to the batch receive
* policy of the consumer configuration.
*
* If the policy was never set before, the batch_receive_policy will be set with the following value:
* {maxNumMessage: -1, maxNumBytes: 10 * 1024 * 1024, timeoutMs: 100}
*/
PULSAR_PUBLIC void pulsar_consumer_configuration_get_batch_receive_policy(
pulsar_consumer_configuration_t *consumer_configuration,
pulsar_consumer_batch_receive_policy_t *batch_receive_policy);

// const CryptoKeyReaderPtr getCryptoKeyReader()
//
Expand Down
24 changes: 18 additions & 6 deletions lib/c/c_ConsumerConfiguration.cc
Original file line number Diff line number Diff line change
Expand Up @@ -251,19 +251,31 @@ pulsar_consumer_regex_subscription_mode pulsar_consumer_configuration_get_regex_
consumer_configuration->consumerConfiguration.getRegexSubscriptionMode();
}

void pulsar_consumer_configuration_set_batch_receive_policy(
int pulsar_consumer_configuration_set_batch_receive_policy(
pulsar_consumer_configuration_t *consumer_configuration,
const pulsar_consumer_batch_receive_policy_t *batch_receive_policy_t) {
pulsar::BatchReceivePolicy batchReceivePolicy(batch_receive_policy_t->maxNumMessage,
if (!batch_receive_policy_t) {
return -1;
}
if (batch_receive_policy_t->maxNumMessages <= 0 && batch_receive_policy_t->maxNumBytes <= 0 &&
batch_receive_policy_t->timeoutMs <= 0) {
return -1;
}
pulsar::BatchReceivePolicy batchReceivePolicy(batch_receive_policy_t->maxNumMessages,
batch_receive_policy_t->maxNumBytes,
batch_receive_policy_t->timeoutMs);
consumer_configuration->consumerConfiguration.setBatchReceivePolicy(batchReceivePolicy);
return 0;
}

pulsar_consumer_batch_receive_policy_t pulsar_consumer_configuration_get_batch_receive_policy(
pulsar_consumer_configuration_t *consumer_configuration) {
void pulsar_consumer_configuration_get_batch_receive_policy(
pulsar_consumer_configuration_t *consumer_configuration, pulsar_consumer_batch_receive_policy_t *policy) {
if (!policy) {
return;
}
pulsar::BatchReceivePolicy batchReceivePolicy =
consumer_configuration->consumerConfiguration.getBatchReceivePolicy();
return {batchReceivePolicy.getMaxNumMessages(), batchReceivePolicy.getMaxNumBytes(),
batchReceivePolicy.getTimeoutMs()};
policy->maxNumMessages = batchReceivePolicy.getMaxNumMessages();
policy->maxNumBytes = batchReceivePolicy.getMaxNumBytes();
policy->timeoutMs = batchReceivePolicy.getTimeoutMs();
}
35 changes: 28 additions & 7 deletions tests/c/c_ConsumerConfigurationTest.cc
Original file line number Diff line number Diff line change
Expand Up @@ -43,13 +43,34 @@ TEST(C_ConsumerConfigurationTest, testCApiConfig) {
ASSERT_EQ(pulsar_consumer_configuration_get_regex_subscription_mode(consumer_conf),
pulsar_consumer_regex_sub_mode_NonPersistentOnly);

pulsar_consumer_batch_receive_policy_t batch_receive_policy{10, 1000, 1000};
pulsar_consumer_configuration_set_batch_receive_policy(consumer_conf, &batch_receive_policy);
pulsar_consumer_batch_receive_policy_t get_batch_receive_policy =
pulsar_consumer_configuration_get_batch_receive_policy(consumer_conf);
ASSERT_EQ(get_batch_receive_policy.maxNumMessage, 10);
ASSERT_EQ(get_batch_receive_policy.maxNumBytes, 1000);
ASSERT_EQ(get_batch_receive_policy.timeoutMs, 1000);
pulsar_consumer_batch_receive_policy_t batch_receive_policy;
pulsar_consumer_configuration_get_batch_receive_policy(consumer_conf, &batch_receive_policy);
ASSERT_EQ(batch_receive_policy.maxNumMessages, -1);
ASSERT_EQ(batch_receive_policy.maxNumBytes, 10 * 1024 * 1024L);
ASSERT_EQ(batch_receive_policy.timeoutMs, 100L);

pulsar_consumer_batch_receive_policy_t new_batch_receive_policy{-1, -1, -1};
ASSERT_EQ(-1, pulsar_consumer_configuration_set_batch_receive_policy(consumer_conf, NULL));
ASSERT_EQ(
-1, pulsar_consumer_configuration_set_batch_receive_policy(consumer_conf, &new_batch_receive_policy));

new_batch_receive_policy.maxNumMessages = 100;
ASSERT_EQ(
0, pulsar_consumer_configuration_set_batch_receive_policy(consumer_conf, &new_batch_receive_policy));
pulsar_consumer_configuration_get_batch_receive_policy(consumer_conf, &batch_receive_policy);
ASSERT_EQ(batch_receive_policy.maxNumMessages, 100);

new_batch_receive_policy.maxNumBytes = 100L * 1024 * 1024 * 1024;
ASSERT_EQ(
0, pulsar_consumer_configuration_set_batch_receive_policy(consumer_conf, &new_batch_receive_policy));
pulsar_consumer_configuration_get_batch_receive_policy(consumer_conf, &batch_receive_policy);
ASSERT_EQ(batch_receive_policy.maxNumBytes, 100L * 1024 * 1024 * 1024);

new_batch_receive_policy.timeoutMs = 365L * 24 * 3600 * 1000;
ASSERT_EQ(
0, pulsar_consumer_configuration_set_batch_receive_policy(consumer_conf, &new_batch_receive_policy));
pulsar_consumer_configuration_get_batch_receive_policy(consumer_conf, &batch_receive_policy);
ASSERT_EQ(batch_receive_policy.timeoutMs, 365L * 24 * 3600 * 1000);

pulsar_consumer_configuration_free(consumer_conf);
}