diff --git a/pulsar-functions/instance/src/main/python/python_instance.py b/pulsar-functions/instance/src/main/python/python_instance.py index 61322049f2128..672134c0689cc 100755 --- a/pulsar-functions/instance/src/main/python/python_instance.py +++ b/pulsar-functions/instance/src/main/python/python_instance.py @@ -150,6 +150,8 @@ def run(self): elif self.instance_config.function_details.retainKeyOrdering: mode = pulsar._pulsar.ConsumerType.KeyShared + nack_args = self.get_negative_ack_args() + position = pulsar._pulsar.InitialPosition.Latest if self.instance_config.function_details.source.subscriptionPosition == Function_pb2.SubscriptionPosition.Value("EARLIEST"): position = pulsar._pulsar.InitialPosition.Earliest @@ -181,7 +183,8 @@ def run(self): message_listener=partial(self.message_listener, self.input_serdes[topic], DEFAULT_SCHEMA), unacked_messages_timeout_ms=int(self.timeout_ms) if self.timeout_ms else None, initial_position=position, - properties=properties + properties=properties, + **nack_args ) for topic, consumer_conf in self.instance_config.function_details.source.inputSpecs.items(): @@ -207,6 +210,7 @@ def run(self): "properties": properties, "crypto_key_reader": crypto_key_reader } + consumer_args.update(nack_args) if consumer_conf.HasField("receiverQueueSize"): consumer_args["receiver_queue_size"] = consumer_conf.receiverQueueSize.value @@ -578,6 +582,26 @@ def get_record_class(self, class_name): except: pass return record_kclass + def get_negative_ack_args(self): + """Build the negative-ack redelivery delay argument for Client.subscribe(). + + Returns a dict to splat into the subscribe() call: either empty, or carrying + negative_ack_redelivery_delay_ms. + + SourceSpec.negativeAckRedeliveryDelayMs is a proto3 scalar with no presence, so an unset field + reads as 0. Only a positive value is forwarded, leaving the client default (60s) in place + otherwise - the same guard the Java runtime applies in JavaInstanceRunnable. + + The argument is omitted rather than passed as None because subscribe() validates it with + _check_type(int, ...) rather than _check_type_or_none, so None would fail for every function + that does not configure it. + """ + delay_ms = self.instance_config.function_details.source.negativeAckRedeliveryDelayMs + if delay_ms <= 0: + return {} + + return {"negative_ack_redelivery_delay_ms": delay_ms} + def get_crypto_reader(self, crypto_spec): crypto_key_reader = None if crypto_spec is not None: diff --git a/pulsar-functions/instance/src/test/python/test_python_instance.py b/pulsar-functions/instance/src/test/python/test_python_instance.py index 3b20ac4b54563..8667f13964b04 100644 --- a/pulsar-functions/instance/src/test/python/test_python_instance.py +++ b/pulsar-functions/instance/src/test/python/test_python_instance.py @@ -360,3 +360,39 @@ def test_batch_builder_reaches_the_producer(self): function_details.sink.producerSpec.batchBuilder = "KEY_BASED" kwargs = self._create_producer_kwargs(function_details) self.assertEqual(kwargs["batching_type"], pulsar.BatchingType.KeyBased) + +class TestNegativeAckRedeliveryDelay(unittest.TestCase): + """Covers SourceSpec.negativeAckRedeliveryDelayMs reaching the consumer. + + The runtime negatively acknowledges on failure but never configured the delay, so the client + default of 60s always applied. The Java runtime guards on > 0 in JavaInstanceRunnable. + """ + + def _instance(self, delay_ms=None): + function_details = Function_pb2.FunctionDetails() + function_details.sink.topic = "test_sink_topic" + if delay_ms is not None: + function_details.source.negativeAckRedeliveryDelayMs = delay_ms + + return PythonInstance('test_instance', 'test_func', '1.0', function_details, 100, 30, + 'user_code', Mock(), Mock(), 'test_cluster', 'test_url', None) + + def test_positive_delay_is_forwarded(self): + args = self._instance(delay_ms=5000).get_negative_ack_args() + self.assertEqual({"negative_ack_redelivery_delay_ms": 5000}, args) + + def test_unset_delay_is_omitted(self): + # proto3 scalar with no presence: unset reads as 0. The argument must be omitted rather than + # sent - subscribe() validates it with _check_type(int), so None would fail for every function + # that does not set it, and 0 would mean immediate redelivery instead of the 60s default. + self.assertEqual({}, self._instance().get_negative_ack_args()) + + def test_explicit_zero_is_omitted(self): + self.assertEqual({}, self._instance(delay_ms=0).get_negative_ack_args()) + + def test_result_is_splattable_into_subscribe_kwargs(self): + # The value is consumed via **nack_args and consumer_args.update(...), so it must be a dict + # with exactly the keyword subscribe() expects. + args = self._instance(delay_ms=250).get_negative_ack_args() + self.assertIsInstance(args, dict) + self.assertEqual(["negative_ack_redelivery_delay_ms"], list(args.keys()))