mirror of
https://github.com/corpnewt/gibMacOS.git
synced 2024-04-21 12:31:48 +00:00
Queue idling fixes
This commit is contained in:
+25
-31
@@ -1,7 +1,4 @@
|
||||
import sys
|
||||
import subprocess
|
||||
import threading
|
||||
import shlex
|
||||
import sys, subprocess, time, threading, shlex
|
||||
try:
|
||||
from Queue import Queue, Empty
|
||||
except:
|
||||
@@ -22,9 +19,16 @@ class Run:
|
||||
pass
|
||||
pipe.close()
|
||||
|
||||
def _create_thread(self, output):
|
||||
# Creates a new queue and thread object to watch based on the output pipe sent
|
||||
q = Queue()
|
||||
t = threading.Thread(target=self._read_output, args=(output, q))
|
||||
t.daemon = True
|
||||
return (q,t)
|
||||
|
||||
def _stream_output(self, comm, shell = False):
|
||||
output = error = ""
|
||||
p = ot = et = None
|
||||
p = None
|
||||
try:
|
||||
if shell and type(comm) is list:
|
||||
comm = " ".join(shlex.quote(x) for x in comm)
|
||||
@@ -32,57 +36,47 @@ class Run:
|
||||
comm = shlex.split(comm)
|
||||
p = subprocess.Popen(comm, shell=shell, stdout=subprocess.PIPE, stderr=subprocess.PIPE, bufsize=0, universal_newlines=True, close_fds=ON_POSIX)
|
||||
# Setup the stdout thread/queue
|
||||
q = Queue()
|
||||
t = threading.Thread(target=self._read_output, args=(p.stdout, q))
|
||||
t.daemon = True # thread dies with the program
|
||||
# Setup the stderr thread/queue
|
||||
qe = Queue()
|
||||
te = threading.Thread(target=self._read_output, args=(p.stderr, qe))
|
||||
te.daemon = True # thread dies with the program
|
||||
q,t = self._create_thread(p.stdout)
|
||||
qe,te = self._create_thread(p.stderr)
|
||||
# Start both threads
|
||||
t.start()
|
||||
te.start()
|
||||
|
||||
while True:
|
||||
c = z = ""
|
||||
try:
|
||||
c = q.get_nowait()
|
||||
except Empty:
|
||||
pass
|
||||
try: c = q.get_nowait()
|
||||
except Empty: pass
|
||||
else:
|
||||
sys.stdout.write(c)
|
||||
output += c
|
||||
sys.stdout.flush()
|
||||
try:
|
||||
z = qe.get_nowait()
|
||||
except Empty:
|
||||
pass
|
||||
try: z = qe.get_nowait()
|
||||
except Empty: pass
|
||||
else:
|
||||
sys.stderr.write(z)
|
||||
error += z
|
||||
sys.stderr.flush()
|
||||
if not c==z=="": continue # Keep going until empty
|
||||
# No output - see if still running
|
||||
p.poll()
|
||||
if c==z=="" and p.returncode != None:
|
||||
if p.returncode != None:
|
||||
# Subprocess ended
|
||||
break
|
||||
# No output, but subprocess still running - stall for 20ms
|
||||
time.sleep(0.02)
|
||||
|
||||
o, e = p.communicate()
|
||||
ot.exit()
|
||||
et.exit()
|
||||
return (output+o, error+e, p.returncode)
|
||||
except:
|
||||
if ot or et:
|
||||
try: ot.exit()
|
||||
except: pass
|
||||
try: et.exit()
|
||||
except: pass
|
||||
if p:
|
||||
return (output, error, p.returncode)
|
||||
o, e = p.communicate()
|
||||
return (output+o, error+e, p.returncode)
|
||||
return ("", "Command not found!", 1)
|
||||
|
||||
def _decode(self, value):
|
||||
def _decode(self, value, encoding="utf-8", errors="ignore"):
|
||||
# Helper method to only decode if bytes type
|
||||
if sys.version_info >= (3,0) and isinstance(value, bytes):
|
||||
return value.decode("utf-8","ignore")
|
||||
return value.decode(encoding,errors)
|
||||
return value
|
||||
|
||||
def _run_command(self, comm, shell = False):
|
||||
|
||||
Reference in New Issue
Block a user