Преглед изворни кода

切换 JoinableQueue 为基于 Redis 队列的 IPC

James Iter пре 8 година
родитељ
комит
7a98668b2e
4 измењених фајлова са 21 додато и 10 уклоњено
  1. 1 0
      docs/todo.md
  2. 10 3
      main.py
  3. 5 4
      models/initialize.py
  4. 5 3
      views/guest.py

+ 1 - 0
docs/todo.md

@@ -69,5 +69,6 @@
 - [ ] 考虑热更新、升级的问题
 - [ ] 加入清除老旧日志的功能
 - [ ] 忽略 guestfs- 为首的日志
+- [ ] 安装脚本、文档加入时间同步环节
 
 

+ 10 - 3
main.py

@@ -18,7 +18,7 @@ from werkzeug.debug import get_current_traceback
 
 from models import Utils
 from models.event_processor import EventProcessor
-from models.initialize import logger, q_ws, Init
+from models.initialize import logger, Init
 import api_route_table
 import views_route_table
 from models import Database as db
@@ -86,8 +86,16 @@ def instantiation_ws_vnc(listen_port, target_host, target_port):
 
 
 def ws_engine_for_vnc():
+
+    logger.info(msg='VNC ws engine is launched.')
+
     while True:
-        payload = q_ws.get()
+        payload = db.r.lpop(app.config['ipc'])
+
+        if payload is None:
+            time.sleep(1)
+            continue
+
         payload = json.loads(payload)
 
         c_pid = os.fork()
@@ -97,7 +105,6 @@ def ws_engine_for_vnc():
         # 因为 WebSocketProxy 使用了 daemon 参数,所以当执行到 ws.start_server() 时,会退出其所在的子进程,
         # 故而这里设置wait来处理结束的子进程的环境,避免出现僵尸进程。
         os.wait()
-        q_ws.task_done()
 
 
 def is_not_need_to_auth(endpoint):

+ 5 - 4
models/initialize.py

@@ -3,7 +3,6 @@
 
 
 import traceback
-from multiprocessing import JoinableQueue
 from flask import Flask
 import logging
 from logging.handlers import TimedRotatingFileHandler
@@ -45,6 +44,7 @@ class Init(object):
         'vnc_port_used_set': 'S:VNCPort:Used',
         'downstream_queue': 'Q:Downstream',
         'upstream_queue': 'Q:Upstream',
+        'ipc_queue': 'Q:IPC',
         'hosts_info': 'H:HostsInfo',
         'compute_nodes_hostname_key': 'S:ComputeNodesHostname',
         'guest_boot_jobs': 'S:GuestBootJobs',
@@ -135,10 +135,11 @@ class Init(object):
         from models import Database as db
         from models import Utils
 
+        logger.info(msg='PS PING PONG engine is launched.')
         while True:
             try:
                 if Utils.exit_flag:
-                    msg = 'Thread pub_sub_ping_pong say bye-bye'
+                    msg = 'Thread PS PING PONG engine say bye-bye'
                     print msg
                     logger.info(msg=msg)
                     return
@@ -157,10 +158,11 @@ class Init(object):
         already_clear = False
         the_time = '03:30'
 
+        logger.info(msg='Clear expire log monitor is launched.')
         while True:
             try:
                 if Utils.exit_flag:
-                    msg = 'Thread clear_expire_monitor_log say bye-bye'
+                    msg = 'Thread clear expire monitor log say bye-bye'
                     print msg
                     logger.info(msg=msg)
                     return
@@ -189,7 +191,6 @@ class Init(object):
                 logger.error(traceback.format_exc())
 
 
-q_ws = JoinableQueue()
 # 预编译效率更高
 regex_sql_str = re.compile('\\\+"')
 regex_dsl_str = re.compile('^\w+:\w+:[\S| ]+$')

+ 5 - 3
views/guest.py

@@ -9,7 +9,9 @@ import requests
 from math import ceil
 import re
 import socket
-from models.initialize import q_ws
+import time
+from models import Database as db
+from models.initialize import config
 
 
 __author__ = 'James Iter'
@@ -192,8 +194,8 @@ def vnc(uuid):
 
     payload = {'listen_port': port, 'target_host': guest_ret['data']['on_host'],
                'target_port': guest_ret['data']['vnc_port']}
-    q_ws.put(json.dumps(payload, ensure_ascii=False))
-    q_ws.join()
+    db.r.rpush(config['ipc'], json.dumps(payload, ensure_ascii=False))
+    time.sleep(1)
 
     return render_template('vnc_lite.html', port=port, password=guest_ret['data']['vnc_password'])