Skip to content

Commit cc4e82b

Browse files
ttypicclaude
andcommitted
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) <noreply@anthropic.com>
1 parent 5158274 commit cc4e82b

6 files changed

Lines changed: 205 additions & 10 deletions

File tree

‎ably/realtime/channel.py‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -739,7 +739,6 @@ def _on_message(self, proto_msg: dict) -> None:
739739
try:
740740
messages = Message.from_encoded_array(proto_msg.get('messages'),
741741
cipher=self.cipher, context=self.__decoding_context)
742-
self.__decoding_context.last_message_id = messages[-1].id
743742
self.__channel_serial = channel_serial
744743
except AblyException as e:
745744
if e.code == 40018: # Delta decode failure - start recovery

‎ably/types/message.py‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -300,12 +300,17 @@ def from_encoded(obj, cipher=None, context=None):
300300
version = obj.get('version', None)
301301

302302
delta_extra = DeltaExtras(extras)
303-
if delta_extra.from_id and delta_extra.from_id != context.last_message_id:
303+
if context and delta_extra.from_id and delta_extra.from_id != context.last_message_id:
304304
raise AblyException(f"Delta message decode failure - previous message not available. "
305305
f"Message id = {id}", 400, 40018)
306306

307307
decoded_data = Message.decode(data, encoding, cipher, context)
308308

309+
# A protocol message can carry several messages, each a delta from the one before it, so
310+
# the base for the next delta is this message rather than the last one of the batch.
311+
if context:
312+
context.last_message_id = id
313+
309314
if action is not None:
310315
try:
311316
action = MessageAction(action)

‎ably/types/mixins.py‎

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -127,4 +127,16 @@ def decode(data, encoding='', cipher=None, context=None):
127127

128128
@classmethod
129129
def from_encoded_array(cls, objs, cipher=None, context=None):
130-
return [cls.from_encoded(obj, cipher=cipher, context=context) for obj in objs]
130+
if context is None:
131+
return [cls.from_encoded(obj, cipher=cipher) for obj in objs]
132+
133+
# Decoding a message advances the context, and a batch that fails part way through is
134+
# discarded whole and replayed by the server, so the context is only kept if every
135+
# message in it decoded. Otherwise the replayed batch would decode against the base
136+
# payload left behind by its own first messages.
137+
base_payload, last_message_id = context.base_payload, context.last_message_id
138+
try:
139+
return [cls.from_encoded(obj, cipher=cipher, context=context) for obj in objs]
140+
except Exception:
141+
context.base_payload, context.last_message_id = base_payload, last_message_id
142+
raise

‎pyproject.toml‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ dependencies = [
4040
[project.optional-dependencies]
4141
oldcrypto = ["pycrypto>=2.6.1,<3.0.0"]
4242
crypto = ["pycryptodome"]
43-
vcdiff = ["vcdiff-decoder>=0.1.0,<0.2.0"]
43+
vcdiff = ["vcdiff-decoder>=0.2.0,<0.3.0"]
4444
dev = [
4545
"pytest>=7.1,<8.0",
4646
"pytest-asyncio>=0.21.0,<0.23.0; python_version=='3.7'",
@@ -55,7 +55,7 @@ dev = [
5555
"pytest-timeout>=2.1.0,<3.0.0",
5656
"async-case>=10.1.0,<11.0.0; python_version=='3.7'",
5757
"tokenize_rt",
58-
"vcdiff-decoder>=0.1.0a1",
58+
"vcdiff-decoder>=0.2.0",
5959
]
6060

6161
[project.scripts]
Lines changed: 179 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,179 @@
1+
"""
2+
Unit tests for delta (vcdiff) message decoding.
3+
4+
A protocol message can carry several messages, each one a delta from the message before it, so the
5+
decoding context is chained per message and a batch that fails part way through leaves it untouched.
6+
"""
7+
8+
import base64
9+
10+
import pytest
11+
12+
from ably import AblyRealtime
13+
from ably.types.channelstate import ChannelState
14+
from ably.types.message import Message
15+
from ably.types.mixins import DecodingContext
16+
from ably.types.options import VCDiffDecoder
17+
from ably.util.exceptions import AblyException
18+
from test.ably.utils import BaseAsyncTestCase
19+
20+
21+
class AppendingDecoder(VCDiffDecoder):
22+
"""A decoder whose output depends on the base, so a wrong base is visible in the result"""
23+
24+
def __init__(self):
25+
self.number_of_calls = 0
26+
27+
def decode(self, delta: bytes, base: bytes) -> bytes:
28+
self.number_of_calls += 1
29+
return base + delta
30+
31+
32+
class FailingDecoder(VCDiffDecoder):
33+
34+
def decode(self, delta: bytes, base: bytes) -> bytes:
35+
raise Exception("Failed to decode delta.")
36+
37+
38+
def full_message(id, data):
39+
return {'id': id, 'name': id, 'data': data, 'encoding': 'utf-8'}
40+
41+
42+
def delta_message(id, delta, from_id):
43+
"""A message carrying `delta` as a vcdiff delta from the message `from_id`"""
44+
return {
45+
'id': id,
46+
'name': id,
47+
'data': base64.b64encode(delta.encode()).decode(),
48+
'encoding': 'utf-8/vcdiff/base64',
49+
'extras': {'delta': {'format': 'vcdiff', 'from': from_id}},
50+
}
51+
52+
53+
class TestDeltaBatchDecoding(BaseAsyncTestCase):
54+
"""Message.from_encoded_array over a protocol message's messages"""
55+
56+
def test_batch_of_chained_deltas_decodes(self):
57+
decoder = AppendingDecoder()
58+
context = DecodingContext(vcdiff_decoder=decoder)
59+
60+
messages = Message.from_encoded_array([
61+
full_message('m:0', 'a'),
62+
delta_message('m:1', 'b', from_id='m:0'),
63+
delta_message('m:2', 'c', from_id='m:1'),
64+
], context=context)
65+
66+
assert [message.data for message in messages] == ['a', 'ab', 'abc']
67+
assert decoder.number_of_calls == 2
68+
assert context.last_message_id == 'm:2'
69+
70+
def test_delta_chains_from_the_previous_batch(self):
71+
context = DecodingContext(vcdiff_decoder=AppendingDecoder())
72+
73+
Message.from_encoded_array([full_message('m:0', 'a')], context=context)
74+
messages = Message.from_encoded_array([delta_message('m:1', 'b', from_id='m:0')],
75+
context=context)
76+
77+
assert [message.data for message in messages] == ['ab']
78+
79+
def test_delta_from_an_unknown_message_is_rejected(self):
80+
context = DecodingContext(vcdiff_decoder=AppendingDecoder())
81+
Message.from_encoded_array([full_message('m:0', 'a')], context=context)
82+
83+
with pytest.raises(AblyException) as excinfo:
84+
Message.from_encoded_array([delta_message('m:9', 'b', from_id='m:8')], context=context)
85+
86+
assert excinfo.value.code == 40018
87+
88+
def test_a_batch_that_fails_leaves_the_context_unchanged(self):
89+
# RTL18: the batch is discarded and replayed by the server from the last channel serial,
90+
# so the base payload has to still be the one the replayed batch decodes against
91+
context = DecodingContext(vcdiff_decoder=AppendingDecoder())
92+
Message.from_encoded_array([full_message('m:0', 'a')], context=context)
93+
base_payload, last_message_id = context.base_payload, context.last_message_id
94+
95+
with pytest.raises(AblyException):
96+
Message.from_encoded_array([
97+
delta_message('m:1', 'b', from_id='m:0'),
98+
delta_message('m:2', 'c', from_id='out-of-order'),
99+
], context=context)
100+
101+
assert context.base_payload == base_payload
102+
assert context.last_message_id == last_message_id
103+
104+
def test_a_batch_whose_decoder_fails_leaves_the_context_unchanged(self):
105+
context = DecodingContext(vcdiff_decoder=FailingDecoder())
106+
Message.from_encoded_array([full_message('m:0', 'a')], context=context)
107+
base_payload, last_message_id = context.base_payload, context.last_message_id
108+
109+
with pytest.raises(AblyException) as excinfo:
110+
Message.from_encoded_array([delta_message('m:1', 'b', from_id='m:0')], context=context)
111+
112+
assert excinfo.value.code == 40018
113+
assert context.base_payload == base_payload
114+
assert context.last_message_id == last_message_id
115+
116+
def test_messages_decode_without_a_decoding_context(self):
117+
# the REST paths (history, for one) decode without a context
118+
messages = Message.from_encoded_array([full_message('m:0', 'a')])
119+
120+
assert [message.data for message in messages] == ['a']
121+
122+
def test_a_delta_without_a_decoding_context_reports_the_missing_decoder(self):
123+
# and not a failure to look up the previous message id in a context that isn't there
124+
with pytest.raises(AblyException) as excinfo:
125+
Message.from_encoded_array([delta_message('m:1', 'b', from_id='m:0')])
126+
127+
assert excinfo.value.code == 40019
128+
129+
130+
class TestChannelDeltaBatch(BaseAsyncTestCase):
131+
"""RealtimeChannel handling of a protocol message carrying several deltas"""
132+
133+
def setup_channel(self, decoder):
134+
ably = AblyRealtime(key='not_a.real:key', auto_connect=False, vcdiff_decoder=decoder)
135+
channel = ably.channels.get('delta')
136+
received = []
137+
# subscribe() would also wait for the channel to attach, which this client cannot do
138+
channel._RealtimeChannel__message_emitter.on(lambda message: received.append(message))
139+
return ably, channel, received
140+
141+
async def test_every_message_of_a_delta_batch_is_emitted(self):
142+
ably, channel, received = self.setup_channel(AppendingDecoder())
143+
144+
channel._on_message({
145+
'action': 15,
146+
'channel': 'delta',
147+
'channelSerial': 'serial-1',
148+
'messages': [
149+
full_message('m:0', 'a'),
150+
delta_message('m:1', 'b', from_id='m:0'),
151+
delta_message('m:2', 'c', from_id='m:1'),
152+
],
153+
})
154+
155+
assert [message.data for message in received] == ['a', 'ab', 'abc']
156+
assert channel._RealtimeChannel__channel_serial == 'serial-1'
157+
await ably.close()
158+
159+
async def test_a_batch_that_fails_to_decode_starts_recovery(self):
160+
ably, channel, received = self.setup_channel(FailingDecoder())
161+
channel._on_message({
162+
'action': 15, 'channel': 'delta', 'channelSerial': 'serial-1',
163+
'messages': [full_message('m:0', 'a')],
164+
})
165+
166+
channel._on_message({
167+
'action': 15,
168+
'channel': 'delta',
169+
'channelSerial': 'serial-2',
170+
'messages': [delta_message('m:1', 'b', from_id='m:0')],
171+
})
172+
173+
# RTL18b/RTL18c: the batch is dropped and the channel reattaches from the serial of the
174+
# last batch it decoded
175+
assert [message.data for message in received] == ['a']
176+
assert channel._RealtimeChannel__channel_serial == 'serial-1'
177+
assert channel.state == ChannelState.ATTACHING
178+
assert channel.error_reason.code == 40018
179+
await ably.close()

‎uv.lock‎

Lines changed: 5 additions & 5 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

0 commit comments

Comments
 (0)