summaryrefslogtreecommitdiff
path: root/src/server
diff options
context:
space:
mode:
authorLV-426 <lv-426@taproot.org.il>2015-04-09 01:29:13 +0300
committerLV-426 <lv-426@taproot.org.il>2015-04-09 02:03:45 +0300
commit033681306c81cdc237db6fd61727117b635fab73 (patch)
tree470b9d8ef06232355d6ba68e4c1b60d7cd717fff /src/server
parentd32bd9d94d710c708484dc8918070fa906e891cb (diff)
rz_kernel: add ThreadPoolExecutor, kernel_heartbeat() + concurrent.futures dependency
Diffstat (limited to 'src/server')
-rw-r--r--src/server/rz_kernel.py39
-rw-r--r--src/server/rz_server.py8
2 files changed, 42 insertions, 5 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):
"""
diff --git a/src/server/rz_server.py b/src/server/rz_server.py
index 3232b18e..e29259b8 100644
--- a/src/server/rz_server.py
+++ b/src/server/rz_server.py
@@ -278,6 +278,8 @@ def init_webapp(cfg, kernel, db_ctl=None):
Initialize webapp:
- call init_rest_interface()
"""
+ global webapp
+
root_path = cfg.root_path
assert os.path.exists(root_path), "root path doesn't exist: %s" % root_path
@@ -328,8 +330,9 @@ def init_signal_handlers():
signal.signal(signal.SIGTERM, signal_handler__exit)
def shutdown():
- user_db.shutdown()
log.info('rz_server: shutting down')
+ user_db.shutdown()
+ webapp.kernel.shutdown()
if __name__ == "__main__":
@@ -363,7 +366,10 @@ if __name__ == "__main__":
log.info('failed initialization, aborting')
exit(-1)
+ # setup kernel
kernel = RZ_Kernel()
+ kernel.start()
+
webapp = init_webapp(cfg, kernel)
ws_srv = init_ws_interface(cfg, kernel, webapp)