Hide keyboard shortcuts

Hot-keys on this page

r m x p   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

1# -*- coding: utf-8 -*- 

2from __future__ import annotations 

3 

4# -- stdlib -- 

5from typing import Any, Callable, Dict 

6import logging 

7 

8# -- third party -- 

9from gevent.event import AsyncResult, Event 

10from gevent.greenlet import Greenlet 

11from gevent.pool import Pool 

12import gevent 

13 

14# -- own -- 

15 

16# -- code -- 

17log = logging.getLogger('core') 

18 

19 

20class Core(object): 

21 core_type = '' 

22 _auto_id = 0 

23 

24 runner: CoreRunner 

25 

26 def __init__(self): 

27 self._auto_id = Core._auto_id 

28 Core._auto_id += 1 

29 

30 self._result = AsyncResult() 

31 self.tasks: Dict[str, Callable[[], None]] = {} 

32 

33 def __repr__(self) -> str: 

34 return f'Core[{self.core_type}{self._auto_id}]' 

35 

36 @property 

37 def result(self): 

38 return self._result 

39 

40 def crash(self, e): 

41 self._result.set_exception(e) 

42 

43 

44class CoreCrashed(Exception): 

45 pass 

46 

47 

48class CoreRunner(object): 

49 def __init__(self, core: Core, paranoid: bool = False): 

50 self.core = core 

51 self.pool = Pool() 

52 self.ready = Event() 

53 self.tasks: Dict[str, Greenlet] = {} 

54 

55 self._paranoid = paranoid 

56 

57 def run(self) -> Any: 

58 core = self.core 

59 

60 core.runner = self 

61 

62 try: 

63 for k, f in core.tasks.items(): 

64 log.debug('Spawning task [%s] for core %s', k, core) 

65 gr = self.pool.spawn(f) 

66 gr.gr_name = k 

67 self.tasks[k] = gr 

68 

69 self.ready.set() 

70 try: 

71 return core.result.get() 

72 except BaseException as e: 

73 raise CoreCrashed(f'{core} crashed') from e 

74 finally: 

75 self.shutdown() 

76 

77 def spawn(self, fn, *args, **kw): 

78 core = self.core 

79 gr = self.pool.spawn(fn, *args, **kw) 

80 gr.gr_name = f'{repr(self.core)}/{fn.__qualname__}' 

81 if self._paranoid: 

82 gr.link_exception(lambda gr: core.crash(gr.exception)) 82 ↛ exitline 82 didn't run the lambda on line 82

83 return gr 

84 

85 def start(self, gr): 

86 core = self.core 

87 self.pool.start(gr) 

88 if self._paranoid: 

89 gr.link_exception(lambda gr: core.crash(gr.exception)) 

90 return gr 

91 

92 def sleep(self, t): 

93 gevent.sleep(t) 

94 

95 def idle(self, prio=0): 

96 gevent.idle(prio) 

97 

98 def shutdown(self) -> None: 

99 self.pool.kill()