Skip to content

Commit b437fb0

Browse files
ctruedenclaude
andcommitted
Keep Python workers alive and report them as tasks
Each environment now gets one resident Python worker, reused across script runs, so imported modules stay imported and objects passed to task.export(...) remain available to later runs. Each run still gets a fresh namespace. Environment builds and script runs appear as SciJava tasks. Canceling a run's task asks the Python task to cancel; if the script has not stopped within a few seconds, its worker is killed, and the next run starts a new one. Supporting changes: * Runs on the same worker are serialized, so that worker output can be attributed to the script that produced it. The wrapper writes an end marker to stderr, so that late stderr lines are not lost or misrouted. * Image outputs are now unlinked after conversion. By Appose convention the service side frees shared memory, and with the worker staying alive, nothing else ever would. The reusable pieces live in a new _internal package, as a prototype of functionality destined for Appose core (BuildListener, ResidentWorker) and a future scijava-appose component (ResidentWorkerService, SciJavaTasks). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 5845ac2 commit b437fb0

9 files changed

Lines changed: 940 additions & 91 deletions

File tree

‎README.md‎

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,31 @@ restarts. Environments are stored in the Appose environments directory
4848
Scripts that share a configuration file share an environment. Editing the file
4949
causes the environment to be updated on the next run.
5050

51-
Each script run starts a fresh Python process.
51+
Builds appear in the application's task list, with their progress, so a
52+
long first build does not look like a hang.
53+
54+
## Python workers
55+
56+
Each environment gets one Python worker process, which stays alive between
57+
script runs. Modules a script imports stay imported, so a second run of a
58+
script that imports e.g. PyTorch starts quickly. Each run still gets a fresh
59+
namespace: variables from one run are not visible to the next.
60+
61+
To keep something expensive, such as a loaded model, across runs, hand it
62+
to `task.export`:
63+
64+
```python
65+
if "model" not in globals():
66+
model = load_model()
67+
task.export(model=model)
68+
```
69+
70+
Runs on the same environment execute one at a time. Each run appears in the
71+
application's task list, where it can be canceled. A script can notice
72+
cancelation by checking `task.cancel_requested`; if it has not stopped a few
73+
seconds after being canceled, its worker process is stopped, and the next
74+
run starts a new one. The worker, and whatever memory (including GPU
75+
memory) it holds, is otherwise released when the application exits.
5276

5377
## Inputs and outputs
5478

‎src/main/java/org/scijava/plugins/scripting/appose/python/ApposePythonScriptEngine.java‎

Lines changed: 126 additions & 88 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,11 @@
4343
import java.util.HashMap;
4444
import java.util.List;
4545
import java.util.Map;
46-
import java.util.concurrent.ConcurrentHashMap;
46+
import java.util.UUID;
47+
import java.util.concurrent.CountDownLatch;
48+
import java.util.concurrent.TimeUnit;
49+
import java.util.concurrent.locks.Lock;
50+
import java.util.function.Consumer;
4751

4852
import javax.script.Bindings;
4953
import javax.script.ScriptEngine;
@@ -55,19 +59,24 @@
5559
import org.apposed.appose.Appose;
5660
import org.apposed.appose.BuildException;
5761
import org.apposed.appose.Builder;
58-
import org.apposed.appose.Environment;
5962
import org.apposed.appose.NDArray;
6063
import org.apposed.appose.Service;
64+
import org.apposed.appose.Service.TaskStatus;
6165
import org.apposed.appose.TaskException;
6266
import org.scijava.Context;
6367
import org.scijava.app.StatusService;
6468
import org.scijava.convert.ConvertService;
6569
import org.scijava.log.LogService;
6670
import org.scijava.module.ModuleItem;
6771
import org.scijava.plugin.Parameter;
72+
import org.scijava.plugins.scripting.appose.python._internal.ResidentWorker;
73+
import org.scijava.plugins.scripting.appose.python._internal.ResidentWorkerService;
74+
import org.scijava.plugins.scripting.appose.python._internal.SciJavaTasks;
6875
import org.scijava.script.AbstractScriptEngine;
6976
import org.scijava.script.ScriptInfo;
7077
import org.scijava.script.ScriptModule;
78+
import org.scijava.task.Task;
79+
import org.scijava.task.TaskService;
7180

7281
/**
7382
* A script engine for Python (CPython, not Jython!), backed by
@@ -79,8 +88,12 @@
7988
* #@script(env="myenv.toml", scheme="pixi.toml")
8089
* </pre>
8190
* <p>
82-
* The engine lazily builds the environment on first use and caches it for
83-
* subsequent calls. Inputs of array-compatible types (e.g., ImgLib2 {@code Img})
91+
* The engine lazily builds the environment on first use, and keeps one Python
92+
* worker process per environment alive across script runs, so that imported
93+
* modules, and objects a script keeps via {@code task.export(...)}, need not
94+
* be loaded again. Each run still gets a fresh namespace. Builds and runs
95+
* appear as SciJava {@link Task}s, and canceling a run's task cancels the
96+
* script. Inputs of array-compatible types (e.g., ImgLib2 {@code Img})
8497
* are automatically converted to Appose {@link NDArray} before being passed to
8598
* Python, and back again on the output side. Both conversions are delegated to
8699
* SciJava's {@link ConvertService}, so they work for any type that has a
@@ -104,19 +117,27 @@ public class ApposePythonScriptEngine extends AbstractScriptEngine {
104117
/** Appose task script that runs the user script; see wrapper.py. */
105118
private static final String WRAPPER_SCRIPT = loadWrapperScript();
106119

107-
/** Built environments, keyed by name, scheme and configuration content. */
108-
private static final Map<String, Environment> ENVIRONMENTS =
109-
new ConcurrentHashMap<>();
120+
/**
121+
* How long to wait, after a run finishes, for the rest of its stderr output
122+
* to arrive.
123+
*/
124+
private static final long OUTPUT_DRAIN_MILLIS = 2000;
110125

111126
@Parameter
112127
private ConvertService convertService;
113128

114129
@Parameter
115130
private LogService log;
116131

132+
@Parameter
133+
private ResidentWorkerService workerService;
134+
117135
@Parameter(required = false)
118136
private StatusService statusService;
119137

138+
@Parameter(required = false)
139+
private TaskService taskService;
140+
120141
public ApposePythonScriptEngine(final Context context) {
121142
context.inject(this);
122143
setLogService(log);
@@ -133,8 +154,8 @@ public Object eval(final String script) throws ScriptException {
133154
final ScriptInfo info = moduleObj instanceof ScriptModule ?
134155
((ScriptModule) moduleObj).getInfo() : null;
135156

136-
// Build (or retrieve from cache) the Appose environment.
137-
final Environment env = buildEnvironment(info);
157+
// Retrieve the resident worker for the script's environment.
158+
final ResidentWorker worker = worker(info);
138159

139160
// Collect declared inputs, converting array-like values to NDArray.
140161
final Map<String, Object> taskInputs = new HashMap<>();
@@ -173,7 +194,7 @@ public Object eval(final String script) throws ScriptException {
173194
taskInputs.put("_appose_array_inputs", arrayInputNames);
174195
taskInputs.put("_appose_outputs", outputNames);
175196

176-
return runTask(env, taskInputs, info);
197+
return runTask(worker, taskInputs, info);
177198
}
178199
finally {
179200
ownedNDArrays.forEach(NDArray::close);
@@ -202,32 +223,45 @@ public Bindings createBindings() {
202223
// -- Helper methods --
203224

204225
/**
205-
* Runs the wrapped script as an Appose task, storing its outputs into the
206-
* engine bindings.
226+
* Runs the wrapped script as an Appose task on the environment's resident
227+
* worker, storing its outputs into the engine bindings.
207228
*
208229
* @return The value of the script's last expression, if any.
209230
*/
210-
private Object runTask(final Environment env,
231+
private Object runTask(final ResidentWorker worker,
211232
final Map<String, Object> taskInputs, final ScriptInfo info)
212233
throws ScriptException
213234
{
214-
// Note: On Windows, importing numpy from a task hangs unless numpy
215-
// was imported during worker initialization.
216-
final Service python = env.python()
217-
.init("try:\n import numpy\nexcept ImportError:\n pass\n");
218-
python.debug(this::forwardWorkerOutput);
219-
boolean started = false;
235+
// Note: Runs on a worker are serialized, because its output cannot
236+
// otherwise be attributed to the script that produced it.
237+
final Lock lock = worker.exclusive();
238+
try {
239+
lock.lockInterruptibly();
240+
}
241+
catch (final InterruptedException e) {
242+
Thread.currentThread().interrupt();
243+
throw new ScriptException("Python script interrupted");
244+
}
245+
final OutputForwarder forwarder = new OutputForwarder();
246+
taskInputs.put("_appose_end_marker", forwarder.endMarker);
247+
worker.debug(forwarder);
248+
Service.Task task = null;
249+
Task progress = null;
220250
try {
221-
final Service.Task task = python.task(WRAPPER_SCRIPT, taskInputs);
222-
started = true;
251+
task = worker.task(WRAPPER_SCRIPT, taskInputs);
223252
task.listen(event -> {
224253
if (event.message != null) log.info("[appose-python] " + event.message);
225-
if (event.maximum > 0 && statusService != null) {
226-
statusService.showStatus((int) event.current, (int) event.maximum,
227-
event.message);
228-
}
229254
});
230-
task.waitFor();
255+
progress = SciJavaTasks.track(taskService, statusService, "Running " +
256+
scriptName(info), task, worker::kill);
257+
try {
258+
task.waitFor();
259+
}
260+
finally {
261+
if (task.status == TaskStatus.COMPLETE || task.status == TaskStatus.FAILED) {
262+
forwarder.awaitEnd();
263+
}
264+
}
231265

232266
// Unmarshal outputs from task.outputs back into engine bindings.
233267
for (final ModuleItem<?> item : info.outputs()) {
@@ -237,37 +271,40 @@ private Object runTask(final Environment env,
237271
}
238272
return unmarshal(task.outputs.get(RETURN_VALUE_KEY), Object.class);
239273
}
274+
catch (final BuildException e) {
275+
throw scriptException("Failed to build Appose environment '" + worker
276+
.name() + "': " + e.getMessage(), e);
277+
}
240278
catch (final TaskException e) {
279+
if (progress != null && progress.isCanceled()) {
280+
// Note: No cause, since ScriptModule reports the innermost cause,
281+
// which would say the worker crashed if it had to be stopped.
282+
throw new ScriptException("Python script canceled");
283+
}
241284
throw scriptException("Python script failed: " + e.getMessage(), e);
242285
}
243286
catch (final InterruptedException e) {
244-
python.kill();
287+
worker.kill();
245288
Thread.currentThread().interrupt();
246289
throw new ScriptException("Python script interrupted");
247290
}
248291
catch (final RuntimeException e) {
249292
// E.g. UncheckedIOException when the worker process fails to launch.
250-
if (python.isAlive()) python.kill();
293+
worker.kill();
251294
throw scriptException("Python script failed: " + e.getMessage(), e);
252295
}
253296
finally {
254-
if (started) shutDown(python);
297+
worker.debug(null);
298+
if (progress != null) progress.finish();
299+
if (statusService != null) statusService.clearStatus();
300+
lock.unlock();
255301
}
256302
}
257303

258-
/**
259-
* Shuts down the worker process, waiting until all of its output has been
260-
* forwarded to the script context's writers.
261-
*/
262-
private void shutDown(final Service python) {
263-
python.close();
264-
try {
265-
python.waitFor();
266-
}
267-
catch (final InterruptedException e) {
268-
python.kill();
269-
Thread.currentThread().interrupt();
270-
}
304+
/** Gets a human-friendly name for the given script. */
305+
private static String scriptName(final ScriptInfo info) {
306+
final String path = info.getPath();
307+
return path == null ? "script" : new File(path).getName();
271308
}
272309

273310
/**
@@ -279,16 +316,34 @@ private void shutDown(final Service python) {
279316
* {@code [SERVICE-n] <INVALID> line}.
280317
* </p>
281318
*/
282-
private void forwardWorkerOutput(final String message) {
283-
final int end = message.indexOf("] ");
284-
if (end < 0) return;
285-
final String prefix = message.substring(0, end);
286-
final String line = message.substring(end + 2);
287-
if (prefix.startsWith("[WORKER-")) {
288-
writeLine(getContext().getErrorWriter(), line);
319+
private class OutputForwarder implements Consumer<String> {
320+
321+
/**
322+
* Line the wrapper script writes to stderr when it is done, since stderr
323+
* lines may still arrive after the task has reported completion.
324+
*/
325+
private final String endMarker = "[appose-python-end " + UUID.randomUUID() +
326+
"]";
327+
private final CountDownLatch ended = new CountDownLatch(1);
328+
329+
@Override
330+
public void accept(final String message) {
331+
final int end = message.indexOf("] ");
332+
if (end < 0) return;
333+
final String prefix = message.substring(0, end);
334+
final String line = message.substring(end + 2);
335+
if (prefix.startsWith("[WORKER-")) {
336+
if (line.equals(endMarker)) ended.countDown();
337+
else writeLine(getContext().getErrorWriter(), line);
338+
}
339+
else if (prefix.startsWith("[SERVICE-") && line.startsWith("<INVALID> ")) {
340+
writeLine(getContext().getWriter(), line.substring(10));
341+
}
289342
}
290-
else if (prefix.startsWith("[SERVICE-") && line.startsWith("<INVALID> ")) {
291-
writeLine(getContext().getWriter(), line.substring(10));
343+
344+
/** Waits until the wrapper script's stderr output has all arrived. */
345+
private void awaitEnd() throws InterruptedException {
346+
ended.await(OUTPUT_DRAIN_MILLIS, TimeUnit.MILLISECONDS);
292347
}
293348
}
294349

@@ -323,7 +378,13 @@ private Object unmarshal(final Object raw, final Class<?> type) {
323378
if (type.isInstance(nd)) return nd;
324379
final Object converted = convertService.convert(nd, type);
325380
if (converted == null) return nd;
326-
if (converted != nd) nd.close();
381+
if (converted != nd) {
382+
// Note: By Appose convention, the service side frees shared memory,
383+
// even when the worker allocated it. The worker also stays alive
384+
// after the run, so nothing else would ever free this block.
385+
nd.shm().unlinkOnClose(true);
386+
nd.close();
387+
}
327388
return converted;
328389
}
329390

@@ -363,12 +424,12 @@ static boolean ownsNDArray(final Object value, final NDArray nd) {
363424
}
364425

365426
/**
366-
* Lazily builds the Appose {@link Environment} described by the {@code env}
367-
* and {@code scheme} attributes of the script's {@code #@script} directive.
427+
* Gets the resident worker for the Appose environment described by the
428+
* {@code env} and {@code scheme} attributes of the script's
429+
* {@code #@script} directive. The environment itself is built lazily, by
430+
* the worker's first task.
368431
*/
369-
private Environment buildEnvironment(final ScriptInfo info)
370-
throws ScriptException
371-
{
432+
private ResidentWorker worker(final ScriptInfo info) throws ScriptException {
372433
final String envRef = info == null ? null : info.get("env");
373434
if (envRef == null) {
374435
throw new ScriptException(
@@ -393,37 +454,14 @@ private Environment buildEnvironment(final ScriptInfo info)
393454
final String envName = envName(envFile);
394455
final String scheme = info.get("scheme");
395456

396-
final String key = envName + "\n" + scheme + "\n" + content;
397-
final Environment cached = ENVIRONMENTS.get(key);
398-
if (cached != null) return cached;
399-
400-
// Note: Builds are serialized, so that concurrent runs of the
401-
// same script do not trample the same environment directory.
402-
synchronized (ENVIRONMENTS) {
403-
final Environment env = ENVIRONMENTS.get(key);
404-
if (env != null) return env;
405-
406-
log.info("[appose-python] Building environment '" + envName +
407-
"' from " + envFile);
408-
try {
409-
Builder<?> builder = Appose.content(content);
410-
if (scheme != null) builder = builder.scheme(scheme);
411-
builder = builder.name(envName)
412-
.subscribeOutput(s -> log.debug(s.trim()))
413-
.subscribeError(s -> log.debug(s.trim()));
414-
if (statusService != null) {
415-
builder = builder.subscribeProgress((title, cur, max) -> statusService
416-
.showStatus((int) cur, (int) max, title));
417-
}
418-
final Environment built = builder.build();
419-
ENVIRONMENTS.put(key, built);
420-
return built;
421-
}
422-
catch (final BuildException | RuntimeException e) {
423-
throw scriptException("Failed to build Appose environment '" +
424-
envName + "' from " + envFile + ": " + e.getMessage(), e);
425-
}
426-
}
457+
return workerService.worker(envName, scheme + "\n" + content, () -> {
458+
Builder<?> builder = Appose.content(content);
459+
if (scheme != null) builder = builder.scheme(scheme);
460+
// Note: On Windows, importing numpy from a task hangs unless numpy
461+
// was imported during worker initialization.
462+
return new ResidentWorker(envName, builder.name(envName), env -> env
463+
.python().init("try:\n import numpy\nexcept ImportError:\n pass\n"));
464+
});
427465
}
428466

429467
/** Resolves an (optionally relative) env file path against the script path. */

0 commit comments

Comments
 (0)