Implement publish/subscribe

This commit is contained in:
Andika Wasisto 2020-05-25 12:35:08 +07:00
parent 5cf90c372f
commit e156c778b9
2 changed files with 82 additions and 15 deletions

View file

@ -36,13 +36,49 @@ Usage Example
ws://localhost:62456/ws ws://localhost:62456/ws
> RANDINT32 > RANDINT32
< 585865374 < 585865374
> RANDUNIFORM > RANDUNIFORM
< 0.70137183786 < 0.70137183786
> RANDNORMAL > RANDNORMAL
< -1.6120135370 < -1.6120135370
> RANDBYTES 16 > RANDBYTES 16
< <20><><EFBFBD><EFBFBD><EFBFBD>9L).<2E>!:<3A><><EFBFBD>@Nh<EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>d<EFBFBD><EFBFBD><EFBFBD><EFBFBD>_q<EFBFBD> < <20><><EFBFBD><EFBFBD><EFBFBD>9L).<2E>!:<3A><><EFBFBD>@Nh<EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>d<EFBFBD><EFBFBD><EFBFBD><EFBFBD>_q<EFBFBD>
> SUBSCRIBEINT32
< 1361330636
< -604581511
< 1510923919
< ...
> UNSUBSCRIBE
< UNSUBSCRIBED
> SUBSCRIBEUNIFORM
< 0.54623951886
< 0.67567578799
< 0.09746421443
< ...
> UNSUBSCRIBE
< UNSUBSCRIBED
> SUBSCRIBENORMAL
< -1.6120135370
< 0.02943381135
< -0.9883458007
< ...
> UNSUBSCRIBE
< UNSUBSCRIBED
> SUBSCRIBEBYTES 8
< <20><><EFBFBD><EFBFBD><EFBFBD>9L)
< .<EFBFBD>!:<EFBFBD><EFBFBD><EFBFBD>@
< Nh<EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
< d<EFBFBD><EFBFBD><EFBFBD><EFBFBD>_q<EFBFBD>
< ...
> UNSUBSCRIBE
< UNSUBSCRIBED
License License
------- -------

View file

@ -18,6 +18,8 @@
# OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE # OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
# SOFTWARE. # SOFTWARE.
import threading
from flask import Flask, request, Response from flask import Flask, request, Response
from flask_sockets import Sockets from flask_sockets import Sockets
from gevent import pywsgi from gevent import pywsgi
@ -59,22 +61,51 @@ def api_randbytes():
@sockets.route('/ws') @sockets.route('/ws')
def ws(websocket): def ws(websocket):
subscribed = [False]
while not websocket.closed: while not websocket.closed:
try: threading.Thread(target=handle_ws_message, args=(websocket.receive(), websocket, subscribed)).start()
message = websocket.receive().strip().lower()
split_message = message.split(' ')
if split_message[0] == 'randint32': def handle_ws_message(message, websocket, subscribed):
websocket.send(str(qng_wrapper.randint32())) try:
elif split_message[0] == 'randuniform': split_message = message.strip().lower().split()
websocket.send(str(qng_wrapper.randuniform())) if split_message[0] == 'randint32':
elif split_message[0] == 'randnormal': websocket.send(str(qng_wrapper.randint32()))
websocket.send(str(qng_wrapper.randnormal())) elif split_message[0] == 'randuniform':
elif split_message[0] == 'randbytes': websocket.send(str(qng_wrapper.randuniform()))
websocket.send(qng_wrapper.randbytes(int(split_message[1]))) elif split_message[0] == 'randnormal':
except ValueError as e: websocket.send(str(qng_wrapper.randnormal()))
websocket.close(code=1003, message=str(e)) elif split_message[0] == 'randbytes':
except Exception as e: length = int(split_message[1])
websocket.close(code=1011, message=str(e)) websocket.send(qng_wrapper.randbytes(length))
elif split_message[0] == 'subscribeint32':
if not subscribed[0]:
subscribed[0] = True
while subscribed[0] and not websocket.closed:
websocket.send(str(qng_wrapper.randint32()))
elif split_message[0] == 'subscribeuniform':
if not subscribed[0]:
subscribed[0] = True
while subscribed[0] and not websocket.closed:
websocket.send(str(qng_wrapper.randuniform()))
elif split_message[0] == 'subscribenormal':
if not subscribed[0]:
subscribed[0] = True
while subscribed[0] and not websocket.closed:
websocket.send(str(qng_wrapper.randnormal()))
elif split_message[0] == 'subscribebytes':
chunk = int(split_message[1])
if not subscribed[0]:
subscribed[0] = True
while subscribed[0] and not websocket.closed:
websocket.send(qng_wrapper.randbytes(chunk))
elif split_message[0] == 'unsubscribe':
subscribed[0] = False
websocket.send('UNSUBSCRIBED')
except (ValueError, BlockingIOError):
pass
except Exception as e:
websocket.close(code=1011, message=str(e))
@app.errorhandler(Exception) @app.errorhandler(Exception)