Skip to content
Closed
Original file line number Diff line number Diff line change
Expand Up @@ -1456,10 +1456,8 @@ private void onClientInitialFlush(
}
}

initialAck += reserved;

final int noAck = (int) (initialSeq - initialAck);
doClientInitialWindow(traceId, noAck, initialBudgetMax);
final int noAck = (int) (initialSeq - initialAck - reserved);
doClientInitialWindow(traceId, noAck, Math.max(initialMax, initialBudgetMax));
}

private void onClientInitialEnd(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1798,16 +1798,21 @@ private void doMergedInitialWindow(
initialMaxRW.value = Integer.MAX_VALUE;
produceStreams.forEach(p -> initialNoAckRW.value = Math.max(p.initialNoAck(), initialNoAckRW.value));
produceStreams.forEach(p -> initialPadRW.value = Math.max(p.initialPad, initialPadRW.value));
produceStreams.forEach(p -> initialMaxRW.value = Math.min(p.initialMax, initialMaxRW.value));
// a partition carrying the traffic holds its acknowledge and grows its maximum, while
// idle partitions keep the maximum they started with, so minimizing the maximum here
// would pair the busiest partition's unacknowledged bytes with an idle partition's
// maximum and shrink this window to nothing; minimize the budget each one offers
produceStreams.forEach(p -> initialMaxRW.value =
Math.min(p.initialMax - p.initialNoAck(), initialMaxRW.value));

maxInitialNoAck = initialNoAckRW.value;
maxInitialPad = initialPadRW.value;
minInitialMax = initialMaxRW.value + maxInitialNoAck;

if (producer != null)
{
initialMaxRW.value = Math.max(producer.initialMax, initialMaxRW.value);
minInitialMax = Math.max(producer.initialMax, minInitialMax);
}

maxInitialNoAck = initialNoAckRW.value;
maxInitialPad = initialPadRW.value;
minInitialMax = initialMaxRW.value;
}

final long newInitialAck = Math.max(initialSeq - maxInitialNoAck, initialAck);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -501,6 +501,16 @@ public void shouldProduceMergedMessageValuesDynamicHashKey() throws Exception
k3po.finish();
}

@Test
@Configuration("cache.options.merged.yaml")
@Specification({
"${app}/merged.produce.message.values.partition.idle/client",
"${app}/unmerged.produce.message.values.partition.idle/server"})
public void shouldProduceMergedMessageValuesPartitionIdle() throws Exception
{
k3po.finish();
}

@Test
@Configuration("cache.options.merged.yaml")
@Specification({
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -828,9 +828,18 @@ private void doMqttWindow(
int padding,
int capabilities)
{
final boolean retainedFlag = hasPublishFlagRetained(publishFlags);
final long newInitialAck = retainedFlag ? Math.min(messages.initialAck, retained.initialAck) : messages.initialAck;
final int newInitialMax = retainedFlag ? Math.max(messages.initialMax, retained.initialMax) : messages.initialMax;
// messages and retained advertise credit differently - messages holds its ack at zero
// and grows its max, retained advances its ack against a fixed max - so minimizing ack
// and max independently would combine the ack of one with the max of the other and cap
// the stream for its whole lifetime; compare the budgets each one actually offers
final int messagesBudget = messages.initialMax - (int)(messages.initialSeq - messages.initialAck);
final int retainedBudget = retained.initialMax - (int)(retained.initialSeq - retained.initialAck);
final int newInitialBudget = retainAvailable ?
Math.min(messagesBudget, retainedBudget) : messagesBudget;
final int newInitialMax = retainAvailable ?
Math.min(messages.initialMax, retained.initialMax) : messages.initialMax;
final long newInitialAck =
Math.max(initialAck, initialSeq - Math.max(newInitialMax - newInitialBudget, 0));

if (MqttKafkaState.initialOpened(messages.state) &&
(!retainAvailable || MqttKafkaState.initialOpened(retained.state)) &&
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1707,7 +1707,9 @@ private void onKafkaFlush(
flushEx != null && flushEx.typeId() == kafkaTypeId ? extension.get(kafkaFlushExRO::tryWrap) : null;
final KafkaMergedFlushExFW kafkaMergedFlushEx =
kafkaFlushEx != null && kafkaFlushEx.kind() == KafkaDataExFW.KIND_MERGED ? kafkaFlushEx.merged() : null;
final Array32FW<KafkaOffsetFW> progress = kafkaMergedFlushEx != null ? kafkaMergedFlushEx.fetch().progress() : null;
final Array32FW<KafkaOffsetFW> progress =
kafkaMergedFlushEx != null && kafkaMergedFlushEx.kind() == KafkaMergedFlushExFW.KIND_FETCH ?
kafkaMergedFlushEx.fetch().progress() : null;

if (progress != null)
{
Expand Down Expand Up @@ -3024,13 +3026,15 @@ protected final void doKafkaData(
int limit,
Flyweight extension)
{
if (kafka != null)
{
doData(kafka, originId, routedId, initialId, initialSeq, initialAck, initialMax,
traceId, authorization, budgetId, flags, reserved, buffer, offset, limit, extension);

doData(kafka, originId, routedId, initialId, initialSeq, initialAck, initialMax,
traceId, authorization, budgetId, flags, reserved, buffer, offset, limit, extension);

initialSeq += reserved;
initialSeq += reserved;

assert initialSeq - padding <= initialAck + initialMax;
assert initialSeq - padding <= initialAck + initialMax;
}
}

protected final void doKafkaData(
Expand All @@ -3042,12 +3046,15 @@ protected final void doKafkaData(
OctetsFW payload,
Flyweight extension)
{
doData(kafka, originId, routedId, initialId, initialSeq, initialAck, initialMax,
traceId, authorization, budgetId, flags, reserved, payload, extension);
if (kafka != null)
{
doData(kafka, originId, routedId, initialId, initialSeq, initialAck, initialMax,
traceId, authorization, budgetId, flags, reserved, payload, extension);

initialSeq += reserved;
initialSeq += reserved;

assert initialSeq <= initialAck + initialMax;
assert initialSeq <= initialAck + initialMax;
}
}

protected void doKafkaData(
Expand All @@ -3060,17 +3067,20 @@ protected void doKafkaData(
Flyweight payload,
Flyweight extension)
{
final DirectBufferEx buffer = payload.buffer();
final int offset = payload.offset();
final int limit = payload.limit();
final int length = limit - offset;
if (kafka != null)
{
final DirectBufferEx buffer = payload.buffer();
final int offset = payload.offset();
final int limit = payload.limit();
final int length = limit - offset;

doData(kafka, originId, routedId, initialId, initialSeq, initialAck, initialMax,
traceId, authorization, budgetId, flags, reserved, buffer, offset, length, extension);
doData(kafka, originId, routedId, initialId, initialSeq, initialAck, initialMax,
traceId, authorization, budgetId, flags, reserved, buffer, offset, length, extension);

initialSeq += reserved;
initialSeq += reserved;

assert initialSeq - padding <= initialAck + initialMax;
assert initialSeq - padding <= initialAck + initialMax;
}
}

private void doKafkaFlush(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,17 @@ public void shouldPublishRetainedMessage() throws Exception
k3po.finish();
}

@Test
@Configuration("proxy.yaml")
@Configure(name = WILL_AVAILABLE_NAME, value = "false")
@Specification({
"${mqtt}/publish.many.messages.retain.available/client",
"${kafka}/publish.many.messages.retain.available/server"})
public void shouldPublishManyMessagesRetainAvailable() throws Exception
{
k3po.finish();
}

@Test
@Configuration("proxy.yaml")
@Configure(name = WILL_AVAILABLE_NAME, value = "false")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,16 @@ public void shouldExpireSessionAfterSignalStreamRestart() throws Exception
k3po.finish();
}

@Test
@Configuration("proxy.yaml")
@Configure(name = PUBLISH_MAX_QOS_NAME, value = "1")
@Specification({
"${kafka}/session.ignore.non.fetch.signal.stream.flush/server"})
public void shouldIgnoreNonFetchSignalStreamFlush() throws Exception
{
k3po.finish();
}

@Test
@Configuration("proxy.yaml")
@Configure(name = PUBLISH_MAX_QOS_NAME, value = "1")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1574,24 +1574,29 @@ private int decodePublishPayload(
}
}

if (ready && reasonCode == SUCCESS)
{
int initialBudget = publisher.initialBudget();
int lengthMax = Math.min(available, server.decodeablePublishPayloadBytes);
int reservedMax = Math.max(publisher.initialPad, Math.min(lengthMax + publisher.initialPad, initialBudget));
final int initialBudget = publisher.initialBudget();
final int lengthMax = Math.min(available, server.decodeablePublishPayloadBytes);

final int maximum = reservedMax;
final int minimum = Math.min(maximum, Math.max(publisher.initialMin, 1024) + publisher.initialPad);
// a frame reserves its payload plus the padding, so the window must hold both; when it
// cannot hold even the padding this is negative, and reserving the padding regardless
// would advance initialSeq past the window the publish stream has been granted
final int sizeMax = Math.min(lengthMax, initialBudget - publisher.initialPad);

int valueClaimed = maximum;
final int maximum = sizeMax + publisher.initialPad;
final int minimum = Math.min(maximum, Math.max(publisher.initialMin, 1024) + publisher.initialPad);

if (canPublish && publisher.debit != null && lengthMax != 0)
{
valueClaimed = publisher.debit.claim(traceId, minimum, maximum);
}
int valueClaimed = maximum;

int sizeClaimed = valueClaimed - publisher.initialPad;
if (ready && reasonCode == SUCCESS && sizeMax >= 0 &&
canPublish && publisher.debit != null && lengthMax != 0)
{
valueClaimed = publisher.debit.claim(traceId, minimum, maximum);
}

final int sizeClaimed = valueClaimed - publisher.initialPad;

if (ready && reasonCode == SUCCESS && sizeClaimed >= 0)
{
final OctetsFW payload = payloadRO
.wrap(payloadBuffer, payloadOffset, payloadOffset + sizeClaimed);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,16 @@ public void shouldPublishMultipleMessages() throws Exception
k3po.finish();
}

@Test
@Configuration("server.yaml")
@Specification({
"${net}/publish.many.messages/client",
"${app}/publish.many.messages/server"})
public void shouldPublishManyMessages() throws Exception
{
k3po.finish();
}

@Test
@Configuration("server.yaml")
@Specification({
Expand Down
Loading
Loading