From 7a1e8dad07d65a6cc4e60a620bfda3c0b0ef4297 Mon Sep 17 00:00:00 2001 From: evgeny Date: Wed, 23 Sep 2026 10:21:17 +0100 Subject: [PATCH] fix: decode a batch of delta messages as a chain A protocol message can carry several messages, each a delta from the one before it, but the decoding context only advanced once per protocol message, so any batch of two or more failed on the second message with "previous message not available" and the whole batch was dropped. - Advance the context's last message id as each message decodes, so deltas chain within a batch as well as across batches, as ably-js does - Restore the context if any message in a batch fails, so the batch the server replays decodes against the right base payload rather than against the output of its own first messages - Skip the delta check when there is no decoding context, as the REST paths pass, instead of raising AttributeError Tests cover chaining within and across batches, out-of-order rejection, the context being left untouched after both failure modes, the context-less paths, and a channel receiving a three-message delta batch. Refs #692 Co-Authored-By: Claude Opus 5 (1M context) --- ably/realtime/channel.py | 1 - ably/types/message.py | 7 +- ably/types/mixins.py | 14 +- pyproject.toml | 4 +- test/ably/realtime/deltadecoding_test.py | 179 +++++++++++++++++++++++ uv.lock | 10 +- 6 files changed, 205 insertions(+), 10 deletions(-) create mode 100644 test/ably/realtime/deltadecoding_test.py diff --git a/ably/realtime/channel.py b/ably/realtime/channel.py index d5aa1a12..24f21d07 100644 --- a/ably/realtime/channel.py +++ b/ably/realtime/channel.py @@ -739,7 +739,6 @@ def _on_message(self, proto_msg: dict) -> None: try: messages = Message.from_encoded_array(proto_msg.get('messages'), cipher=self.cipher, context=self.__decoding_context) - self.__decoding_context.last_message_id = messages[-1].id self.__channel_serial = channel_serial except AblyException as e: if e.code == 40018: # Delta decode failure - start recovery diff --git a/ably/types/message.py b/ably/types/message.py index 2442a587..387b1ec2 100644 --- a/ably/types/message.py +++ b/ably/types/message.py @@ -300,12 +300,17 @@ def from_encoded(obj, cipher=None, context=None): version = obj.get('version', None) delta_extra = DeltaExtras(extras) - if delta_extra.from_id and delta_extra.from_id != context.last_message_id: + if context and delta_extra.from_id and delta_extra.from_id != context.last_message_id: raise AblyException(f"Delta message decode failure - previous message not available. " f"Message id = {id}", 400, 40018) decoded_data = Message.decode(data, encoding, cipher, context) + # A protocol message can carry several messages, each a delta from the one before it, so + # the base for the next delta is this message rather than the last one of the batch. + if context: + context.last_message_id = id + if action is not None: try: action = MessageAction(action) diff --git a/ably/types/mixins.py b/ably/types/mixins.py index 2d2b6041..48808d77 100644 --- a/ably/types/mixins.py +++ b/ably/types/mixins.py @@ -127,4 +127,16 @@ def decode(data, encoding='', cipher=None, context=None): @classmethod def from_encoded_array(cls, objs, cipher=None, context=None): - return [cls.from_encoded(obj, cipher=cipher, context=context) for obj in objs] + if context is None: + return [cls.from_encoded(obj, cipher=cipher) for obj in objs] + + # Decoding a message advances the context, and a batch that fails part way through is + # discarded whole and replayed by the server, so the context is only kept if every + # message in it decoded. Otherwise the replayed batch would decode against the base + # payload left behind by its own first messages. + base_payload, last_message_id = context.base_payload, context.last_message_id + try: + return [cls.from_encoded(obj, cipher=cipher, context=context) for obj in objs] + except Exception: + context.base_payload, context.last_message_id = base_payload, last_message_id + raise diff --git a/pyproject.toml b/pyproject.toml index ea49bc8f..de31bd5f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -40,7 +40,7 @@ dependencies = [ [project.optional-dependencies] oldcrypto = ["pycrypto>=2.6.1,<3.0.0"] crypto = ["pycryptodome"] -vcdiff = ["vcdiff-decoder>=0.1.0,<0.2.0"] +vcdiff = ["vcdiff-decoder>=0.2.0,<1.0.0"] dev = [ "pytest>=7.1,<8.0", "pytest-asyncio>=0.21.0,<0.23.0; python_version=='3.7'", @@ -55,7 +55,7 @@ dev = [ "pytest-timeout>=2.1.0,<3.0.0", "async-case>=10.1.0,<11.0.0; python_version=='3.7'", "tokenize_rt", - "vcdiff-decoder>=0.1.0a1", + "vcdiff-decoder>=0.2.0", ] [project.scripts] diff --git a/test/ably/realtime/deltadecoding_test.py b/test/ably/realtime/deltadecoding_test.py new file mode 100644 index 00000000..f2930dfb --- /dev/null +++ b/test/ably/realtime/deltadecoding_test.py @@ -0,0 +1,179 @@ +""" +Unit tests for delta (vcdiff) message decoding. + +A protocol message can carry several messages, each one a delta from the message before it, so the +decoding context is chained per message and a batch that fails part way through leaves it untouched. +""" + +import base64 + +import pytest + +from ably import AblyRealtime +from ably.types.channelstate import ChannelState +from ably.types.message import Message +from ably.types.mixins import DecodingContext +from ably.types.options import VCDiffDecoder +from ably.util.exceptions import AblyException +from test.ably.utils import BaseAsyncTestCase + + +class AppendingDecoder(VCDiffDecoder): + """A decoder whose output depends on the base, so a wrong base is visible in the result""" + + def __init__(self): + self.number_of_calls = 0 + + def decode(self, delta: bytes, base: bytes) -> bytes: + self.number_of_calls += 1 + return base + delta + + +class FailingDecoder(VCDiffDecoder): + + def decode(self, delta: bytes, base: bytes) -> bytes: + raise Exception("Failed to decode delta.") + + +def full_message(id, data): + return {'id': id, 'name': id, 'data': data, 'encoding': 'utf-8'} + + +def delta_message(id, delta, from_id): + """A message carrying `delta` as a vcdiff delta from the message `from_id`""" + return { + 'id': id, + 'name': id, + 'data': base64.b64encode(delta.encode()).decode(), + 'encoding': 'utf-8/vcdiff/base64', + 'extras': {'delta': {'format': 'vcdiff', 'from': from_id}}, + } + + +class TestDeltaBatchDecoding(BaseAsyncTestCase): + """Message.from_encoded_array over a protocol message's messages""" + + def test_batch_of_chained_deltas_decodes(self): + decoder = AppendingDecoder() + context = DecodingContext(vcdiff_decoder=decoder) + + messages = Message.from_encoded_array([ + full_message('m:0', 'a'), + delta_message('m:1', 'b', from_id='m:0'), + delta_message('m:2', 'c', from_id='m:1'), + ], context=context) + + assert [message.data for message in messages] == ['a', 'ab', 'abc'] + assert decoder.number_of_calls == 2 + assert context.last_message_id == 'm:2' + + def test_delta_chains_from_the_previous_batch(self): + context = DecodingContext(vcdiff_decoder=AppendingDecoder()) + + Message.from_encoded_array([full_message('m:0', 'a')], context=context) + messages = Message.from_encoded_array([delta_message('m:1', 'b', from_id='m:0')], + context=context) + + assert [message.data for message in messages] == ['ab'] + + def test_delta_from_an_unknown_message_is_rejected(self): + context = DecodingContext(vcdiff_decoder=AppendingDecoder()) + Message.from_encoded_array([full_message('m:0', 'a')], context=context) + + with pytest.raises(AblyException) as excinfo: + Message.from_encoded_array([delta_message('m:9', 'b', from_id='m:8')], context=context) + + assert excinfo.value.code == 40018 + + def test_a_batch_that_fails_leaves_the_context_unchanged(self): + # RTL18: the batch is discarded and replayed by the server from the last channel serial, + # so the base payload has to still be the one the replayed batch decodes against + context = DecodingContext(vcdiff_decoder=AppendingDecoder()) + Message.from_encoded_array([full_message('m:0', 'a')], context=context) + base_payload, last_message_id = context.base_payload, context.last_message_id + + with pytest.raises(AblyException): + Message.from_encoded_array([ + delta_message('m:1', 'b', from_id='m:0'), + delta_message('m:2', 'c', from_id='out-of-order'), + ], context=context) + + assert context.base_payload == base_payload + assert context.last_message_id == last_message_id + + def test_a_batch_whose_decoder_fails_leaves_the_context_unchanged(self): + context = DecodingContext(vcdiff_decoder=FailingDecoder()) + Message.from_encoded_array([full_message('m:0', 'a')], context=context) + base_payload, last_message_id = context.base_payload, context.last_message_id + + with pytest.raises(AblyException) as excinfo: + Message.from_encoded_array([delta_message('m:1', 'b', from_id='m:0')], context=context) + + assert excinfo.value.code == 40018 + assert context.base_payload == base_payload + assert context.last_message_id == last_message_id + + def test_messages_decode_without_a_decoding_context(self): + # the REST paths (history, for one) decode without a context + messages = Message.from_encoded_array([full_message('m:0', 'a')]) + + assert [message.data for message in messages] == ['a'] + + def test_a_delta_without_a_decoding_context_reports_the_missing_decoder(self): + # and not a failure to look up the previous message id in a context that isn't there + with pytest.raises(AblyException) as excinfo: + Message.from_encoded_array([delta_message('m:1', 'b', from_id='m:0')]) + + assert excinfo.value.code == 40019 + + +class TestChannelDeltaBatch(BaseAsyncTestCase): + """RealtimeChannel handling of a protocol message carrying several deltas""" + + def setup_channel(self, decoder): + ably = AblyRealtime(key='not_a.real:key', auto_connect=False, vcdiff_decoder=decoder) + channel = ably.channels.get('delta') + received = [] + # subscribe() would also wait for the channel to attach, which this client cannot do + channel._RealtimeChannel__message_emitter.on(lambda message: received.append(message)) + return ably, channel, received + + async def test_every_message_of_a_delta_batch_is_emitted(self): + ably, channel, received = self.setup_channel(AppendingDecoder()) + + channel._on_message({ + 'action': 15, + 'channel': 'delta', + 'channelSerial': 'serial-1', + 'messages': [ + full_message('m:0', 'a'), + delta_message('m:1', 'b', from_id='m:0'), + delta_message('m:2', 'c', from_id='m:1'), + ], + }) + + assert [message.data for message in received] == ['a', 'ab', 'abc'] + assert channel._RealtimeChannel__channel_serial == 'serial-1' + await ably.close() + + async def test_a_batch_that_fails_to_decode_starts_recovery(self): + ably, channel, received = self.setup_channel(FailingDecoder()) + channel._on_message({ + 'action': 15, 'channel': 'delta', 'channelSerial': 'serial-1', + 'messages': [full_message('m:0', 'a')], + }) + + channel._on_message({ + 'action': 15, + 'channel': 'delta', + 'channelSerial': 'serial-2', + 'messages': [delta_message('m:1', 'b', from_id='m:0')], + }) + + # RTL18b/RTL18c: the batch is dropped and the channel reattaches from the serial of the + # last batch it decoded + assert [message.data for message in received] == ['a'] + assert channel._RealtimeChannel__channel_serial == 'serial-1' + assert channel.state == ChannelState.ATTACHING + assert channel.error_reason.code == 40018 + await ably.close() diff --git a/uv.lock b/uv.lock index 2c5a0857..f84cb3c1 100644 --- a/uv.lock +++ b/uv.lock @@ -79,8 +79,8 @@ requires-dist = [ { name = "respx", marker = "python_full_version >= '3.8' and extra == 'dev'", specifier = ">=0.22.0,<0.23.0" }, { name = "ruff", marker = "extra == 'dev'", specifier = ">=0.14.0,<1.0.0" }, { name = "tokenize-rt", marker = "extra == 'dev'" }, - { name = "vcdiff-decoder", marker = "extra == 'dev'", specifier = ">=0.1.0a1" }, - { name = "vcdiff-decoder", marker = "extra == 'vcdiff'", specifier = ">=0.1.0,<0.2.0" }, + { name = "vcdiff-decoder", marker = "extra == 'dev'", specifier = ">=0.2.0" }, + { name = "vcdiff-decoder", marker = "extra == 'vcdiff'", specifier = ">=0.2.0,<0.3.0" }, { name = "websockets", marker = "python_full_version == '3.7.*'", specifier = ">=10.0,<12.0" }, { name = "websockets", marker = "python_full_version == '3.8.*'", specifier = ">=12.0,<15.0" }, { name = "websockets", marker = "python_full_version >= '3.9'", specifier = ">=15.0,<16.0" }, @@ -1553,11 +1553,11 @@ wheels = [ [[package]] name = "vcdiff-decoder" -version = "0.1.0" +version = "0.2.0" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/4c/17/fb4e840967b9e734e45fef6b61280bac49aa40da675a031958010707c31b/vcdiff_decoder-0.1.0.tar.gz", hash = "sha256:905d9c39fd451331301652c16b19505c16d323446fa4dffa745b2855aff5fe69", size = 18613, upload-time = "2025-09-19T17:15:14.909Z" } +sdist = { url = "https://files.pythonhosted.org/packages/15/bd/1cb5f75efe8b88f8ce6f5bc7a11d1d0da5821c8ec7aafc2736732c701a4f/vcdiff_decoder-0.2.0.tar.gz", hash = "sha256:c89a76fd46d879e5a4195187c6fc9be09827251004726a0840dd9ef8f59f5e74", size = 18248, upload-time = "2026-09-23T13:05:58.438Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/10/d5/0d1153f2dbaa02a11b2491d26b59f08e203409dacb91853c26c13bc28cb6/vcdiff_decoder-0.1.0-py3-none-any.whl", hash = "sha256:42f4e3d77b3bd4be881853858ee471a11d6a474fda375482d589b8576b91318f", size = 26333, upload-time = "2025-09-19T17:15:13.611Z" }, + { url = "https://files.pythonhosted.org/packages/6f/81/5aace343f579809a50e152eac66458ad159c4fae9e06bef438ddb0ebaf92/vcdiff_decoder-0.2.0-py3-none-any.whl", hash = "sha256:0be7163871d0c79fb57b0c9f0a994ee151c4e348a442d3d508da6975983b91aa", size = 19277, upload-time = "2026-09-23T13:05:57.346Z" }, ] [[package]]