From 9ad2670d507f87170e4d7d65c018b3b739c0d5ae Mon Sep 17 00:00:00 2001 From: Ankit Kumar Date: Wed, 22 Jul 2026 21:08:39 +0530 Subject: [PATCH] fix(binding-kafka): guard budget credit against released creditor slot in cache client produce fan KafkaCacheClientProduceFan.onClientInitialData/onClientInitialFlush called creditor.credit(traceId, partitionIndex, reserved) unconditionally, even after onClientFanInitialClosed() had already released the budget slot and reset partitionIndex to NO_CREDITOR_INDEX (-1L). A member stream whose produce reply was closed via onClientFanReplyEnd/onClientFanReplyAbort keeps its own initial open (only the reply side is forwarded to members), so a still-producing app client can deliver another DATA/FLUSH frame that reaches the fan and credits a released slot. DefaultBudgetCreditor.credit computes a garbage negative buffer offset from the sentinel index and throws an IndexOutOfBoundsException that terminates the whole EngineWorker agent thread, dropping every unrelated connection on that worker. Guard both credit() call sites on partitionIndex != NO_CREDITOR_INDEX, and mirror the #1303 stream-level guard (!KafkaState.initialClosed(state)) onto the FLUSH path, which had been left with a bare else since that fix only covered DATA. Fixes #2179 Co-Authored-By: Claude Sonnet 5 --- .../KafkaCacheClientProduceFactory.java | 13 ++- .../kafka/internal/stream/CacheProduceIT.java | 11 +++ .../client.rpt | 90 +++++++++++++++++++ .../server.rpt | 77 ++++++++++++++++ .../kafka/streams/application/ProduceIT.java | 9 ++ 5 files changed, 197 insertions(+), 3 deletions(-) create mode 100644 specs/binding-kafka.spec/src/main/scripts/io/aklivity/zilla/specs/binding/kafka/streams/application/produce/message.value.after.fan.reply.abort/client.rpt create mode 100644 specs/binding-kafka.spec/src/main/scripts/io/aklivity/zilla/specs/binding/kafka/streams/application/produce/message.value.after.fan.reply.abort/server.rpt diff --git a/runtime/binding-kafka/src/main/java/io/aklivity/zilla/runtime/binding/kafka/internal/stream/KafkaCacheClientProduceFactory.java b/runtime/binding-kafka/src/main/java/io/aklivity/zilla/runtime/binding/kafka/internal/stream/KafkaCacheClientProduceFactory.java index fffa793b770..514bb5e3b69 100644 --- a/runtime/binding-kafka/src/main/java/io/aklivity/zilla/runtime/binding/kafka/internal/stream/KafkaCacheClientProduceFactory.java +++ b/runtime/binding-kafka/src/main/java/io/aklivity/zilla/runtime/binding/kafka/internal/stream/KafkaCacheClientProduceFactory.java @@ -779,7 +779,10 @@ private void onClientInitialData( onClientFanMemberClosed(traceId, stream); } - creditor.credit(traceId, partitionIndex, reserved); + if (partitionIndex != NO_CREDITOR_INDEX) + { + creditor.credit(traceId, partitionIndex, reserved); + } } private void onClientInitialFlush( @@ -823,7 +826,11 @@ stream.valueMark, stream.valueLimit, now().toEpochMilli(), stream.initialId, PRO stream.cleanupClient(traceId, error); onClientFanMemberClosed(traceId, stream); } - creditor.credit(traceId, partitionIndex, reserved); + + if (partitionIndex != NO_CREDITOR_INDEX) + { + creditor.credit(traceId, partitionIndex, reserved); + } } private void flushClientFanInitialIfNecessary( @@ -1394,7 +1401,7 @@ private void onClientInitialFlush( doClientReplyAbortIfNecessary(traceId); fan.onClientFanMemberClosed(traceId, this); } - else + else if (!KafkaState.initialClosed(state)) { fan.onClientInitialFlush(this, flush); } diff --git a/runtime/binding-kafka/src/test/java/io/aklivity/zilla/runtime/binding/kafka/internal/stream/CacheProduceIT.java b/runtime/binding-kafka/src/test/java/io/aklivity/zilla/runtime/binding/kafka/internal/stream/CacheProduceIT.java index 1f7704a55a3..aec2d2c362e 100644 --- a/runtime/binding-kafka/src/test/java/io/aklivity/zilla/runtime/binding/kafka/internal/stream/CacheProduceIT.java +++ b/runtime/binding-kafka/src/test/java/io/aklivity/zilla/runtime/binding/kafka/internal/stream/CacheProduceIT.java @@ -213,6 +213,17 @@ public void shouldSendMessageValueNull() throws Exception k3po.finish(); } + @Test + @Configuration("cache.yaml") + @Specification({ + "${app}/message.value.after.fan.reply.abort/client", + "${app}/message.value.after.fan.reply.abort/server"}) + @ScriptProperty("serverAddress \"zilla://streams/app1\"") + public void shouldSendMessageValueAfterFanReplyAbort() throws Exception + { + k3po.finish(); + } + @Test @Configuration("cache.yaml") @Specification({ diff --git a/specs/binding-kafka.spec/src/main/scripts/io/aklivity/zilla/specs/binding/kafka/streams/application/produce/message.value.after.fan.reply.abort/client.rpt b/specs/binding-kafka.spec/src/main/scripts/io/aklivity/zilla/specs/binding/kafka/streams/application/produce/message.value.after.fan.reply.abort/client.rpt new file mode 100644 index 00000000000..81523cbac8b --- /dev/null +++ b/specs/binding-kafka.spec/src/main/scripts/io/aklivity/zilla/specs/binding/kafka/streams/application/produce/message.value.after.fan.reply.abort/client.rpt @@ -0,0 +1,90 @@ +# +# Copyright 2021-2026 Aklivity Inc. +# +# Aklivity licenses this file to you under the Apache License, +# version 2.0 (the "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at: +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations +# under the License. +# + +property deltaMillis 0L + +connect "zilla://streams/app0" + option zilla:window 8192 + option zilla:transmission "half-duplex" + +write zilla:begin.ext ${kafka:beginEx() + .typeId(zilla:id("kafka")) + .meta() + .topic("test") + .build() + .build()} + +connected + +read zilla:begin.ext ${kafka:beginEx() + .typeId(zilla:id("kafka")) + .meta() + .topic("test") + .build() + .build()} + +read zilla:data.ext ${kafka:dataEx() + .typeId(zilla:id("kafka")) + .meta() + .partition(0, 177) + .build() + .build()} + +read notify ROUTED_BROKER_CLIENT + +connect await ROUTED_BROKER_CLIENT + "zilla://streams/app0" + option zilla:window 8192 + option zilla:transmission "half-duplex" + option zilla:affinity 0xb1 + +write zilla:begin.ext ${kafka:beginEx() + .typeId(zilla:id("kafka")) + .produce() + .topic("test") + .partition(0) + .build() + .build()} + +connected + +read zilla:begin.ext ${kafka:beginEx() + .typeId(zilla:id("kafka")) + .produce() + .topic("test") + .partition(0) + .build() + .build()} + +write zilla:data.ext ${kafka:dataEx() + .typeId(zilla:id("kafka")) + .produce() + .timestamp(1716424650323) + .build() + .build()} +write "Hello, world" +write flush + +read aborted + +write zilla:data.ext ${kafka:dataEx() + .typeId(zilla:id("kafka")) + .produce() + .timestamp(1716424650323) + .build() + .build()} +write "Hello, again" +write flush diff --git a/specs/binding-kafka.spec/src/main/scripts/io/aklivity/zilla/specs/binding/kafka/streams/application/produce/message.value.after.fan.reply.abort/server.rpt b/specs/binding-kafka.spec/src/main/scripts/io/aklivity/zilla/specs/binding/kafka/streams/application/produce/message.value.after.fan.reply.abort/server.rpt new file mode 100644 index 00000000000..b784aa4ed92 --- /dev/null +++ b/specs/binding-kafka.spec/src/main/scripts/io/aklivity/zilla/specs/binding/kafka/streams/application/produce/message.value.after.fan.reply.abort/server.rpt @@ -0,0 +1,77 @@ +# +# Copyright 2021-2026 Aklivity Inc. +# +# Aklivity licenses this file to you under the Apache License, +# version 2.0 (the "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at: +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations +# under the License. +# + +property serverAddress "zilla://streams/app0" + +accept ${serverAddress} + option zilla:window 8192 + option zilla:transmission "duplex" + option zilla:update "handshake" + +accepted + +read zilla:begin.ext ${kafka:beginEx() + .typeId(zilla:id("kafka")) + .meta() + .topic("test") + .build() + .build()} + +write zilla:begin.ext ${kafka:beginEx() + .typeId(zilla:id("kafka")) + .meta() + .topic("test") + .build() + .build()} + +connected + +write zilla:data.ext ${kafka:dataEx() + .typeId(zilla:id("kafka")) + .meta() + .partition(0, 177) + .build() + .build()} +write flush + +accepted + +read zilla:begin.ext ${kafka:beginEx() + .typeId(zilla:id("kafka")) + .produce() + .topic("test") + .partition(0) + .build() + .build()} + +write zilla:begin.ext ${kafka:beginEx() + .typeId(zilla:id("kafka")) + .produce() + .topic("test") + .partition(0) + .build() + .build()} + +connected + +read zilla:data.ext ${kafka:matchDataEx() + .typeId(zilla:id("kafka")) + .produce() + .build() + .build()} +read "Hello, world" + +write abort diff --git a/specs/binding-kafka.spec/src/test/java/io/aklivity/zilla/specs/binding/kafka/streams/application/ProduceIT.java b/specs/binding-kafka.spec/src/test/java/io/aklivity/zilla/specs/binding/kafka/streams/application/ProduceIT.java index 913b44f4803..47f2446a43c 100644 --- a/specs/binding-kafka.spec/src/test/java/io/aklivity/zilla/specs/binding/kafka/streams/application/ProduceIT.java +++ b/specs/binding-kafka.spec/src/test/java/io/aklivity/zilla/specs/binding/kafka/streams/application/ProduceIT.java @@ -47,6 +47,15 @@ public void shouldReconnectIdleNoErrorParallel() throws Exception k3po.finish(); } + @Test + @Specification({ + "${app}/message.value.after.fan.reply.abort/client", + "${app}/message.value.after.fan.reply.abort/server"}) + public void shouldSendMessageValueAfterFanReplyAbort() throws Exception + { + k3po.finish(); + } + @Test @Specification({ "${app}/topic.missing/client",