suite.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460
  1. from __future__ import absolute_import, print_function, unicode_literals
  2. import inspect
  3. import platform
  4. import random
  5. import socket
  6. import sys
  7. from collections import OrderedDict, defaultdict, namedtuple
  8. from itertools import count
  9. from time import sleep
  10. from celery import VERSION_BANNER, chain, group, uuid
  11. from celery.exceptions import TimeoutError
  12. from celery.five import items, monotonic, range, values
  13. from celery.utils.debug import blockdetection
  14. from celery.utils.text import pluralize, truncate
  15. from celery.utils.timeutils import humanize_seconds
  16. from .app import (
  17. marker, _marker, add, any_, collect_ids, exiting, ids, kill, sleeping,
  18. sleeping_ignore_limits, any_returning, print_unicode,
  19. )
  20. from .data import BIG, SMALL
  21. from .fbi import FBI
  22. BANNER = """\
  23. Celery stress-suite v{version}
  24. {platform}
  25. [config]
  26. .> broker: {conninfo}
  27. [toc: {total} {TESTS} total]
  28. {toc}
  29. """
  30. F_PROGRESS = """\
  31. {0.index}: {0.test.__name__}({0.iteration}/{0.total_iterations}) \
  32. rep#{0.repeats} runtime: {runtime}/{elapsed} \
  33. """
  34. Progress = namedtuple('Progress', (
  35. 'test', 'iteration', 'total_iterations',
  36. 'index', 'repeats', 'runtime', 'elapsed', 'completed',
  37. ))
  38. Inf = float('Inf')
  39. def assert_equal(a, b):
  40. assert a == b, '{0!r} != {1!r}'.format(a, b)
  41. class StopSuite(Exception):
  42. pass
  43. def pstatus(p):
  44. runtime = monotonic() - p.runtime
  45. elapsed = monotonic() - p.elapsed
  46. return F_PROGRESS.format(
  47. p,
  48. runtime=humanize_seconds(runtime, now=runtime),
  49. elapsed=humanize_seconds(elapsed, now=elapsed),
  50. )
  51. class Speaker(object):
  52. def __init__(self, gap=5.0):
  53. self.gap = gap
  54. self.last_noise = monotonic() - self.gap * 2
  55. def beep(self):
  56. now = monotonic()
  57. if now - self.last_noise >= self.gap:
  58. self.emit()
  59. self.last_noise = now
  60. def emit(self):
  61. print('\a', file=sys.stderr, end='')
  62. def testgroup(*funs):
  63. return OrderedDict((fun.__name__, fun) for fun in funs)
  64. class BaseSuite(object):
  65. def __init__(self, app, block_timeout=30 * 60):
  66. self.app = app
  67. self.connerrors = self.app.connection().recoverable_connection_errors
  68. self.block_timeout = block_timeout
  69. self.progress = None
  70. self.speaker = Speaker()
  71. self.fbi = FBI(app)
  72. self.init_groups()
  73. def init_groups(self):
  74. acc = defaultdict(list)
  75. for attr in dir(self):
  76. if not _is_descriptor(self, attr):
  77. meth = getattr(self, attr)
  78. try:
  79. groups = meth.__func__.__testgroup__
  80. except AttributeError:
  81. pass
  82. else:
  83. for g in groups:
  84. acc[g].append(meth)
  85. # sort the tests by the order in which they are defined in the class
  86. for g in values(acc):
  87. g[:] = sorted(g, key=lambda m: m.__func__.__testsort__)
  88. self.groups = dict(
  89. (name, testgroup(*tests)) for name, tests in items(acc)
  90. )
  91. def run(self, names=None, iterations=50, offset=0,
  92. numtests=None, list_all=False, repeat=0, group='all',
  93. diag=False, no_join=False, **kw):
  94. self.no_join = no_join
  95. self.fbi.enable(diag)
  96. tests = self.filtertests(group, names)[offset:numtests or None]
  97. if list_all:
  98. return print(self.testlist(tests))
  99. print(self.banner(tests))
  100. print('+ Enabling events')
  101. self.app.control.enable_events()
  102. it = count() if repeat == Inf else range(int(repeat) or 1)
  103. for i in it:
  104. marker(
  105. 'Stresstest suite start (repetition {0})'.format(i + 1),
  106. '+',
  107. )
  108. for j, test in enumerate(tests):
  109. self.runtest(test, iterations, j + 1, i + 1)
  110. marker(
  111. 'Stresstest suite end (repetition {0})'.format(i + 1),
  112. '+',
  113. )
  114. def filtertests(self, group, names):
  115. tests = self.groups[group]
  116. try:
  117. return ([tests[n] for n in names] if names
  118. else list(values(tests)))
  119. except KeyError as exc:
  120. raise KeyError('Unknown test name: {0}'.format(exc))
  121. def testlist(self, tests):
  122. return ',\n'.join(
  123. '.> {0}) {1}'.format(i + 1, t.__name__)
  124. for i, t in enumerate(tests)
  125. )
  126. def banner(self, tests):
  127. app = self.app
  128. return BANNER.format(
  129. app='{0}:0x{1:x}'.format(app.main or '__main__', id(app)),
  130. version=VERSION_BANNER,
  131. conninfo=app.connection().as_uri(),
  132. platform=platform.platform(),
  133. toc=self.testlist(tests),
  134. TESTS=pluralize(len(tests), 'test'),
  135. total=len(tests),
  136. )
  137. def runtest(self, fun, n=50, index=0, repeats=1):
  138. n = getattr(fun, '__iterations__', None) or n
  139. print('{0}: [[[{1}({2})]]]'.format(repeats, fun.__name__, n))
  140. with blockdetection(self.block_timeout):
  141. with self.fbi.investigation():
  142. runtime = elapsed = monotonic()
  143. i = 0
  144. failed = False
  145. self.progress = Progress(
  146. fun, i, n, index, repeats, elapsed, runtime, 0,
  147. )
  148. _marker.delay(pstatus(self.progress))
  149. try:
  150. for i in range(n):
  151. runtime = monotonic()
  152. self.progress = Progress(
  153. fun, i + 1, n, index, repeats, runtime, elapsed, 0,
  154. )
  155. try:
  156. fun()
  157. except StopSuite:
  158. raise
  159. except Exception as exc:
  160. print('-> {0!r}'.format(exc))
  161. import traceback
  162. print(traceback.format_exc())
  163. print(pstatus(self.progress))
  164. else:
  165. print(pstatus(self.progress))
  166. except Exception:
  167. failed = True
  168. self.speaker.beep()
  169. raise
  170. finally:
  171. print('{0} {1} iterations in {2}'.format(
  172. 'failed after' if failed else 'completed',
  173. i + 1, humanize_seconds(monotonic() - elapsed),
  174. ))
  175. if not failed:
  176. self.progress = Progress(
  177. fun, i + 1, n, index, repeats, runtime, elapsed, 1,
  178. )
  179. def missing_results(self, r):
  180. return [res.id for res in r if res.id not in res.backend._cache]
  181. def join(self, r, propagate=False, max_retries=10, **kwargs):
  182. if self.no_join:
  183. return
  184. received = []
  185. def on_result(task_id, value):
  186. received.append(task_id)
  187. for i in range(max_retries) if max_retries else count(0):
  188. received[:] = []
  189. try:
  190. return r.get(callback=on_result, propagate=propagate, **kwargs)
  191. except (socket.timeout, TimeoutError) as exc:
  192. waiting_for = self.missing_results(r)
  193. self.speaker.beep()
  194. marker(
  195. 'Still waiting for {0}/{1}: [{2}]: {3!r}'.format(
  196. len(r) - len(received), len(r),
  197. truncate(', '.join(waiting_for)), exc), '!',
  198. )
  199. self.fbi.diag(waiting_for)
  200. except self.connerrors as exc:
  201. self.speaker.beep()
  202. marker('join: connection lost: {0!r}'.format(exc), '!')
  203. raise StopSuite('Test failed: Missing task results')
  204. def dump_progress(self):
  205. return pstatus(self.progress) if self.progress else 'No test running'
  206. _creation_counter = count(0)
  207. def testcase(*groups, **kwargs):
  208. if not groups:
  209. raise ValueError('@testcase requires at least one group name')
  210. def _mark_as_case(fun):
  211. fun.__testgroup__ = groups
  212. fun.__testsort__ = next(_creation_counter)
  213. fun.__iterations__ = kwargs.get('iterations')
  214. return fun
  215. return _mark_as_case
  216. def _is_descriptor(obj, attr):
  217. try:
  218. cattr = getattr(obj.__class__, attr)
  219. except AttributeError:
  220. pass
  221. else:
  222. return not inspect.ismethod(cattr) and hasattr(cattr, '__get__')
  223. return False
  224. class Suite(BaseSuite):
  225. @testcase('all', 'green', 'redis', iterations=1)
  226. def chain(self):
  227. c = add.s(4, 4) | add.s(8) | add.s(16)
  228. assert_equal(self.join(c()), 32)
  229. @testcase('all', 'green', 'redis', iterations=1)
  230. def chaincomplex(self):
  231. c = (
  232. add.s(2, 2) | (
  233. add.s(4) | add.s(8) | add.s(16)
  234. ) |
  235. group(add.s(i) for i in range(4))
  236. )
  237. res = c()
  238. assert_equal(res.get(), [32, 33, 34, 35])
  239. @testcase('all', 'green', 'redis', iterations=1)
  240. def parentids_chain(self, num=248):
  241. c = chain(ids.si(i) for i in range(num))
  242. c.freeze()
  243. res = c()
  244. res.get(timeout=5)
  245. self.assert_ids(res, num - 1)
  246. @testcase('all', 'green', 'redis', iterations=1)
  247. def parentids_group(self):
  248. g = ids.si(1) | ids.si(2) | group(ids.si(i) for i in range(2, 50))
  249. res = g()
  250. expected_root_id = res.parent.parent.id
  251. expected_parent_id = res.parent.id
  252. values = res.get(timeout=5)
  253. for i, r in enumerate(values):
  254. root_id, parent_id, value = r
  255. assert_equal(root_id, expected_root_id)
  256. assert_equal(parent_id, expected_parent_id)
  257. assert_equal(value, i + 2)
  258. def assert_ids(self, res, size):
  259. i, root = size, res
  260. while root.parent:
  261. root = root.parent
  262. node = res
  263. while node:
  264. root_id, parent_id, value = node.get(timeout=5)
  265. assert_equal(value, i)
  266. assert_equal(root_id, root.id)
  267. if node.parent:
  268. assert_equal(parent_id, node.parent.id)
  269. node = node.parent
  270. i -= 1
  271. @testcase('redis', iterations=1)
  272. def parentids_chord(self):
  273. self.assert_parentids_chord()
  274. self.assert_parentids_chord(uuid(), uuid())
  275. def assert_parentids_chord(self, base_root=None, base_parent=None):
  276. g = (
  277. ids.si(1) |
  278. ids.si(2) |
  279. group(ids.si(i) for i in range(3, 50)) |
  280. collect_ids.s(i=50) |
  281. ids.si(51)
  282. )
  283. g.freeze(root_id=base_root, parent_id=base_parent)
  284. res = g.apply_async(root_id=base_root, parent_id=base_parent)
  285. expected_root_id = base_root or res.parent.parent.parent.id
  286. root_id, parent_id, value = res.get(timeout=5)
  287. assert_equal(value, 51)
  288. assert_equal(root_id, expected_root_id)
  289. assert_equal(parent_id, res.parent.id)
  290. prev, (root_id, parent_id, value) = res.parent.get(timeout=5)
  291. assert_equal(value, 50)
  292. assert_equal(root_id, expected_root_id)
  293. assert_equal(parent_id, res.parent.parent.id)
  294. for i, p in enumerate(prev):
  295. root_id, parent_id, value = p
  296. assert_equal(root_id, expected_root_id)
  297. assert_equal(parent_id, res.parent.parent.id)
  298. root_id, parent_id, value = res.parent.parent.get(timeout=5)
  299. assert_equal(value, 2)
  300. assert_equal(parent_id, res.parent.parent.parent.id)
  301. assert_equal(root_id, expected_root_id)
  302. root_id, parent_id, value = res.parent.parent.parent.get(timeout=5)
  303. assert_equal(value, 1)
  304. assert_equal(root_id, expected_root_id)
  305. assert_equal(parent_id, base_parent)
  306. @testcase('all', 'green')
  307. def manyshort(self):
  308. self.join(group(add.s(i, i) for i in range(1000))(),
  309. timeout=10, propagate=True)
  310. @testcase('all', 'green', iterations=1)
  311. def unicodetask(self):
  312. self.join(group(print_unicode.s() for _ in range(5))(),
  313. timeout=1, propagate=True)
  314. @testcase('all')
  315. def always_timeout(self):
  316. self.join(
  317. group(sleeping.s(1).set(time_limit=0.1)
  318. for _ in range(100))(),
  319. timeout=10, propagate=True,
  320. )
  321. @testcase('all')
  322. def termbysig(self):
  323. self._evil_groupmember(kill)
  324. @testcase('green')
  325. def group_with_exit(self):
  326. self._evil_groupmember(exiting)
  327. @testcase('all')
  328. def timelimits(self):
  329. self._evil_groupmember(sleeping, 2, time_limit=1)
  330. @testcase('all')
  331. def timelimits_soft(self):
  332. self._evil_groupmember(sleeping_ignore_limits, 2,
  333. soft_time_limit=1, time_limit=1.1)
  334. @testcase('all')
  335. def alwayskilled(self):
  336. g = group(kill.s() for _ in range(10))
  337. self.join(g(), timeout=10)
  338. @testcase('all', 'green')
  339. def alwaysexits(self):
  340. g = group(exiting.s() for _ in range(10))
  341. self.join(g(), timeout=10)
  342. def _evil_groupmember(self, evil_t, *eargs, **opts):
  343. g1 = group(add.s(2, 2).set(**opts), evil_t.s(*eargs).set(**opts),
  344. add.s(4, 4).set(**opts), add.s(8, 8).set(**opts))
  345. g2 = group(add.s(3, 3).set(**opts), add.s(5, 5).set(**opts),
  346. evil_t.s(*eargs).set(**opts), add.s(7, 7).set(**opts))
  347. self.join(g1(), timeout=10)
  348. self.join(g2(), timeout=10)
  349. @testcase('all', 'green')
  350. def bigtasksbigvalue(self):
  351. g = group(any_returning.s(BIG, sleep=0.3) for i in range(8))
  352. r = g()
  353. try:
  354. self.join(r, timeout=10)
  355. finally:
  356. # very big values so remove results from backend
  357. try:
  358. r.forget()
  359. except NotImplementedError:
  360. pass
  361. @testcase('all', 'green')
  362. def bigtasks(self, wait=None):
  363. self._revoketerm(wait, False, False, BIG)
  364. @testcase('all', 'green')
  365. def smalltasks(self, wait=None):
  366. self._revoketerm(wait, False, False, SMALL)
  367. @testcase('all')
  368. def revoketermfast(self, wait=None):
  369. self._revoketerm(wait, True, False, SMALL)
  370. @testcase('all')
  371. def revoketermslow(self, wait=5):
  372. self._revoketerm(wait, True, True, BIG)
  373. def _revoketerm(self, wait=None, terminate=True,
  374. joindelay=True, data=BIG):
  375. g = group(any_.s(data, sleep=wait) for i in range(8))
  376. r = g()
  377. if terminate:
  378. if joindelay:
  379. sleep(random.choice(range(4)))
  380. r.revoke(terminate=True)
  381. self.join(r, timeout=10)