summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorLV-426 <lv-426@taproot.org.il>2015-04-09 01:33:41 +0300
committerLV-426 <lv-426@taproot.org.il>2015-04-09 02:03:45 +0300
commitb246145698eef78bc1df67eba933ca722ea0dbdc (patch)
treeb02bb550c8fa122a57e523eb4bada7e514c801b8
parent660e44838648b1aff0a4b3b2c755df2661fa8b93 (diff)
rz_api_websocket: ws reader subscription messages, ack & nak
-rw-r--r--src/server/rz_api_websocket.py70
1 files changed, 55 insertions, 15 deletions
diff --git a/src/server/rz_api_websocket.py b/src/server/rz_api_websocket.py
index 4412476f..48c62fb0 100644
--- a/src/server/rz_api_websocket.py
+++ b/src/server/rz_api_websocket.py
@@ -7,8 +7,8 @@ import logging
from socketio.mixins import BroadcastMixin
from socketio.namespace import BaseNamespace
import traceback
-
from model.graph import Attr_Diff, Topo_Diff
+from rz_kernel import RZDoc_Exception__not_found
log = logging.getLogger('rhizi')
@@ -27,27 +27,53 @@ class WebSocket_Graph_NS(BaseNamespace, BroadcastMixin):
# FIXME: impl
pass
- def multicast_msg(self, msg_name, *args):
- self.socket.server.log_multicast(msg_name)
- try:
- super(WebSocket_Graph_NS, self).broadcast_event_not_me(msg_name, *args)
- except Exception as e:
- log.error(e.message)
- log.error(traceback.print_exc())
-
def _log_conn(self, prefix_msg):
rmt_addr = self.environ['REMOTE_ADDR']
rmt_port = self.environ['REMOTE_PORT']
sid = self.environ['socketio'].sessid
log.info('ws: %s: sid: %s, remote-socket: %s:%s' % (prefix_msg, sid, rmt_addr, rmt_port))
- def recv_connect(self):
- self._log_conn('conn open')
- super(WebSocket_Graph_NS, self).recv_connect() # super called despite being empty
+ def _on_rzdoc_subscribe_common(self, data_dict, is_subscribe=None):
+ rzdoc_name_raw = data_dict['rzdoc_name']
- def recv_disconnect(self):
- self._log_conn('conn close')
- super(WebSocket_Graph_NS, self).recv_disconnect()
+ # FIXME: non-flask dep. sanitization
+ # rzdoc_name = sanitize_input__rzdoc_name(rzdoc_name_raw)
+ rzdoc_name = rzdoc_name_raw
+
+ rmt_addr = self.environ['REMOTE_ADDR']
+ rmt_port = self.environ['REMOTE_PORT']
+ remote_socket_addr = (rmt_addr, rmt_port)
+ socket = self.socket
+
+ kernel = self.request.kernel
+ msg_name = 'rzdoc_subscribe' if is_subscribe else 'rzdoc_unsubscribe'
+ try:
+ if is_subscribe:
+ kernel.rzdoc__reader_subscribe(remote_socket_addr=remote_socket_addr,
+ rzdoc_name=rzdoc_name,
+ socket=socket)
+ self.ack(msg_name)
+ else:
+ kernel.rzdoc__reader_unsubscribe(remote_socket_addr=remote_socket_addr,
+ rzdoc_name=rzdoc_name,
+ socket=socket)
+ self.ack(msg_name)
+ except RZDoc_Exception__not_found:
+ self.nak(msg_name)
+
+ def ack(self, acked_msg_name):
+ self.emit('ack', acked_msg_name)
+
+ def nak(self, acked_msg_name):
+ self.emit('nak', acked_msg_name)
+
+ def multicast_msg(self, msg_name, *args):
+ self.socket.server.log_multicast(msg_name)
+ try:
+ super(WebSocket_Graph_NS, self).broadcast_event_not_me(msg_name, *args)
+ except Exception as e:
+ log.error(e.message)
+ log.error(traceback.print_exc())
def on_diff_commit__topo(self, json_data):
json_dict = json.loads(json_data)
@@ -78,3 +104,17 @@ class WebSocket_Graph_NS(BaseNamespace, BroadcastMixin):
# commit_ret may not be the same
return self.multicast_msg('diff_commit__attr', attr_diff, commit_ret)
+ def on_rzdoc_subscribe(self, data_dict):
+ return self._on_rzdoc_subscribe_common(data_dict, is_subscribe=True)
+
+ def on_rzdoc_unsubscribe(self, data_dict):
+ return self._on_rzdoc_subscribe_common(data_dict, is_subscribe=False)
+
+ def recv_connect(self):
+ self._log_conn('conn open')
+ super(WebSocket_Graph_NS, self).recv_connect() # super called despite being empty
+
+ def recv_disconnect(self):
+ self._log_conn('conn close')
+ super(WebSocket_Graph_NS, self).recv_disconnect()
+