Skip to content
Merged
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
4 changes: 2 additions & 2 deletions tests/test_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -522,7 +522,7 @@ async def handler(msg: dict[str, int]) -> None:
# Every later message still reached the handler: the application kept consuming after
# the signal that used to kill it. This is the regression, not merely "the loop is alive".
assert seen == n_messages
assert [m["id"] for m in processed] == [1, 2]
assert sorted(m["id"] for m in processed) == [1, 2]


async def test_real_kafka_direct_ack_from_handler_is_refused(kafka_bootstrap_servers: str) -> None:
Expand Down Expand Up @@ -565,7 +565,7 @@ async def handler(msg: dict[str, int], message: KafkaMessage) -> None:
# Message 0 was refused and never completed; 1 and 2 processed normally.
assert len(errors) == 1
assert "Do not call `message.ack()`" in errors[0]
assert [m["id"] for m in processed] == [1, 2]
assert sorted(m["id"] for m in processed) == [1, 2]


async def _committed_offsets(bootstrap_servers: str, group_id: str) -> dict[int, int]:
Expand Down
Loading