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
Original file line number Diff line number Diff line change
Expand Up @@ -3587,7 +3587,7 @@ private int onDecodeConnectWillMessage(
final int payloadSize = Math.min(payloadAvailable, session.initialBudget() - headerSize);
final OctetsFW willPayload = payloadRO.wrap(buffer, connectPayloadLimit, connectPayloadLimit + payloadSize);

willMessageBuffer.putBytes(headerSize, willPayload.buffer(), willPayload.offset(), willPayload.limit());
willMessageBuffer.putBytes(headerSize, willPayload.buffer(), willPayload.offset(), willPayload.sizeof());

final int deferred = willPayloadBytes - payloadSize;
final int dataFlags = deferred > 0 ? FLAG_INIT : FLAG_INIT | FLAG_FIN;
Expand Down Expand Up @@ -3642,12 +3642,19 @@ private int onDecodeConnectWillPayload(
int offset,
int limit)
{
final OctetsFW payload = payloadRO.wrap(buffer, offset, limit);
assert willPayloadDeferred >= 0;
final int flags = willPayloadDeferred - payload.sizeof() > 0 ? FLAG_CONT : FLAG_FIN;

// bounded by whatever is currently available in buffer (network-limited), the session
// window (backpressure-limited), and what is left of the declared will payload - bytes
// buffered beyond that boundary belong to the next control packet, not to the will
final int payloadAvailable = Math.min(limit - offset, willPayloadDeferred);
final int payloadSize = Math.min(payloadAvailable, session.initialBudget());
final OctetsFW willPayload = payloadRO.wrap(buffer, offset, offset + payloadSize);

final int flags = willPayloadDeferred - payloadSize > 0 ? FLAG_CONT : FLAG_FIN;

final int publishedWillSize = session.doSessionData(traceId, flags,
payload.buffer(), offset, limit, 0, EMPTY_OCTETS);
willPayload.buffer(), willPayload.offset(), willPayload.limit(), 0, EMPTY_OCTETS);
willPayloadDeferred -= publishedWillSize;
willPayloadBytes -= publishedWillSize;

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

@Test
@Configuration("server.yaml")
@Specification({
"${net}/session.will.message.disconnect.while.deferred/client",
"${app}/session.will.message.disconnect.while.deferred/server"})
@Configure(name = ENGINE_BUFFER_SLOT_CAPACITY_NAME, value = "16384")
public void shouldDisconnectWhileDeferred() throws Exception
{
k3po.finish();
}

@Test
@Configuration("server.yaml")
@Specification({
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
#
# 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.
#

connect "zilla://streams/app0"
option zilla:window 8192
option zilla:transmission "duplex"

write zilla:begin.ext ${mqtt:beginEx()
.typeId(zilla:id("mqtt"))
.session()
.flags("WILL", "CLEAN_START")
.clientId("one")
.build()
.build()}

read zilla:begin.ext ${mqtt:matchBeginEx()
.typeId(zilla:id("mqtt"))
.session()
.flags("WILL", "CLEAN_START")
.subscribeQosMax(2)
.publishQosMax(2)
.packetSizeMax(33792)
.capabilities("RETAIN", "WILDCARD", "SUBSCRIPTION_IDS", "SHARED_SUBSCRIPTIONS")
.clientId("one")
.build()
.build()}

connected

write option zilla:flags "init"
write zilla:data.ext ${mqtt:dataEx()
.typeId(zilla:id("mqtt"))
.session()
.deferred(2129)
.kind("WILL")
.build()
.build()}
write ${mqtt:will()
.topic("uagv/v2/kaercher/1.533-002.0-020156/connection")
.flags("RETAIN")
.payloadSize(10240)
.build()}
${mqtt:randomBytes(8111)}
write flush

write option zilla:flags "fin"
write ${mqtt:seededBytes(2129)}
write flush

write close
read closed
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
#
# 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.
#

accept "zilla://streams/app0"
option zilla:window 0
option zilla:transmission "duplex"

accepted

read zilla:begin.ext ${mqtt:matchBeginEx()
.typeId(zilla:id("mqtt"))
.session()
.flags("WILL", "CLEAN_START")
.clientId("one")
.build()
.build()}

write zilla:begin.ext ${mqtt:beginEx()
.typeId(zilla:id("mqtt"))
.session()
.flags("WILL", "CLEAN_START")
.subscribeQosMax(2)
.publishQosMax(2)
.packetSizeMax(33792)
.capabilities("RETAIN", "WILDCARD", "SUBSCRIPTION_IDS", "SHARED_SUBSCRIPTIONS")
.clientId("one")
.build()
.build()}

connected

# zero window granted at accept, so the whole CONNECT plus the pipelined
# PINGREQ behind it is still buffered when window is finally granted
read option zilla:window 8192
read zilla:data.ext ${mqtt:dataEx()
.typeId(zilla:id("mqtt"))
.session()
.deferred(2129)
.kind("WILL")
.build()
.build()}
read ${mqtt:will()
.topic("uagv/v2/kaercher/1.533-002.0-020156/connection")
.flags("RETAIN")
.payloadSize(10240)
.build()}
read [0..8111]

# exact-match the deferred remainder - the will payload ends here, and the four
# DISCONNECT bytes buffered immediately behind it must not be forwarded as part
# of it, otherwise the produced session record is longer than declared
read ${mqtt:seededBytes(2129)}

# the pipelined DISCONNECT is decoded once the will payload boundary is
# respected, so the session ends cleanly rather than being aborted
read closed
write close
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
#
# 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.
#

connect "zilla://streams/net0"
option zilla:window 65536
option zilla:transmission "duplex"
option zilla:byteorder "network"

connected

# CONNECT and DISCONNECT are pipelined into a single network write, so the will
# payload is not the last thing buffered when the deferred remainder is
# forwarded - the four DISCONNECT bytes behind it belong to the next control
# packet and must not be swallowed as will payload
write [0x10 0xc8 0x50] # CONNECT
[0x00 0x04] "MQTT" # protocol name
[0x05] # protocol version
[0x26] # flags = will retain, will flag, clean start
[0x00 0x0a] # keep alive = 10s
[0x05] # properties
[0x27] 33792 # maximum packet size = 33792
[0x00 0x03] "one" # client id
[0x00] # will properties
[0x00 0x2e] "uagv/v2/kaercher/1.533-002.0-020156/connection"
# will topic
[0x28 0x00] ${mqtt:randomBytes(8111)} # will payload, round 1
${mqtt:seededBytes(2129)} # will payload, round 2 (deferred remainder) -
# seeded so the reader can exact-match it
[0xe0 0x02] # DISCONNECT
[0x00] # normal disconnect
[0x00] # properties = none

# the pipelined DISCONNECT is decoded once the deferred will payload has been
# forwarded, so the connection is ended rather than left hanging
read closed
write close
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
#
# 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.
#

accept "zilla://streams/net0"
option zilla:window 65536
option zilla:transmission "duplex"
option zilla:byteorder "network"

accepted
connected

read [0x10 0xc8 0x50] # CONNECT
[0x00 0x04] "MQTT" # protocol name
[0x05] # protocol version
[0x26] # flags = will retain, will flag, clean start
[0x00 0x0a] # keep alive = 10s
[0x05] # properties
[0x27] 33792 # maximum packet size = 33792
[0x00 0x03] "one" # client id
[0x00] # will properties
[0x00 0x2e] "uagv/v2/kaercher/1.533-002.0-020156/connection"
# will topic
[0x28 0x00] [0..10240] # will payload

read [0xe0 0x02] # DISCONNECT
[0x00] # normal disconnect
[0x00] # properties = none

write close
read closed
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,15 @@ public void shouldAbortWhileDeferred() throws Exception
k3po.finish();
}

@Test
@Specification({
"${app}/session.will.message.disconnect.while.deferred/client",
"${app}/session.will.message.disconnect.while.deferred/server"})
public void shouldDisconnectWhileDeferred() throws Exception
{
k3po.finish();
}

@Test
@Specification({
"${app}/session.will.message.zero.window.on.connect/client",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,15 @@ public void shouldAbortWhileDeferred() throws Exception
k3po.finish();
}

@Test
@Specification({
"${net}/session.will.message.disconnect.while.deferred/client",
"${net}/session.will.message.disconnect.while.deferred/server"})
public void shouldDisconnectWhileDeferred() throws Exception
{
k3po.finish();
}

@Test
@Specification({
"${net}/session.will.message.zero.window.on.connect/client",
Expand Down
Loading