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
1 change: 0 additions & 1 deletion ably/realtime/channel.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 6 additions & 1 deletion ably/types/message.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
14 changes: 13 additions & 1 deletion ably/types/mixins.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
4 changes: 2 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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'",
Expand All @@ -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]
Expand Down
179 changes: 179 additions & 0 deletions test/ably/realtime/deltadecoding_test.py
Original file line number Diff line number Diff line change
@@ -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()
10 changes: 5 additions & 5 deletions uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading