Skip to content

Commit

Permalink
Remove stream on group leave (#492)
Browse files Browse the repository at this point in the history
  • Loading branch information
akrambek authored Oct 5, 2023
1 parent b96513e commit 226bcba
Show file tree
Hide file tree
Showing 8 changed files with 926 additions and 83 deletions.

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -176,4 +176,14 @@ public void shouldIgnoreHeartbeatBeforeHandshakeComplete() throws Exception
{
k3po.finish();
}

@Test
@Configuration("client.yaml")
@Specification({
"${app}/rebalance.multiple.members.with.same.group.id/client",
"${net}/rebalance.multiple.members.with.same.group.id/server"})
public void shouldRebalanceMultipleMembersWithSameGroupId() throws Exception
{
k3po.finish();
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
#
# Copyright 2021-2023 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 "half-duplex"

write zilla:begin.ext ${kafka:beginEx()
.typeId(zilla:id("kafka"))
.group()
.groupId("test")
.protocol("highlander")
.timeout(45000)
.build()
.build()}

connected

read zilla:begin.ext ${kafka:matchBeginEx()
.typeId(zilla:id("kafka"))
.group()
.groupId("test")
.protocol("highlander")
.timeout(30000)
.build()
.build()}

read advised zilla:flush ${kafka:flushEx()
.typeId(zilla:id("kafka"))
.group()
.leaderId("memberId-1")
.memberId("memberId-1")
.members("memberId-1")
.build()
.build()}

write zilla:data.empty
write flush

read zilla:data.empty

write close
read closed

read notify ROUTED_BROKER_SERVER

connect await ROUTED_BROKER_SERVER "zilla://streams/app0"
option zilla:window 8192
option zilla:transmission "half-duplex"

write zilla:begin.ext ${kafka:beginEx()
.typeId(zilla:id("kafka"))
.group()
.groupId("test")
.protocol("highlander")
.timeout(45000)
.build()
.build()}

connected

read zilla:begin.ext ${kafka:matchBeginEx()
.typeId(zilla:id("kafka"))
.group()
.groupId("test")
.protocol("highlander")
.timeout(30000)
.build()
.build()}

read advised zilla:flush ${kafka:flushEx()
.typeId(zilla:id("kafka"))
.group()
.leaderId("memberId-2")
.memberId("memberId-2")
.members("memberId-2")
.build()
.build()}

write zilla:data.empty
write flush

read zilla:data.empty
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
#
# Copyright 2021-2023 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 "half-duplex"

accepted

read zilla:begin.ext ${kafka:matchBeginEx()
.typeId(zilla:id("kafka"))
.group()
.groupId("test")
.protocol("highlander")
.timeout(45000)
.build()
.build()}

connected

write zilla:begin.ext ${kafka:beginEx()
.typeId(zilla:id("kafka"))
.group()
.groupId("test")
.protocol("highlander")
.timeout(30000)
.build()
.build()}
write flush

write advise zilla:flush ${kafka:flushEx()
.typeId(zilla:id("kafka"))
.group()
.leaderId("memberId-1")
.memberId("memberId-1")
.members("memberId-1")
.build()
.build()}

read zilla:data.empty

write zilla:data.empty
write flush

read closed
write close

accepted

read zilla:begin.ext ${kafka:matchBeginEx()
.typeId(zilla:id("kafka"))
.group()
.groupId("test")
.protocol("highlander")
.timeout(45000)
.build()
.build()}

connected

write zilla:begin.ext ${kafka:beginEx()
.typeId(zilla:id("kafka"))
.group()
.groupId("test")
.protocol("highlander")
.timeout(30000)
.build()
.build()}
write flush

write advise zilla:flush ${kafka:flushEx()
.typeId(zilla:id("kafka"))
.group()
.leaderId("memberId-2")
.memberId("memberId-2")
.members("memberId-2")
.build()
.build()}

read zilla:data.empty

write zilla:data.empty
write flush
Loading

0 comments on commit 226bcba

Please sign in to comment.