diff options
| author | LV-426 <lv-426@taproot.org.il> | 2015-04-09 01:29:13 +0300 |
|---|---|---|
| committer | LV-426 <lv-426@taproot.org.il> | 2015-04-09 02:03:45 +0300 |
| commit | 033681306c81cdc237db6fd61727117b635fab73 (patch) | |
| tree | 470b9d8ef06232355d6ba68e4c1b60d7cd717fff /src/server/rz_kernel.py | |
| parent | d32bd9d94d710c708484dc8918070fa906e891cb (diff) | |
rz_kernel: add ThreadPoolExecutor, kernel_heartbeat() + concurrent.futures dependency
Diffstat (limited to 'src/server/rz_kernel.py')
| -rw-r--r-- | src/server/rz_kernel.py | 39 |
1 files changed, 35 insertions, 4 deletions
diff --git a/src/server/rz_kernel.py b/src/server/rz_kernel.py index 53749595..060c22b3 100644 --- a/src/server/rz_kernel.py +++ b/src/server/rz_kernel.py @@ -1,19 +1,21 @@ """ Rhizi kernel, home to core operation login """ -import json +from functools import wraps import logging import traceback +from concurrent.futures import ThreadPoolExecutor from db_op import DBO_diff_commit__attr, DBO_block_chain__commit, DBO_rzdoc__create, \ - DBO_rzdoc__lookup_by_name, DBO_rzdoc__clone, DBO_rzdoc__delete, DBO_rzdoc__list,\ + DBO_rzdoc__lookup_by_name, DBO_rzdoc__clone, DBO_rzdoc__delete, DBO_rzdoc__list, \ DBO_block_chain__init from db_op import DBO_diff_commit__topo from model.graph import Topo_Diff from model.model import RZDoc -import neo4j_util from neo4j_qt import QT_RZDOC_NS_Filter, QT_RZDOC_Meta_NS_Filter - +import neo4j_util +from collections import defaultdict +import time log = logging.getLogger('rhizi') @@ -79,6 +81,7 @@ def for_all_public_functions(decorator): class RZ_Kernel(object): """ RZ kernel: + - activation/shutdown via start(), shutdown() - all public methods decorated with deco__exception_log """ @@ -86,6 +89,34 @@ class RZ_Kernel(object): self.db_ctl = None self.rzdoc_reader_assoc_map = defaultdict(list) self.cache__rzdoc_name_to_rzdoc = {} + self.should_stop = False + self.heartbeat_period_sec = 5 + + def start(self): + + def kernel_heartbeat(): + """ + Handle periodic server tasks: + - manage rzdoc subscriber lists + """ + while False == self.should_stop: + + for rzdoc, r_assoc_set in self.rzdoc_reader_assoc_map.items(): + for r_assoc in r_assoc_set: + if r_assoc.err_count__IO > 3: + r_assoc_set.remove(r_assoc) + log.info('rz_kernel: evicting reader: IO error count exceeded limit: remote-addr: %s, rzdoc: %s' % (r_assoc.remote_socket_addr, + rzdoc.name)) + time.sleep(self.heartbeat_period_sec) + + self.executor = ThreadPoolExecutor(max_workers=8) + self.executor.submit(kernel_heartbeat) + log.info('rz_kernel: on-line') + + def shutdown(self): + self.should_stop = True + self.executor.shutdown() + log.info('rz_kernel: shutting down') def cache_lookup__rzdoc(self, rzdoc_name): """ |
