serialization.py 17 KB


  1. """Convert between bytestreams and higher-level AMQP types.
  2. 2007-11-05 Barry Pederson <bp@barryp.org>
  3. """
  4. # Copyright (C) 2007 Barry Pederson <bp@barryp.org>
  5. from __future__ import absolute_import, unicode_literals
  6. import calendar
  7. import sys
  8. from datetime import datetime
  9. from decimal import Decimal
  10. from io import BytesIO
  11. from .exceptions import FrameSyntaxError
  12. from .five import int_types, items, long_t, string, string_t
  13. from .platform import pack, unpack_from
  14. from .spec import Basic
  15. from .utils import bytes_to_str as pstr_t
  16. from .utils import str_to_bytes
  17. ftype_t = chr if sys.version_info[0] == 3 else None
  18. ILLEGAL_TABLE_TYPE = """\
  19. Table type {0!r} not handled by amqp.
  20. """
  21. ILLEGAL_TABLE_TYPE_WITH_KEY = """\
  22. Table type {0!r} for key {1!r} not handled by amqp. [value: {2!r}]
  23. """
  24. ILLEGAL_TABLE_TYPE_WITH_VALUE = """\
  25. Table type {0!r} not handled by amqp. [value: {1!r}]
  26. """
  27. def _read_item(buf, offset=0, unpack_from=unpack_from, ftype_t=ftype_t):
  28. ftype = ftype_t(buf[offset]) if ftype_t else buf[offset]
  29. offset += 1
  30. # 'S': long string
  31. if ftype == 'S':
  32. slen, = unpack_from('>I', buf, offset)
  33. offset += 4
  34. val = pstr_t(buf[offset:offset + slen])
  35. offset += slen
  36. # 's': short string
  37. elif ftype == 's':
  38. slen, = unpack_from('>B', buf, offset)
  39. offset += 1
  40. val = pstr_t(buf[offset:offset + slen])
  41. offset += slen
  42. # 'x': Bytes Array
  43. elif ftype == 'x':
  44. blen, = unpack_from('>I', buf, offset)
  45. offset += 4
  46. val = buf[offset:offset + blen]
  47. offset += blen
  48. # 'b': short-short int
  49. elif ftype == 'b':
  50. val, = unpack_from('>B', buf, offset)
  51. offset += 1
  52. # 'B': short-short unsigned int
  53. elif ftype == 'B':
  54. val, = unpack_from('>b', buf, offset)
  55. offset += 1
  56. # 'U': short int
  57. elif ftype == 'U':
  58. val, = unpack_from('>h', buf, offset)
  59. offset += 2
  60. # 'u': short unsigned int
  61. elif ftype == 'u':
  62. val, = unpack_from('>H', buf, offset)
  63. offset += 2
  64. # 'I': long int
  65. elif ftype == 'I':
  66. val, = unpack_from('>i', buf, offset)
  67. offset += 4
  68. # 'i': long unsigned int
  69. elif ftype == 'i':
  70. val, = unpack_from('>I', buf, offset)
  71. offset += 4
  72. # 'L': long long int
  73. elif ftype == 'L':
  74. val, = unpack_from('>q', buf, offset)
  75. offset += 8
  76. # 'l': long long unsigned int
  77. elif ftype == 'l':
  78. val, = unpack_from('>Q', buf, offset)
  79. offset += 8
  80. # 'f': float
  81. elif ftype == 'f':
  82. val, = unpack_from('>f', buf, offset)
  83. offset += 4
  84. # 'd': double
  85. elif ftype == 'd':
  86. val, = unpack_from('>d', buf, offset)
  87. offset += 8
  88. # 'D': decimal
  89. elif ftype == 'D':
  90. d, = unpack_from('>B', buf, offset)
  91. offset += 1
  92. n, = unpack_from('>i', buf, offset)
  93. offset += 4
  94. val = Decimal(n) / Decimal(10 ** d)
  95. # 'F': table
  96. elif ftype == 'F':
  97. tlen, = unpack_from('>I', buf, offset)
  98. offset += 4
  99. limit = offset + tlen
  100. val = {}
  101. while offset < limit:
  102. keylen, = unpack_from('>B', buf, offset)
  103. offset += 1
  104. key = pstr_t(buf[offset:offset + keylen])
  105. offset += keylen
  106. val[key], offset = _read_item(buf, offset)
  107. # 'A': array
  108. elif ftype == 'A':
  109. alen, = unpack_from('>I', buf, offset)
  110. offset += 4
  111. limit = offset + alen
  112. val = []
  113. while offset < limit:
  114. v, offset = _read_item(buf, offset)
  115. val.append(v)
  116. # 't' (bool)
  117. elif ftype == 't':
  118. val, = unpack_from('>B', buf, offset)
  119. val = bool(val)
  120. offset += 1
  121. # 'T': timestamp
  122. elif ftype == 'T':
  123. val, = unpack_from('>Q', buf, offset)
  124. offset += 8
  125. val = datetime.utcfromtimestamp(val)
  126. # 'V': void
  127. elif ftype == 'V':
  128. val = None
  129. else:
  130. raise FrameSyntaxError(
  131. 'Unknown value in table: {0!r} ({1!r})'.format(
  132. ftype, type(ftype)))
  133. return val, offset
  134. def loads(format, buf, offset=0,
  135. ord=ord, unpack_from=unpack_from,
  136. _read_item=_read_item, pstr_t=pstr_t):
  137. """Deserialize amqp format.
  138. bit = b
  139. octet = o
  140. short = B
  141. long = l
  142. long long = L
  143. float = f
  144. shortstr = s
  145. longstr = S
  146. table = F
  147. array = A
  148. timestamp = T
  149. """
  150. bitcount = bits = 0
  151. values = []
  152. append = values.append
  153. format = pstr_t(format)
  154. for p in format:
  155. if p == 'b':
  156. if not bitcount:
  157. bits = ord(buf[offset:offset + 1])
  158. bitcount = 8
  159. val = (bits & 1) == 1
  160. bits >>= 1
  161. bitcount -= 1
  162. offset += 1
  163. elif p == 'o':
  164. bitcount = bits = 0
  165. val, = unpack_from('>B', buf, offset)
  166. offset += 1
  167. elif p == 'B':
  168. bitcount = bits = 0
  169. val, = unpack_from('>H', buf, offset)
  170. offset += 2
  171. elif p == 'l':
  172. bitcount = bits = 0
  173. val, = unpack_from('>I', buf, offset)
  174. offset += 4
  175. elif p == 'L':
  176. bitcount = bits = 0
  177. val, = unpack_from('>Q', buf, offset)
  178. offset += 8
  179. elif p == 'f':
  180. bitcount = bits = 0
  181. val, = unpack_from('>f', buf, offset)
  182. offset += 4
  183. elif p == 's':
  184. bitcount = bits = 0
  185. slen, = unpack_from('B', buf, offset)
  186. offset += 1
  187. val = buf[offset:offset + slen].decode('utf-8', 'surrogatepass')
  188. offset += slen
  189. elif p == 'S':
  190. bitcount = bits = 0
  191. slen, = unpack_from('>I', buf, offset)
  192. offset += 4
  193. val = buf[offset:offset + slen].decode('utf-8', 'surrogatepass')
  194. offset += slen
  195. elif p == 'x':
  196. blen, = unpack_from('>I', buf, offset)
  197. offset += 4
  198. val = buf[offset:offset + blen]
  199. offset += blen
  200. elif p == 'F':
  201. bitcount = bits = 0
  202. tlen, = unpack_from('>I', buf, offset)
  203. offset += 4
  204. limit = offset + tlen
  205. val = {}
  206. while offset < limit:
  207. keylen, = unpack_from('>B', buf, offset)
  208. offset += 1
  209. key = pstr_t(buf[offset:offset + keylen])
  210. offset += keylen
  211. val[key], offset = _read_item(buf, offset)
  212. elif p == 'A':
  213. bitcount = bits = 0
  214. alen, = unpack_from('>I', buf, offset)
  215. offset += 4
  216. limit = offset + alen
  217. val = []
  218. while offset < limit:
  219. aval, offset = _read_item(buf, offset)
  220. val.append(aval)
  221. elif p == 'T':
  222. bitcount = bits = 0
  223. val, = unpack_from('>Q', buf, offset)
  224. offset += 8
  225. val = datetime.utcfromtimestamp(val)
  226. else:
  227. raise FrameSyntaxError(ILLEGAL_TABLE_TYPE.format(p))
  228. append(val)
  229. return values, offset
  230. def _flushbits(bits, write, pack=pack):
  231. if bits:
  232. write(pack('B' * len(bits), *bits))
  233. bits[:] = []
  234. return 0
  235. def dumps(format, values):
  236. """Serialize AMQP arguments.
  237. Notes:
  238. bit = b
  239. octet = o
  240. short = B
  241. long = l
  242. long long = L
  243. shortstr = s
  244. longstr = S
  245. byte array = x
  246. table = F
  247. array = A
  248. """
  249. bitcount = 0
  250. bits = []
  251. out = BytesIO()
  252. write = out.write
  253. format = pstr_t(format)
  254. for i, val in enumerate(values):
  255. p = format[i]
  256. if p == 'b':
  257. val = 1 if val else 0
  258. shift = bitcount % 8
  259. if shift == 0:
  260. bits.append(0)
  261. bits[-1] |= (val << shift)
  262. bitcount += 1
  263. elif p == 'o':
  264. bitcount = _flushbits(bits, write)
  265. write(pack('B', val))
  266. elif p == 'B':
  267. bitcount = _flushbits(bits, write)
  268. write(pack('>H', int(val)))
  269. elif p == 'l':
  270. bitcount = _flushbits(bits, write)
  271. write(pack('>I', val))
  272. elif p == 'L':
  273. bitcount = _flushbits(bits, write)
  274. write(pack('>Q', val))
  275. elif p == 'f':
  276. bitcount = _flushbits(bits, write)
  277. write(pack('>f', val))
  278. elif p == 's':
  279. val = val or ''
  280. bitcount = _flushbits(bits, write)
  281. if isinstance(val, string):
  282. val = val.encode('utf-8', 'surrogatepass')
  283. write(pack('B', len(val)))
  284. write(val)
  285. elif p == 'S' or p == 'x':
  286. val = val or ''
  287. bitcount = _flushbits(bits, write)
  288. if isinstance(val, string):
  289. val = val.encode('utf-8', 'surrogatepass')
  290. write(pack('>I', len(val)))
  291. write(val)
  292. elif p == 'F':
  293. bitcount = _flushbits(bits, write)
  294. _write_table(val or {}, write, bits)
  295. elif p == 'A':
  296. bitcount = _flushbits(bits, write)
  297. _write_array(val or [], write, bits)
  298. elif p == 'T':
  299. write(pack('>Q', long_t(calendar.timegm(val.utctimetuple()))))
  300. _flushbits(bits, write)
  301. return out.getvalue()
  302. def _write_table(d, write, bits, pack=pack):
  303. out = BytesIO()
  304. twrite = out.write
  305. for k, v in items(d):
  306. if isinstance(k, string):
  307. k = k.encode('utf-8', 'surrogatepass')
  308. twrite(pack('B', len(k)))
  309. twrite(k)
  310. try:
  311. _write_item(v, twrite, bits)
  312. except ValueError:
  313. raise FrameSyntaxError(
  314. ILLEGAL_TABLE_TYPE_WITH_KEY.format(type(v), k, v))
  315. table_data = out.getvalue()
  316. write(pack('>I', len(table_data)))
  317. write(table_data)
  318. def _write_array(l, write, bits, pack=pack):
  319. out = BytesIO()
  320. awrite = out.write
  321. for v in l:
  322. try:
  323. _write_item(v, awrite, bits)
  324. except ValueError:
  325. raise FrameSyntaxError(
  326. ILLEGAL_TABLE_TYPE_WITH_VALUE.format(type(v), v))
  327. array_data = out.getvalue()
  328. write(pack('>I', len(array_data)))
  329. write(array_data)
  330. def _write_item(v, write, bits, pack=pack,
  331. string_t=string_t, bytes=bytes, string=string, bool=bool,
  332. float=float, int_types=int_types, Decimal=Decimal,
  333. datetime=datetime, dict=dict, list=list, tuple=tuple,
  334. None_t=None):
  335. if isinstance(v, (string_t, bytes)):
  336. if isinstance(v, string):
  337. v = v.encode('utf-8', 'surrogatepass')
  338. write(pack('>cI', b'S', len(v)))
  339. write(v)
  340. elif isinstance(v, bool):
  341. write(pack('>cB', b't', int(v)))
  342. elif isinstance(v, float):
  343. write(pack('>cd', b'd', v))
  344. elif isinstance(v, int_types):
  345. if v > 2147483647 or v < -2147483647:
  346. write(pack('>cq', b'L', v))
  347. else:
  348. write(pack('>ci', b'I', v))
  349. elif isinstance(v, Decimal):
  350. sign, digits, exponent = v.as_tuple()
  351. v = 0
  352. for d in digits:
  353. v = (v * 10) + d
  354. if sign:
  355. v = -v
  356. write(pack('>cBi', b'D', -exponent, v))
  357. elif isinstance(v, datetime):
  358. write(
  359. pack('>cQ', b'T', long_t(calendar.timegm(v.utctimetuple()))))
  360. elif isinstance(v, dict):
  361. write(b'F')
  362. _write_table(v, write, bits)
  363. elif isinstance(v, (list, tuple)):
  364. write(b'A')
  365. _write_array(v, write, bits)
  366. elif v is None_t:
  367. write(b'V')
  368. else:
  369. raise ValueError()
  370. def decode_properties_basic(buf, offset=0,
  371. unpack_from=unpack_from, pstr_t=pstr_t):
  372. """Decode basic properties."""
  373. properties = {}
  374. flags, = unpack_from('>H', buf, offset)
  375. offset += 2
  376. if flags & 0x8000:
  377. slen, = unpack_from('>B', buf, offset)
  378. offset += 1
  379. properties['content_type'] = pstr_t(buf[offset:offset + slen])
  380. offset += slen
  381. if flags & 0x4000:
  382. slen, = unpack_from('>B', buf, offset)
  383. offset += 1
  384. properties['content_encoding'] = pstr_t(buf[offset:offset + slen])
  385. offset += slen
  386. if flags & 0x2000:
  387. _f, offset = loads('F', buf, offset)
  388. properties['application_headers'], = _f
  389. if flags & 0x1000:
  390. properties['delivery_mode'], = unpack_from('>B', buf, offset)
  391. offset += 1
  392. if flags & 0x0800:
  393. properties['priority'], = unpack_from('>B', buf, offset)
  394. offset += 1
  395. if flags & 0x0400:
  396. slen, = unpack_from('>B', buf, offset)
  397. offset += 1
  398. properties['correlation_id'] = pstr_t(buf[offset:offset + slen])
  399. offset += slen
  400. if flags & 0x0200:
  401. slen, = unpack_from('>B', buf, offset)
  402. offset += 1
  403. properties['reply_to'] = pstr_t(buf[offset:offset + slen])
  404. offset += slen
  405. if flags & 0x0100:
  406. slen, = unpack_from('>B', buf, offset)
  407. offset += 1
  408. properties['expiration'] = pstr_t(buf[offset:offset + slen])
  409. offset += slen
  410. if flags & 0x0080:
  411. slen, = unpack_from('>B', buf, offset)
  412. offset += 1
  413. properties['message_id'] = pstr_t(buf[offset:offset + slen])
  414. offset += slen
  415. if flags & 0x0040:
  416. properties['timestamp'], = unpack_from('>Q', buf, offset)
  417. offset += 8
  418. if flags & 0x0020:
  419. slen, = unpack_from('>B', buf, offset)
  420. offset += 1
  421. properties['type'] = pstr_t(buf[offset:offset + slen])
  422. offset += slen
  423. if flags & 0x0010:
  424. slen, = unpack_from('>B', buf, offset)
  425. offset += 1
  426. properties['user_id'] = pstr_t(buf[offset:offset + slen])
  427. offset += slen
  428. if flags & 0x0008:
  429. slen, = unpack_from('>B', buf, offset)
  430. offset += 1
  431. properties['app_id'] = pstr_t(buf[offset:offset + slen])
  432. offset += slen
  433. if flags & 0x0004:
  434. slen, = unpack_from('>B', buf, offset)
  435. offset += 1
  436. properties['cluster_id'] = pstr_t(buf[offset:offset + slen])
  437. offset += slen
  438. return properties, offset
  439. PROPERTY_CLASSES = {
  440. Basic.CLASS_ID: decode_properties_basic,
  441. }
  442. class GenericContent(object):
  443. """Abstract base class for AMQP content.
  444. Subclasses should override the PROPERTIES attribute.
  445. """
  446. CLASS_ID = None
  447. PROPERTIES = [('dummy', 's')]
  448. def __init__(self, frame_method=None, frame_args=None, **props):
  449. self.frame_method = frame_method
  450. self.frame_args = frame_args
  451. self.properties = props
  452. self._pending_chunks = []
  453. self.body_received = 0
  454. self.body_size = 0
  455. self.ready = False
  456. def __getattr__(self, name):
  457. # Look for additional properties in the 'properties'
  458. # dictionary, and if present - the 'delivery_info' dictionary.
  459. if name == '__setstate__':
  460. # Allows pickling/unpickling to work
  461. raise AttributeError('__setstate__')
  462. if name in self.properties:
  463. return self.properties[name]
  464. raise AttributeError(name)
  465. def _load_properties(self, class_id, buf, offset=0,
  466. classes=PROPERTY_CLASSES, unpack_from=unpack_from):
  467. """Load AMQP properties.
  468. Given the raw bytes containing the property-flags and property-list
  469. from a content-frame-header, parse and insert into a dictionary
  470. stored in this object as an attribute named 'properties'.
  471. """
  472. # Read 16-bit shorts until we get one with a low bit set to zero
  473. props, offset = classes[class_id](buf, offset)
  474. self.properties = props
  475. return offset
  476. def _serialize_properties(self):
  477. """Serialize AMQP properties.
  478. Serialize the 'properties' attribute (a dictionary) into
  479. the raw bytes making up a set of property flags and a
  480. property list, suitable for putting into a content frame header.
  481. """
  482. shift = 15
  483. flag_bits = 0
  484. flags = []
  485. sformat, svalues = [], []
  486. props = self.properties
  487. for key, proptype in self.PROPERTIES:
  488. val = props.get(key, None)
  489. if val is not None:
  490. if shift == 0:
  491. flags.append(flag_bits)
  492. flag_bits = 0
  493. shift = 15
  494. flag_bits |= (1 << shift)
  495. if proptype != 'bit':
  496. sformat.append(str_to_bytes(proptype))
  497. svalues.append(val)
  498. shift -= 1
  499. flags.append(flag_bits)
  500. result = BytesIO()
  501. write = result.write
  502. for flag_bits in flags:
  503. write(pack('>H', flag_bits))
  504. write(dumps(b''.join(sformat), svalues))
  505. return result.getvalue()
  506. def inbound_header(self, buf, offset=0):
  507. class_id, self.body_size = unpack_from('>HxxQ', buf, offset)
  508. offset += 12
  509. self._load_properties(class_id, buf, offset)
  510. if not self.body_size:
  511. self.ready = True
  512. return offset
  513. def inbound_body(self, buf):
  514. chunks = self._pending_chunks
  515. self.body_received += len(buf)
  516. if self.body_received >= self.body_size:
  517. if chunks:
  518. chunks.append(buf)
  519. self.body = bytes().join(chunks)
  520. chunks[:] = []
  521. else:
  522. self.body = buf
  523. self.ready = True
  524. else:
  525. chunks.append(buf)