Repository navigation
Expand file tree
/
Copy pathrun.py
More file actions
202 lines (171 loc) · 7.29 KB
/
Copy pathrun.py
File metadata and controls
202 lines (171 loc) · 7.29 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
from __future__ import annotations
import shlex
import sys
from io import BytesIO
from subprocess import PIPE, STDOUT, Popen
from refinery.lib.meta import metavars
from refinery.lib.structures import MemoryFile
from refinery.lib.types import Param
from refinery.units import Arg, RefineryPartialResult, Unit
class run(Unit):
"""
Turns any other shell command into a (frame-compatible) refinery unit.
Using this unit is only required when you require framing features (variables or multi-chunk
processing) for running an external program.
Data is processed by feeding it to the standard input of a process spawned from the given
command line, and then reading the standard output of that process as the result of the
operation. The main purpose of this unit is to allow using the syntax from `refinery.lib.frame`
with other command line tools. By default, the unit streams the output from the executed
command as individual outputs, but the `buffer` option can be set to buffer all output of a
single execution. The format string expression `{}` or `{0}` can be used as one of the
arguments passed to the external command to represent the incoming data. In this case, the data
will not be sent to the standard input device of the new process.
"""
_JOIN_TIME = 2
_WAIT_TIME = 0.01
def __init__(
self, *commandline: Param[str, Arg.String(nargs='...', metavar='(all remaining)', help=(
'All remaining command line tokens form an arbitrary command line to be executed. Use'
' format string syntax to insert meta variables and incoming data chunks.'))],
stream: Param[bool, Arg.Switch('-s',
help='Stream the command output rather than buffering it.')] = False,
noinput: Param[bool, Arg.Switch('-x', help='Do not send any input to the new process.')] = False,
errors: Param[bool, Arg.Switch('-m', help=(
'Merge stdout and stderr. By default, the standard error stream of the coupled command'
' is forwarded to the logger, i.e. it is only visible if -v is also specified.'
))] = False,
timeout: Param[float, Arg.Double('-t', metavar='T', help=(
'Optionally set an execution timeout as a floating point number in seconds.'
))] = 0.0
):
if not commandline:
raise ValueError('you need to provide a command line.')
super().__init__(
commandline=commandline, errors=errors, noinput=noinput, stream=stream, timeout=timeout)
def process(self, data):
meta = metavars(data)
used = set()
commandline = [
meta.format_str(cmd, self.codec, [data], None, used=used)
for cmd in self.args.commandline
]
if self.args.noinput:
self.log_info('sending no input to process stdin')
data = None
stream: bool = self.args.stream
merge: bool = self.args.errors
timeout: int = self.args.timeout
if posix := 'posix' in sys.builtin_module_names:
args = shlex.join(commandline)
else:
args = ' '.join(F'"{a}"' for a in commandline)
self.log_info(args)
process = Popen(args, shell=True,
stdin=PIPE, stdout=PIPE, stderr=STDOUT if merge else PIPE, close_fds=posix)
if not stream and not timeout and not merge:
out, err = process.communicate(data)
for line in err.splitlines():
self.log_info(line)
yield out
return
from queue import Empty, Queue
from threading import Event, Thread
from time import monotonic, sleep
start = 0
result = None
_jt = self._JOIN_TIME
_wt = self._WAIT_TIME
qerr: Queue[bytes] = Queue()
qout: Queue[bytes] = Queue()
done = Event()
def adapter(stream: BytesIO, queue: Queue[bytes], event: Event):
while not event.is_set():
out = stream.read1()
if out:
queue.put(out)
else:
break
stream.close()
recvout = Thread(target=adapter, args=(process.stdout, qout, done), daemon=True)
recvout.start()
if not merge:
recverr = Thread(target=adapter, args=(process.stderr, qerr, done), daemon=True)
recverr.start()
else:
recverr = None
if stdin := process.stdin:
if data:
stdin.write(data)
stdin.close()
start = monotonic()
if not stream or timeout:
result = MemoryFile()
def queue_read(q: Queue[bytes]):
try:
return q.get_nowait()
except Empty:
return None
errbuf = MemoryFile()
errobj = None
while True:
out = queue_read(qout)
err = None
if not merge:
err = queue_read(qerr)
if err and self.log_info():
errbuf.write(err)
errbuf.seek(0)
lines = errbuf.readlines()
errbuf.seek(0)
errbuf.truncate()
if lines:
if not (done.is_set() or lines[~0].endswith(B'\n')):
errbuf.write(lines.pop())
for line in lines:
if line := line.rstrip(B'\n'):
self.log_info(line)
if out:
if not stream or timeout:
if result is not None:
result.write(out)
if stream:
yield out
if done.is_set():
if recverr is not None and recverr.is_alive():
self.log_warn('stderr receiver thread zombied')
if recvout.is_alive():
self.log_warn('stdout receiver thread zombied')
break
elif not err and not out:
if process.poll() is None:
sleep(_wt)
else:
if recverr is not None:
recverr.join(_jt)
recvout.join(_jt)
done.set()
elif timeout:
assert result is not None
if monotonic() - start > timeout:
self.log_info('terminating process after timeout expired')
done.set()
process.terminate()
for wait in range(4):
if process.poll() is not None:
break
sleep(_wt)
else:
self.log_warn('process termination may have failed')
if recverr is not None:
recverr.join(_jt)
recvout.join(_jt)
if not len(result):
errobj = RuntimeError('timeout reached, process had no output')
else:
errobj = RefineryPartialResult(
'timeout reached, returning all collected output',
partial=result.getvalue())
if errobj is not None:
raise errobj
if result is not None:
yield result.getvalue()