wireprotoframing.py
861 lines
| 28.8 KiB
| text/x-python
|
PythonLexer
/ mercurial / wireprotoframing.py
Gregory Szorc
|
r37069 | # wireprotoframing.py - unified framing protocol for wire protocol | ||
# | ||||
# Copyright 2018 Gregory Szorc <gregory.szorc@gmail.com> | ||||
# | ||||
# This software may be used and distributed according to the terms of the | ||||
# GNU General Public License version 2 or any later version. | ||||
# This file contains functionality to support the unified frame-based wire | ||||
# protocol. For details about the protocol, see | ||||
# `hg help internals.wireprotocol`. | ||||
from __future__ import absolute_import | ||||
import struct | ||||
Gregory Szorc
|
r37070 | from .i18n import _ | ||
Gregory Szorc
|
r37079 | from .thirdparty import ( | ||
attr, | ||||
Gregory Szorc
|
r37306 | cbor, | ||
Gregory Szorc
|
r37079 | ) | ||
Gregory Szorc
|
r37069 | from . import ( | ||
Gregory Szorc
|
r37070 | error, | ||
Gregory Szorc
|
r37069 | util, | ||
) | ||||
Yuya Nishihara
|
r37102 | from .utils import ( | ||
stringutil, | ||||
) | ||||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37304 | FRAME_HEADER_SIZE = 8 | ||
Gregory Szorc
|
r37069 | DEFAULT_MAX_FRAME_SIZE = 32768 | ||
Gregory Szorc
|
r37304 | STREAM_FLAG_BEGIN_STREAM = 0x01 | ||
STREAM_FLAG_END_STREAM = 0x02 | ||||
STREAM_FLAG_ENCODING_APPLIED = 0x04 | ||||
STREAM_FLAGS = { | ||||
b'stream-begin': STREAM_FLAG_BEGIN_STREAM, | ||||
b'stream-end': STREAM_FLAG_END_STREAM, | ||||
b'encoded': STREAM_FLAG_ENCODING_APPLIED, | ||||
} | ||||
Gregory Szorc
|
r37308 | FRAME_TYPE_COMMAND_REQUEST = 0x01 | ||
Gregory Szorc
|
r37069 | FRAME_TYPE_COMMAND_DATA = 0x03 | ||
Gregory Szorc
|
r37073 | FRAME_TYPE_BYTES_RESPONSE = 0x04 | ||
FRAME_TYPE_ERROR_RESPONSE = 0x05 | ||||
Gregory Szorc
|
r37078 | FRAME_TYPE_TEXT_OUTPUT = 0x06 | ||
Gregory Szorc
|
r37307 | FRAME_TYPE_PROGRESS = 0x07 | ||
Gregory Szorc
|
r37304 | FRAME_TYPE_STREAM_SETTINGS = 0x08 | ||
Gregory Szorc
|
r37069 | |||
FRAME_TYPES = { | ||||
Gregory Szorc
|
r37308 | b'command-request': FRAME_TYPE_COMMAND_REQUEST, | ||
Gregory Szorc
|
r37069 | b'command-data': FRAME_TYPE_COMMAND_DATA, | ||
Gregory Szorc
|
r37073 | b'bytes-response': FRAME_TYPE_BYTES_RESPONSE, | ||
b'error-response': FRAME_TYPE_ERROR_RESPONSE, | ||||
Gregory Szorc
|
r37078 | b'text-output': FRAME_TYPE_TEXT_OUTPUT, | ||
Gregory Szorc
|
r37307 | b'progress': FRAME_TYPE_PROGRESS, | ||
Gregory Szorc
|
r37304 | b'stream-settings': FRAME_TYPE_STREAM_SETTINGS, | ||
Gregory Szorc
|
r37069 | } | ||
Gregory Szorc
|
r37308 | FLAG_COMMAND_REQUEST_NEW = 0x01 | ||
FLAG_COMMAND_REQUEST_CONTINUATION = 0x02 | ||||
FLAG_COMMAND_REQUEST_MORE_FRAMES = 0x04 | ||||
FLAG_COMMAND_REQUEST_EXPECT_DATA = 0x08 | ||||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37308 | FLAGS_COMMAND_REQUEST = { | ||
b'new': FLAG_COMMAND_REQUEST_NEW, | ||||
b'continuation': FLAG_COMMAND_REQUEST_CONTINUATION, | ||||
b'more': FLAG_COMMAND_REQUEST_MORE_FRAMES, | ||||
b'have-data': FLAG_COMMAND_REQUEST_EXPECT_DATA, | ||||
Gregory Szorc
|
r37069 | } | ||
FLAG_COMMAND_DATA_CONTINUATION = 0x01 | ||||
FLAG_COMMAND_DATA_EOS = 0x02 | ||||
FLAGS_COMMAND_DATA = { | ||||
b'continuation': FLAG_COMMAND_DATA_CONTINUATION, | ||||
b'eos': FLAG_COMMAND_DATA_EOS, | ||||
} | ||||
Gregory Szorc
|
r37073 | FLAG_BYTES_RESPONSE_CONTINUATION = 0x01 | ||
FLAG_BYTES_RESPONSE_EOS = 0x02 | ||||
FLAGS_BYTES_RESPONSE = { | ||||
b'continuation': FLAG_BYTES_RESPONSE_CONTINUATION, | ||||
b'eos': FLAG_BYTES_RESPONSE_EOS, | ||||
} | ||||
FLAG_ERROR_RESPONSE_PROTOCOL = 0x01 | ||||
FLAG_ERROR_RESPONSE_APPLICATION = 0x02 | ||||
FLAGS_ERROR_RESPONSE = { | ||||
b'protocol': FLAG_ERROR_RESPONSE_PROTOCOL, | ||||
b'application': FLAG_ERROR_RESPONSE_APPLICATION, | ||||
} | ||||
Gregory Szorc
|
r37069 | # Maps frame types to their available flags. | ||
FRAME_TYPE_FLAGS = { | ||||
Gregory Szorc
|
r37308 | FRAME_TYPE_COMMAND_REQUEST: FLAGS_COMMAND_REQUEST, | ||
Gregory Szorc
|
r37069 | FRAME_TYPE_COMMAND_DATA: FLAGS_COMMAND_DATA, | ||
Gregory Szorc
|
r37073 | FRAME_TYPE_BYTES_RESPONSE: FLAGS_BYTES_RESPONSE, | ||
FRAME_TYPE_ERROR_RESPONSE: FLAGS_ERROR_RESPONSE, | ||||
Gregory Szorc
|
r37078 | FRAME_TYPE_TEXT_OUTPUT: {}, | ||
Gregory Szorc
|
r37307 | FRAME_TYPE_PROGRESS: {}, | ||
Gregory Szorc
|
r37304 | FRAME_TYPE_STREAM_SETTINGS: {}, | ||
Gregory Szorc
|
r37069 | } | ||
Gregory Szorc
|
r37308 | ARGUMENT_RECORD_HEADER = struct.Struct(r'<HH') | ||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37079 | @attr.s(slots=True) | ||
class frameheader(object): | ||||
"""Represents the data in a frame header.""" | ||||
length = attr.ib() | ||||
requestid = attr.ib() | ||||
Gregory Szorc
|
r37304 | streamid = attr.ib() | ||
streamflags = attr.ib() | ||||
Gregory Szorc
|
r37079 | typeid = attr.ib() | ||
flags = attr.ib() | ||||
@attr.s(slots=True) | ||||
class frame(object): | ||||
"""Represents a parsed frame.""" | ||||
requestid = attr.ib() | ||||
Gregory Szorc
|
r37304 | streamid = attr.ib() | ||
streamflags = attr.ib() | ||||
Gregory Szorc
|
r37079 | typeid = attr.ib() | ||
flags = attr.ib() | ||||
payload = attr.ib() | ||||
Gregory Szorc
|
r37304 | def makeframe(requestid, streamid, streamflags, typeid, flags, payload): | ||
Gregory Szorc
|
r37069 | """Assemble a frame into a byte array.""" | ||
# TODO assert size of payload. | ||||
frame = bytearray(FRAME_HEADER_SIZE + len(payload)) | ||||
Gregory Szorc
|
r37075 | # 24 bits length | ||
# 16 bits request id | ||||
Gregory Szorc
|
r37304 | # 8 bits stream id | ||
# 8 bits stream flags | ||||
Gregory Szorc
|
r37075 | # 4 bits type | ||
# 4 bits flags | ||||
Gregory Szorc
|
r37069 | l = struct.pack(r'<I', len(payload)) | ||
frame[0:3] = l[0:3] | ||||
Gregory Szorc
|
r37304 | struct.pack_into(r'<HBB', frame, 3, requestid, streamid, streamflags) | ||
frame[7] = (typeid << 4) | flags | ||||
frame[8:] = payload | ||||
Gregory Szorc
|
r37069 | |||
return frame | ||||
def makeframefromhumanstring(s): | ||||
Gregory Szorc
|
r37075 | """Create a frame from a human readable string | ||
Gregory Szorc
|
r37306 | DANGER: NOT SAFE TO USE WITH UNTRUSTED INPUT BECAUSE OF POTENTIAL | ||
eval() USAGE. DO NOT USE IN CORE. | ||||
Gregory Szorc
|
r37075 | Strings have the form: | ||
Gregory Szorc
|
r37304 | <request-id> <stream-id> <stream-flags> <type> <flags> <payload> | ||
Gregory Szorc
|
r37069 | |||
This can be used by user-facing applications and tests for creating | ||||
frames easily without having to type out a bunch of constants. | ||||
Gregory Szorc
|
r37304 | Request ID and stream IDs are integers. | ||
Gregory Szorc
|
r37075 | |||
Gregory Szorc
|
r37304 | Stream flags, frame type, and flags can be specified by integer or | ||
named constant. | ||||
Gregory Szorc
|
r37075 | |||
Gregory Szorc
|
r37069 | Flags can be delimited by `|` to bitwise OR them together. | ||
Gregory Szorc
|
r37306 | |||
If the payload begins with ``cbor:``, the following string will be | ||||
evaluated as Python code and the resulting object will be fed into | ||||
a CBOR encoder. Otherwise, the payload is interpreted as a Python | ||||
byte string literal. | ||||
Gregory Szorc
|
r37069 | """ | ||
Gregory Szorc
|
r37304 | fields = s.split(b' ', 5) | ||
requestid, streamid, streamflags, frametype, frameflags, payload = fields | ||||
Gregory Szorc
|
r37075 | |||
requestid = int(requestid) | ||||
Gregory Szorc
|
r37304 | streamid = int(streamid) | ||
finalstreamflags = 0 | ||||
for flag in streamflags.split(b'|'): | ||||
if flag in STREAM_FLAGS: | ||||
finalstreamflags |= STREAM_FLAGS[flag] | ||||
else: | ||||
finalstreamflags |= int(flag) | ||||
Gregory Szorc
|
r37069 | |||
if frametype in FRAME_TYPES: | ||||
frametype = FRAME_TYPES[frametype] | ||||
else: | ||||
frametype = int(frametype) | ||||
finalflags = 0 | ||||
validflags = FRAME_TYPE_FLAGS[frametype] | ||||
for flag in frameflags.split(b'|'): | ||||
if flag in validflags: | ||||
finalflags |= validflags[flag] | ||||
else: | ||||
finalflags |= int(flag) | ||||
Gregory Szorc
|
r37306 | if payload.startswith(b'cbor:'): | ||
payload = cbor.dumps(stringutil.evalpython(payload[5:]), canonical=True) | ||||
else: | ||||
payload = stringutil.unescapestr(payload) | ||||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37304 | return makeframe(requestid=requestid, streamid=streamid, | ||
streamflags=finalstreamflags, typeid=frametype, | ||||
Gregory Szorc
|
r37080 | flags=finalflags, payload=payload) | ||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37070 | def parseheader(data): | ||
"""Parse a unified framing protocol frame header from a buffer. | ||||
The header is expected to be in the buffer at offset 0 and the | ||||
buffer is expected to be large enough to hold a full header. | ||||
""" | ||||
# 24 bits payload length (little endian) | ||||
Gregory Szorc
|
r37304 | # 16 bits request ID | ||
# 8 bits stream ID | ||||
# 8 bits stream flags | ||||
Gregory Szorc
|
r37070 | # 4 bits frame type | ||
# 4 bits frame flags | ||||
# ... payload | ||||
framelength = data[0] + 256 * data[1] + 16384 * data[2] | ||||
Gregory Szorc
|
r37304 | requestid, streamid, streamflags = struct.unpack_from(r'<HBB', data, 3) | ||
typeflags = data[7] | ||||
Gregory Szorc
|
r37070 | |||
frametype = (typeflags & 0xf0) >> 4 | ||||
frameflags = typeflags & 0x0f | ||||
Gregory Szorc
|
r37304 | return frameheader(framelength, requestid, streamid, streamflags, | ||
frametype, frameflags) | ||||
Gregory Szorc
|
r37070 | |||
def readframe(fh): | ||||
"""Read a unified framing protocol frame from a file object. | ||||
Returns a 3-tuple of (type, flags, payload) for the decoded frame or | ||||
None if no frame is available. May raise if a malformed frame is | ||||
seen. | ||||
""" | ||||
header = bytearray(FRAME_HEADER_SIZE) | ||||
readcount = fh.readinto(header) | ||||
if readcount == 0: | ||||
return None | ||||
if readcount != FRAME_HEADER_SIZE: | ||||
raise error.Abort(_('received incomplete frame: got %d bytes: %s') % | ||||
(readcount, header)) | ||||
Gregory Szorc
|
r37079 | h = parseheader(header) | ||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37079 | payload = fh.read(h.length) | ||
if len(payload) != h.length: | ||||
Gregory Szorc
|
r37070 | raise error.Abort(_('frame length error: expected %d; got %d') % | ||
Gregory Szorc
|
r37079 | (h.length, len(payload))) | ||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37304 | return frame(h.requestid, h.streamid, h.streamflags, h.typeid, h.flags, | ||
payload) | ||||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37308 | def createcommandframes(stream, requestid, cmd, args, datafh=None, | ||
maxframesize=DEFAULT_MAX_FRAME_SIZE): | ||||
Gregory Szorc
|
r37069 | """Create frames necessary to transmit a request to run a command. | ||
This is a generator of bytearrays. Each item represents a frame | ||||
ready to be sent over the wire to a peer. | ||||
""" | ||||
Gregory Szorc
|
r37308 | data = {b'name': cmd} | ||
Gregory Szorc
|
r37069 | if args: | ||
Gregory Szorc
|
r37308 | data[b'args'] = args | ||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37308 | data = cbor.dumps(data, canonical=True) | ||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37308 | offset = 0 | ||
while True: | ||||
flags = 0 | ||||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37308 | # Must set new or continuation flag. | ||
if not offset: | ||||
flags |= FLAG_COMMAND_REQUEST_NEW | ||||
else: | ||||
flags |= FLAG_COMMAND_REQUEST_CONTINUATION | ||||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37308 | # Data frames is set on all frames. | ||
if datafh: | ||||
flags |= FLAG_COMMAND_REQUEST_EXPECT_DATA | ||||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37308 | payload = data[offset:offset + maxframesize] | ||
offset += len(payload) | ||||
if len(payload) == maxframesize and offset < len(data): | ||||
flags |= FLAG_COMMAND_REQUEST_MORE_FRAMES | ||||
Gregory Szorc
|
r37303 | yield stream.makeframe(requestid=requestid, | ||
Gregory Szorc
|
r37308 | typeid=FRAME_TYPE_COMMAND_REQUEST, | ||
Gregory Szorc
|
r37303 | flags=flags, | ||
payload=payload) | ||||
Gregory Szorc
|
r37069 | |||
Gregory Szorc
|
r37308 | if not (flags & FLAG_COMMAND_REQUEST_MORE_FRAMES): | ||
break | ||||
Gregory Szorc
|
r37069 | if datafh: | ||
while True: | ||||
data = datafh.read(DEFAULT_MAX_FRAME_SIZE) | ||||
done = False | ||||
if len(data) == DEFAULT_MAX_FRAME_SIZE: | ||||
flags = FLAG_COMMAND_DATA_CONTINUATION | ||||
else: | ||||
flags = FLAG_COMMAND_DATA_EOS | ||||
assert datafh.read(1) == b'' | ||||
done = True | ||||
Gregory Szorc
|
r37303 | yield stream.makeframe(requestid=requestid, | ||
typeid=FRAME_TYPE_COMMAND_DATA, | ||||
flags=flags, | ||||
payload=data) | ||||
Gregory Szorc
|
r37069 | |||
if done: | ||||
break | ||||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37303 | def createbytesresponseframesfrombytes(stream, requestid, data, | ||
Gregory Szorc
|
r37073 | maxframesize=DEFAULT_MAX_FRAME_SIZE): | ||
"""Create a raw frame to send a bytes response from static bytes input. | ||||
Returns a generator of bytearrays. | ||||
""" | ||||
# Simple case of a single frame. | ||||
if len(data) <= maxframesize: | ||||
Gregory Szorc
|
r37303 | yield stream.makeframe(requestid=requestid, | ||
typeid=FRAME_TYPE_BYTES_RESPONSE, | ||||
flags=FLAG_BYTES_RESPONSE_EOS, | ||||
payload=data) | ||||
Gregory Szorc
|
r37073 | return | ||
offset = 0 | ||||
while True: | ||||
chunk = data[offset:offset + maxframesize] | ||||
offset += len(chunk) | ||||
done = offset == len(data) | ||||
if done: | ||||
flags = FLAG_BYTES_RESPONSE_EOS | ||||
else: | ||||
flags = FLAG_BYTES_RESPONSE_CONTINUATION | ||||
Gregory Szorc
|
r37303 | yield stream.makeframe(requestid=requestid, | ||
typeid=FRAME_TYPE_BYTES_RESPONSE, | ||||
flags=flags, | ||||
payload=chunk) | ||||
Gregory Szorc
|
r37073 | |||
if done: | ||||
break | ||||
Gregory Szorc
|
r37303 | def createerrorframe(stream, requestid, msg, protocol=False, application=False): | ||
Gregory Szorc
|
r37073 | # TODO properly handle frame size limits. | ||
assert len(msg) <= DEFAULT_MAX_FRAME_SIZE | ||||
flags = 0 | ||||
if protocol: | ||||
flags |= FLAG_ERROR_RESPONSE_PROTOCOL | ||||
if application: | ||||
flags |= FLAG_ERROR_RESPONSE_APPLICATION | ||||
Gregory Szorc
|
r37303 | yield stream.makeframe(requestid=requestid, | ||
typeid=FRAME_TYPE_ERROR_RESPONSE, | ||||
flags=flags, | ||||
payload=msg) | ||||
Gregory Szorc
|
r37073 | |||
Gregory Szorc
|
r37303 | def createtextoutputframe(stream, requestid, atoms): | ||
Gregory Szorc
|
r37078 | """Create a text output frame to render text to people. | ||
``atoms`` is a 3-tuple of (formatting string, args, labels). | ||||
The formatting string contains ``%s`` tokens to be replaced by the | ||||
corresponding indexed entry in ``args``. ``labels`` is an iterable of | ||||
formatters to be applied at rendering time. In terms of the ``ui`` | ||||
class, each atom corresponds to a ``ui.write()``. | ||||
""" | ||||
bytesleft = DEFAULT_MAX_FRAME_SIZE | ||||
atomchunks = [] | ||||
for (formatting, args, labels) in atoms: | ||||
if len(args) > 255: | ||||
raise ValueError('cannot use more than 255 formatting arguments') | ||||
if len(labels) > 255: | ||||
raise ValueError('cannot use more than 255 labels') | ||||
# TODO look for localstr, other types here? | ||||
if not isinstance(formatting, bytes): | ||||
raise ValueError('must use bytes formatting strings') | ||||
for arg in args: | ||||
if not isinstance(arg, bytes): | ||||
raise ValueError('must use bytes for arguments') | ||||
for label in labels: | ||||
if not isinstance(label, bytes): | ||||
raise ValueError('must use bytes for labels') | ||||
# Formatting string must be UTF-8. | ||||
formatting = formatting.decode(r'utf-8', r'replace').encode(r'utf-8') | ||||
# Arguments must be UTF-8. | ||||
args = [a.decode(r'utf-8', r'replace').encode(r'utf-8') for a in args] | ||||
# Labels must be ASCII. | ||||
labels = [l.decode(r'ascii', r'strict').encode(r'ascii') | ||||
for l in labels] | ||||
if len(formatting) > 65535: | ||||
raise ValueError('formatting string cannot be longer than 64k') | ||||
if any(len(a) > 65535 for a in args): | ||||
raise ValueError('argument string cannot be longer than 64k') | ||||
if any(len(l) > 255 for l in labels): | ||||
raise ValueError('label string cannot be longer than 255 bytes') | ||||
chunks = [ | ||||
struct.pack(r'<H', len(formatting)), | ||||
struct.pack(r'<BB', len(labels), len(args)), | ||||
struct.pack(r'<' + r'B' * len(labels), *map(len, labels)), | ||||
struct.pack(r'<' + r'H' * len(args), *map(len, args)), | ||||
] | ||||
chunks.append(formatting) | ||||
chunks.extend(labels) | ||||
chunks.extend(args) | ||||
atom = b''.join(chunks) | ||||
atomchunks.append(atom) | ||||
bytesleft -= len(atom) | ||||
if bytesleft < 0: | ||||
raise ValueError('cannot encode data in a single frame') | ||||
Gregory Szorc
|
r37303 | yield stream.makeframe(requestid=requestid, | ||
typeid=FRAME_TYPE_TEXT_OUTPUT, | ||||
flags=0, | ||||
payload=b''.join(atomchunks)) | ||||
class stream(object): | ||||
"""Represents a logical unidirectional series of frames.""" | ||||
Gregory Szorc
|
r37304 | def __init__(self, streamid, active=False): | ||
self.streamid = streamid | ||||
self._active = False | ||||
Gregory Szorc
|
r37303 | def makeframe(self, requestid, typeid, flags, payload): | ||
"""Create a frame to be sent out over this stream. | ||||
Only returns the frame instance. Does not actually send it. | ||||
""" | ||||
Gregory Szorc
|
r37304 | streamflags = 0 | ||
if not self._active: | ||||
streamflags |= STREAM_FLAG_BEGIN_STREAM | ||||
self._active = True | ||||
return makeframe(requestid, self.streamid, streamflags, typeid, flags, | ||||
payload) | ||||
def ensureserverstream(stream): | ||||
if stream.streamid % 2: | ||||
raise error.ProgrammingError('server should only write to even ' | ||||
'numbered streams; %d is not even' % | ||||
stream.streamid) | ||||
Gregory Szorc
|
r37078 | |||
Gregory Szorc
|
r37070 | class serverreactor(object): | ||
"""Holds state of a server handling frame-based protocol requests. | ||||
This class is the "brain" of the unified frame-based protocol server | ||||
component. While the protocol is stateless from the perspective of | ||||
requests/commands, something needs to track which frames have been | ||||
received, what frames to expect, etc. This class is that thing. | ||||
Instances are modeled as a state machine of sorts. Instances are also | ||||
reactionary to external events. The point of this class is to encapsulate | ||||
the state of the connection and the exchange of frames, not to perform | ||||
work. Instead, callers tell this class when something occurs, like a | ||||
frame arriving. If that activity is worthy of a follow-up action (say | ||||
*run a command*), the return value of that handler will say so. | ||||
I/O and CPU intensive operations are purposefully delegated outside of | ||||
this class. | ||||
Consumers are expected to tell instances when events occur. They do so by | ||||
calling the various ``on*`` methods. These methods return a 2-tuple | ||||
describing any follow-up action(s) to take. The first element is the | ||||
name of an action to perform. The second is a data structure (usually | ||||
a dict) specific to that action that contains more information. e.g. | ||||
if the server wants to send frames back to the client, the data structure | ||||
will contain a reference to those frames. | ||||
Valid actions that consumers can be instructed to take are: | ||||
Gregory Szorc
|
r37073 | sendframes | ||
Indicates that frames should be sent to the client. The ``framegen`` | ||||
key contains a generator of frames that should be sent. The server | ||||
assumes that all frames are sent to the client. | ||||
Gregory Szorc
|
r37070 | error | ||
Indicates that an error occurred. Consumer should probably abort. | ||||
runcommand | ||||
Indicates that the consumer should run a wire protocol command. Details | ||||
of the command to run are given in the data structure. | ||||
wantframe | ||||
Indicates that nothing of interest happened and the server is waiting on | ||||
more frames from the client before anything interesting can be done. | ||||
Gregory Szorc
|
r37074 | |||
noop | ||||
Indicates no additional action is required. | ||||
Gregory Szorc
|
r37076 | |||
Known Issues | ||||
------------ | ||||
There are no limits to the number of partially received commands or their | ||||
size. A malicious client could stream command request data and exhaust the | ||||
server's memory. | ||||
Partially received commands are not acted upon when end of input is | ||||
reached. Should the server error if it receives a partial request? | ||||
Should the client send a message to abort a partially transmitted request | ||||
to facilitate graceful shutdown? | ||||
Active requests that haven't been responded to aren't tracked. This means | ||||
that if we receive a command and instruct its dispatch, another command | ||||
with its request ID can come in over the wire and there will be a race | ||||
between who responds to what. | ||||
Gregory Szorc
|
r37070 | """ | ||
Gregory Szorc
|
r37074 | def __init__(self, deferoutput=False): | ||
"""Construct a new server reactor. | ||||
``deferoutput`` can be used to indicate that no output frames should be | ||||
instructed to be sent until input has been exhausted. In this mode, | ||||
events that would normally generate output frames (such as a command | ||||
response being ready) will instead defer instructing the consumer to | ||||
send those frames. This is useful for half-duplex transports where the | ||||
sender cannot receive until all data has been transmitted. | ||||
""" | ||||
self._deferoutput = deferoutput | ||||
Gregory Szorc
|
r37070 | self._state = 'idle' | ||
Gregory Szorc
|
r37305 | self._nextoutgoingstreamid = 2 | ||
Gregory Szorc
|
r37074 | self._bufferedframegens = [] | ||
Gregory Szorc
|
r37304 | # stream id -> stream instance for all active streams from the client. | ||
self._incomingstreams = {} | ||||
Gregory Szorc
|
r37305 | self._outgoingstreams = {} | ||
Gregory Szorc
|
r37076 | # request id -> dict of commands that are actively being received. | ||
self._receivingcommands = {} | ||||
Gregory Szorc
|
r37081 | # Request IDs that have been received and are actively being processed. | ||
# Once all output for a request has been sent, it is removed from this | ||||
# set. | ||||
self._activecommands = set() | ||||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37079 | def onframerecv(self, frame): | ||
Gregory Szorc
|
r37070 | """Process a frame that has been received off the wire. | ||
Returns a dict with an ``action`` key that details what action, | ||||
if any, the consumer should take next. | ||||
""" | ||||
Gregory Szorc
|
r37304 | if not frame.streamid % 2: | ||
self._state = 'errored' | ||||
return self._makeerrorresult( | ||||
_('received frame with even numbered stream ID: %d') % | ||||
frame.streamid) | ||||
if frame.streamid not in self._incomingstreams: | ||||
if not frame.streamflags & STREAM_FLAG_BEGIN_STREAM: | ||||
self._state = 'errored' | ||||
return self._makeerrorresult( | ||||
_('received frame on unknown inactive stream without ' | ||||
'beginning of stream flag set')) | ||||
self._incomingstreams[frame.streamid] = stream(frame.streamid) | ||||
if frame.streamflags & STREAM_FLAG_ENCODING_APPLIED: | ||||
# TODO handle decoding frames | ||||
self._state = 'errored' | ||||
raise error.ProgrammingError('support for decoding stream payloads ' | ||||
'not yet implemented') | ||||
if frame.streamflags & STREAM_FLAG_END_STREAM: | ||||
del self._incomingstreams[frame.streamid] | ||||
Gregory Szorc
|
r37070 | handlers = { | ||
'idle': self._onframeidle, | ||||
Gregory Szorc
|
r37076 | 'command-receiving': self._onframecommandreceiving, | ||
Gregory Szorc
|
r37070 | 'errored': self._onframeerrored, | ||
} | ||||
meth = handlers.get(self._state) | ||||
if not meth: | ||||
raise error.ProgrammingError('unhandled state: %s' % self._state) | ||||
Gregory Szorc
|
r37079 | return meth(frame) | ||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37303 | def onbytesresponseready(self, stream, requestid, data): | ||
Gregory Szorc
|
r37073 | """Signal that a bytes response is ready to be sent to the client. | ||
The raw bytes response is passed as an argument. | ||||
""" | ||||
Gregory Szorc
|
r37304 | ensureserverstream(stream) | ||
Gregory Szorc
|
r37081 | def sendframes(): | ||
Gregory Szorc
|
r37303 | for frame in createbytesresponseframesfrombytes(stream, requestid, | ||
data): | ||||
Gregory Szorc
|
r37081 | yield frame | ||
self._activecommands.remove(requestid) | ||||
result = sendframes() | ||||
Gregory Szorc
|
r37074 | |||
if self._deferoutput: | ||||
Gregory Szorc
|
r37081 | self._bufferedframegens.append(result) | ||
Gregory Szorc
|
r37074 | return 'noop', {} | ||
else: | ||||
return 'sendframes', { | ||||
Gregory Szorc
|
r37081 | 'framegen': result, | ||
Gregory Szorc
|
r37074 | } | ||
def oninputeof(self): | ||||
"""Signals that end of input has been received. | ||||
No more frames will be received. All pending activity should be | ||||
completed. | ||||
""" | ||||
Gregory Szorc
|
r37076 | # TODO should we do anything about in-flight commands? | ||
Gregory Szorc
|
r37074 | if not self._deferoutput or not self._bufferedframegens: | ||
return 'noop', {} | ||||
# If we buffered all our responses, emit those. | ||||
def makegen(): | ||||
for gen in self._bufferedframegens: | ||||
for frame in gen: | ||||
yield frame | ||||
Gregory Szorc
|
r37073 | return 'sendframes', { | ||
Gregory Szorc
|
r37074 | 'framegen': makegen(), | ||
Gregory Szorc
|
r37073 | } | ||
Gregory Szorc
|
r37303 | def onapplicationerror(self, stream, requestid, msg): | ||
Gregory Szorc
|
r37304 | ensureserverstream(stream) | ||
Gregory Szorc
|
r37073 | return 'sendframes', { | ||
Gregory Szorc
|
r37303 | 'framegen': createerrorframe(stream, requestid, msg, | ||
application=True), | ||||
Gregory Szorc
|
r37073 | } | ||
Gregory Szorc
|
r37305 | def makeoutputstream(self): | ||
"""Create a stream to be used for sending data to the client.""" | ||||
streamid = self._nextoutgoingstreamid | ||||
self._nextoutgoingstreamid += 2 | ||||
s = stream(streamid) | ||||
self._outgoingstreams[streamid] = s | ||||
return s | ||||
Gregory Szorc
|
r37070 | def _makeerrorresult(self, msg): | ||
return 'error', { | ||||
'message': msg, | ||||
} | ||||
Gregory Szorc
|
r37076 | def _makeruncommandresult(self, requestid): | ||
entry = self._receivingcommands[requestid] | ||||
Gregory Szorc
|
r37308 | |||
if not entry['requestdone']: | ||||
self._state = 'errored' | ||||
raise error.ProgrammingError('should not be called without ' | ||||
'requestdone set') | ||||
Gregory Szorc
|
r37076 | del self._receivingcommands[requestid] | ||
if self._receivingcommands: | ||||
self._state = 'command-receiving' | ||||
else: | ||||
self._state = 'idle' | ||||
Gregory Szorc
|
r37308 | # Decode the payloads as CBOR. | ||
entry['payload'].seek(0) | ||||
request = cbor.load(entry['payload']) | ||||
if b'name' not in request: | ||||
self._state = 'errored' | ||||
return self._makeerrorresult( | ||||
_('command request missing "name" field')) | ||||
if b'args' not in request: | ||||
request[b'args'] = {} | ||||
Gregory Szorc
|
r37081 | assert requestid not in self._activecommands | ||
self._activecommands.add(requestid) | ||||
Gregory Szorc
|
r37070 | return 'runcommand', { | ||
Gregory Szorc
|
r37076 | 'requestid': requestid, | ||
Gregory Szorc
|
r37308 | 'command': request[b'name'], | ||
'args': request[b'args'], | ||||
Gregory Szorc
|
r37076 | 'data': entry['data'].getvalue() if entry['data'] else None, | ||
Gregory Szorc
|
r37070 | } | ||
def _makewantframeresult(self): | ||||
return 'wantframe', { | ||||
'state': self._state, | ||||
} | ||||
Gregory Szorc
|
r37308 | def _validatecommandrequestframe(self, frame): | ||
new = frame.flags & FLAG_COMMAND_REQUEST_NEW | ||||
continuation = frame.flags & FLAG_COMMAND_REQUEST_CONTINUATION | ||||
if new and continuation: | ||||
self._state = 'errored' | ||||
return self._makeerrorresult( | ||||
_('received command request frame with both new and ' | ||||
'continuation flags set')) | ||||
if not new and not continuation: | ||||
self._state = 'errored' | ||||
return self._makeerrorresult( | ||||
_('received command request frame with neither new nor ' | ||||
'continuation flags set')) | ||||
Gregory Szorc
|
r37079 | def _onframeidle(self, frame): | ||
Gregory Szorc
|
r37070 | # The only frame type that should be received in this state is a | ||
# command request. | ||||
Gregory Szorc
|
r37308 | if frame.typeid != FRAME_TYPE_COMMAND_REQUEST: | ||
Gregory Szorc
|
r37070 | self._state = 'errored' | ||
return self._makeerrorresult( | ||||
Gregory Szorc
|
r37308 | _('expected command request frame; got %d') % frame.typeid) | ||
res = self._validatecommandrequestframe(frame) | ||||
if res: | ||||
return res | ||||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37079 | if frame.requestid in self._receivingcommands: | ||
Gregory Szorc
|
r37076 | self._state = 'errored' | ||
return self._makeerrorresult( | ||||
Gregory Szorc
|
r37079 | _('request with ID %d already received') % frame.requestid) | ||
Gregory Szorc
|
r37076 | |||
Gregory Szorc
|
r37081 | if frame.requestid in self._activecommands: | ||
self._state = 'errored' | ||||
Gregory Szorc
|
r37308 | return self._makeerrorresult( | ||
_('request with ID %d is already active') % frame.requestid) | ||||
new = frame.flags & FLAG_COMMAND_REQUEST_NEW | ||||
moreframes = frame.flags & FLAG_COMMAND_REQUEST_MORE_FRAMES | ||||
expectingdata = frame.flags & FLAG_COMMAND_REQUEST_EXPECT_DATA | ||||
Gregory Szorc
|
r37081 | |||
Gregory Szorc
|
r37308 | if not new: | ||
self._state = 'errored' | ||||
return self._makeerrorresult( | ||||
_('received command request frame without new flag set')) | ||||
payload = util.bytesio() | ||||
payload.write(frame.payload) | ||||
Gregory Szorc
|
r37076 | |||
Gregory Szorc
|
r37079 | self._receivingcommands[frame.requestid] = { | ||
Gregory Szorc
|
r37308 | 'payload': payload, | ||
Gregory Szorc
|
r37076 | 'data': None, | ||
Gregory Szorc
|
r37308 | 'requestdone': not moreframes, | ||
'expectingdata': bool(expectingdata), | ||||
Gregory Szorc
|
r37076 | } | ||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37308 | # This is the final frame for this request. Dispatch it. | ||
if not moreframes and not expectingdata: | ||||
Gregory Szorc
|
r37079 | return self._makeruncommandresult(frame.requestid) | ||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37308 | assert moreframes or expectingdata | ||
self._state = 'command-receiving' | ||||
return self._makewantframeresult() | ||||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37079 | def _onframecommandreceiving(self, frame): | ||
Gregory Szorc
|
r37308 | if frame.typeid == FRAME_TYPE_COMMAND_REQUEST: | ||
# Process new command requests as such. | ||||
if frame.flags & FLAG_COMMAND_REQUEST_NEW: | ||||
return self._onframeidle(frame) | ||||
res = self._validatecommandrequestframe(frame) | ||||
if res: | ||||
return res | ||||
Gregory Szorc
|
r37076 | |||
# All other frames should be related to a command that is currently | ||||
Gregory Szorc
|
r37081 | # receiving but is not active. | ||
if frame.requestid in self._activecommands: | ||||
self._state = 'errored' | ||||
return self._makeerrorresult( | ||||
_('received frame for request that is still active: %d') % | ||||
frame.requestid) | ||||
Gregory Szorc
|
r37079 | if frame.requestid not in self._receivingcommands: | ||
Gregory Szorc
|
r37070 | self._state = 'errored' | ||
Gregory Szorc
|
r37076 | return self._makeerrorresult( | ||
_('received frame for request that is not receiving: %d') % | ||||
Gregory Szorc
|
r37079 | frame.requestid) | ||
Gregory Szorc
|
r37076 | |||
Gregory Szorc
|
r37079 | entry = self._receivingcommands[frame.requestid] | ||
Gregory Szorc
|
r37076 | |||
Gregory Szorc
|
r37308 | if frame.typeid == FRAME_TYPE_COMMAND_REQUEST: | ||
moreframes = frame.flags & FLAG_COMMAND_REQUEST_MORE_FRAMES | ||||
expectingdata = bool(frame.flags & FLAG_COMMAND_REQUEST_EXPECT_DATA) | ||||
if entry['requestdone']: | ||||
self._state = 'errored' | ||||
return self._makeerrorresult( | ||||
_('received command request frame when request frames ' | ||||
'were supposedly done')) | ||||
if expectingdata != entry['expectingdata']: | ||||
Gregory Szorc
|
r37076 | self._state = 'errored' | ||
Gregory Szorc
|
r37308 | return self._makeerrorresult( | ||
_('mismatch between expect data flag and previous frame')) | ||||
entry['payload'].write(frame.payload) | ||||
Gregory Szorc
|
r37076 | |||
Gregory Szorc
|
r37308 | if not moreframes: | ||
entry['requestdone'] = True | ||||
if not moreframes and not expectingdata: | ||||
return self._makeruncommandresult(frame.requestid) | ||||
return self._makewantframeresult() | ||||
Gregory Szorc
|
r37076 | |||
Gregory Szorc
|
r37079 | elif frame.typeid == FRAME_TYPE_COMMAND_DATA: | ||
Gregory Szorc
|
r37076 | if not entry['expectingdata']: | ||
self._state = 'errored' | ||||
return self._makeerrorresult(_( | ||||
'received command data frame for request that is not ' | ||||
Gregory Szorc
|
r37079 | 'expecting data: %d') % frame.requestid) | ||
Gregory Szorc
|
r37076 | |||
if entry['data'] is None: | ||||
entry['data'] = util.bytesio() | ||||
Gregory Szorc
|
r37079 | return self._handlecommanddataframe(frame, entry) | ||
Gregory Szorc
|
r37308 | else: | ||
Gregory Szorc
|
r37070 | self._state = 'errored' | ||
Gregory Szorc
|
r37308 | return self._makeerrorresult(_( | ||
'received unexpected frame type: %d') % frame.typeid) | ||||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37079 | def _handlecommanddataframe(self, frame, entry): | ||
assert frame.typeid == FRAME_TYPE_COMMAND_DATA | ||||
Gregory Szorc
|
r37070 | |||
# TODO support streaming data instead of buffering it. | ||||
Gregory Szorc
|
r37079 | entry['data'].write(frame.payload) | ||
Gregory Szorc
|
r37070 | |||
Gregory Szorc
|
r37079 | if frame.flags & FLAG_COMMAND_DATA_CONTINUATION: | ||
Gregory Szorc
|
r37070 | return self._makewantframeresult() | ||
Gregory Szorc
|
r37079 | elif frame.flags & FLAG_COMMAND_DATA_EOS: | ||
Gregory Szorc
|
r37076 | entry['data'].seek(0) | ||
Gregory Szorc
|
r37079 | return self._makeruncommandresult(frame.requestid) | ||
Gregory Szorc
|
r37070 | else: | ||
self._state = 'errored' | ||||
return self._makeerrorresult(_('command data frame without ' | ||||
'flags')) | ||||
Gregory Szorc
|
r37079 | def _onframeerrored(self, frame): | ||
Gregory Szorc
|
r37070 | return self._makeerrorresult(_('server already errored')) | ||