Skip to content

Commit d3c83a2

Browse files
ctruedenclaude
andcommitted
Separate environment building from worker processes
ResidentWorker built its own environment, conflating two lifetimes: an environment is built once and can serve any number of worker processes, while a worker is one process with its own imports, exported state and memory. Two workers on one environment would each have built it. LazyEnvironment now builds an environment once, reporting to its build listeners, and serializes builds of same-named environments, which share a directory on disk. ResidentWorker is one process running in a LazyEnvironment. ResidentWorkerService owns the policy of how many workers an environment gets, which is still exactly one. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent b437fb0 commit d3c83a2

7 files changed

Lines changed: 278 additions & 77 deletions

File tree

‎README.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,9 @@ if "model" not in globals():
6767
task.export(model=model)
6868
```
6969

70-
Runs on the same environment execute one at a time. Each run appears in the
70+
Runs on the same worker execute one at a time; since each environment
71+
currently gets a single worker, so do runs on the same environment, while
72+
runs on different environments proceed in parallel. Each run appears in the
7173
application's task list, where it can be canceled. A script can notice
7274
cancelation by checking `task.cancel_requested`; if it has not stopped a few
7375
seconds after being canceled, its worker process is stopped, and the next

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

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -273,7 +273,7 @@ private Object runTask(final ResidentWorker worker,
273273
}
274274
catch (final BuildException e) {
275275
throw scriptException("Failed to build Appose environment '" + worker
276-
.name() + "': " + e.getMessage(), e);
276+
.environment().name() + "': " + e.getMessage(), e);
277277
}
278278
catch (final TaskException e) {
279279
if (progress != null && progress.isCanceled()) {
@@ -455,13 +455,13 @@ private ResidentWorker worker(final ScriptInfo info) throws ScriptException {
455455
final String scheme = info.get("scheme");
456456

457457
return workerService.worker(envName, scheme + "\n" + content, () -> {
458-
Builder<?> builder = Appose.content(content);
459-
if (scheme != null) builder = builder.scheme(scheme);
458+
final Builder<?> builder = Appose.content(content);
459+
return scheme == null ? builder : builder.scheme(scheme);
460+
},
460461
// Note: On Windows, importing numpy from a task hangs unless numpy
461462
// 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-
});
463+
env -> env.python().init(
464+
"try:\n import numpy\nexcept ImportError:\n pass\n"));
465465
}
466466

467467
/** Resolves an (optionally relative) env file path against the script path. */
Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
/*
2+
* #%L
3+
* Python scripting language plugin backed by Appose.
4+
* %%
5+
* Copyright (C) 2026 SciJava developers.
6+
* %%
7+
* Redistribution and use in source and binary forms, with or without
8+
* modification, are permitted provided that the following conditions are met:
9+
*
10+
* 1. Redistributions of source code must retain the above copyright notice,
11+
* this list of conditions and the following disclaimer.
12+
* 2. Redistributions in binary form must reproduce the above copyright notice,
13+
* this list of conditions and the following disclaimer in the documentation
14+
* and/or other materials provided with the distribution.
15+
*
16+
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS"
17+
* AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
18+
* IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE
19+
* ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDERS OR CONTRIBUTORS BE
20+
* LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR
21+
* CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF
22+
* SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS
23+
* INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN
24+
* CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE)
25+
* ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE
26+
* POSSIBILITY OF SUCH DAMAGE.
27+
* #L%
28+
*/
29+
30+
package org.scijava.plugins.scripting.appose.python._internal;
31+
32+
import java.util.List;
33+
import java.util.Map;
34+
import java.util.concurrent.ConcurrentHashMap;
35+
import java.util.concurrent.CopyOnWriteArrayList;
36+
37+
import org.apposed.appose.BuildException;
38+
import org.apposed.appose.Builder;
39+
import org.apposed.appose.Environment;
40+
41+
/**
42+
* An Appose environment that is built on first use, once, and then shared by
43+
* any number of {@link ResidentWorker}s.
44+
* <p>
45+
* Builds of environments with the same name are serialized, even across
46+
* instances, because they share a directory on disk: e.g. when the
47+
* environment's configuration changes while an older build is still running.
48+
* </p>
49+
*/
50+
public class LazyEnvironment {
51+
52+
/** Locks serializing builds, keyed by environment name. */
53+
private static final Map<String, Object> BUILD_LOCKS =
54+
new ConcurrentHashMap<>();
55+
56+
private final String name;
57+
private final Builder<?> builder;
58+
private final List<BuildListener> listeners = new CopyOnWriteArrayList<>();
59+
60+
private volatile Environment env;
61+
62+
/**
63+
* @param name Name of the environment, as reported to listeners, and as
64+
* given to the builder.
65+
* @param builder Builder for the environment; built on first use.
66+
*/
67+
public LazyEnvironment(final String name, final Builder<?> builder) {
68+
this.name = name;
69+
// Note: Subscribe once, here, and dispatch to the current listeners,
70+
// because builders accumulate subscribers with every call.
71+
this.builder = builder.name(name) //
72+
.subscribeProgress((title, current, maximum) -> listeners.forEach(
73+
l -> l.buildProgress(name, title, current, maximum))) //
74+
.subscribeOutput(text -> listeners.forEach(l -> l.buildOutput(name,
75+
text))) //
76+
.subscribeError(text -> listeners.forEach(l -> l.buildError(name,
77+
text)));
78+
}
79+
80+
public String name() {
81+
return name;
82+
}
83+
84+
public LazyEnvironment addBuildListener(final BuildListener listener) {
85+
listeners.add(listener);
86+
return this;
87+
}
88+
89+
/** Whether the environment has been built by this instance. */
90+
public boolean isBuilt() {
91+
return env != null;
92+
}
93+
94+
/** Gets the environment, building it if this is its first use. */
95+
public Environment get() throws BuildException {
96+
if (env != null) return env;
97+
synchronized (BUILD_LOCKS.computeIfAbsent(name, k -> new Object())) {
98+
if (env != null) return env;
99+
listeners.forEach(l -> l.buildStarted(name));
100+
try {
101+
env = builder.build();
102+
}
103+
catch (final BuildException | RuntimeException exc) {
104+
listeners.forEach(l -> l.buildFinished(name, exc));
105+
throw exc;
106+
}
107+
listeners.forEach(l -> l.buildFinished(name, null));
108+
return env;
109+
}
110+
}
111+
}

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

Lines changed: 14 additions & 49 deletions
Original file line numberDiff line numberDiff line change
@@ -31,23 +31,19 @@
3131

3232
import java.io.IOException;
3333
import java.io.UncheckedIOException;
34-
import java.util.List;
3534
import java.util.Map;
36-
import java.util.concurrent.CopyOnWriteArrayList;
3735
import java.util.concurrent.locks.Lock;
3836
import java.util.concurrent.locks.ReentrantLock;
3937
import java.util.function.Consumer;
4038
import java.util.function.Function;
4139

4240
import org.apposed.appose.BuildException;
43-
import org.apposed.appose.Builder;
4441
import org.apposed.appose.Environment;
4542
import org.apposed.appose.Service;
4643

4744
/**
48-
* An Appose worker that stays resident across tasks: its environment is
49-
* built on first use, its worker process is started once and reused, and
50-
* both are released only on request.
45+
* An Appose worker process that stays resident across tasks: it is started
46+
* on first use, reused afterward, and ended only on request.
5147
* <p>
5248
* Keeping the worker alive is what makes repeated runs fast: modules the
5349
* worker has imported stay imported, and objects a task hands to
@@ -57,46 +53,35 @@
5753
* free it.
5854
* </p>
5955
* <p>
60-
* If the worker dies (e.g. it crashes, or {@link #kill()} is called), the
61-
* next task transparently starts a new one.
56+
* A worker runs in a {@link LazyEnvironment}, which it builds if needed; any
57+
* number of workers may share one environment. If the worker dies (e.g. it
58+
* crashes, or {@link #kill()} is called), the next task transparently starts
59+
* a new process.
6260
* </p>
6361
*/
6462
public class ResidentWorker implements AutoCloseable {
6563

66-
private final String name;
67-
private final Builder<?> builder;
64+
private final LazyEnvironment environment;
6865
private final Function<Environment, Service> launcher;
69-
private final List<BuildListener> listeners = new CopyOnWriteArrayList<>();
7066
private final Lock exclusive = new ReentrantLock();
7167

7268
private volatile Consumer<String> debug;
73-
private Environment env;
7469
private Service service;
7570

7671
/**
77-
* @param name Name of the environment, as reported to listeners.
78-
* @param builder Builder for the environment; built on first use.
72+
* @param environment Environment the worker runs in.
7973
* @param launcher Creates the worker's service from the built environment,
8074
* e.g. {@code env -> env.python().init("import numpy")}.
8175
*/
82-
public ResidentWorker(final String name, final Builder<?> builder,
76+
public ResidentWorker(final LazyEnvironment environment,
8377
final Function<Environment, Service> launcher)
8478
{
85-
this.name = name;
79+
this.environment = environment;
8680
this.launcher = launcher;
87-
// Note: Subscribe once, here, and dispatch to the current listeners,
88-
// because builders accumulate subscribers with every call.
89-
this.builder = builder //
90-
.subscribeProgress((title, current, maximum) -> listeners.forEach(
91-
l -> l.buildProgress(name, title, current, maximum))) //
92-
.subscribeOutput(text -> listeners.forEach(l -> l.buildOutput(name,
93-
text))) //
94-
.subscribeError(text -> listeners.forEach(l -> l.buildError(name,
95-
text)));
9681
}
9782

98-
public String name() {
99-
return name;
83+
public LazyEnvironment environment() {
84+
return environment;
10085
}
10186

10287
/**
@@ -108,11 +93,6 @@ public Lock exclusive() {
10893
return exclusive;
10994
}
11095

111-
public ResidentWorker addBuildListener(final BuildListener listener) {
112-
listeners.add(listener);
113-
return this;
114-
}
115-
11696
/**
11797
* Sets where the worker's debug output goes: stderr lines, non-protocol
11898
* stdout lines and protocol traffic, as described by
@@ -123,21 +103,6 @@ public void debug(final Consumer<String> listener) {
123103
debug = listener;
124104
}
125105

126-
/** Gets the environment, building it if this is its first use. */
127-
public synchronized Environment environment() throws BuildException {
128-
if (env != null) return env;
129-
listeners.forEach(l -> l.buildStarted(name));
130-
try {
131-
env = builder.build();
132-
}
133-
catch (final BuildException | RuntimeException exc) {
134-
listeners.forEach(l -> l.buildFinished(name, exc));
135-
throw exc;
136-
}
137-
listeners.forEach(l -> l.buildFinished(name, null));
138-
return env;
139-
}
140-
141106
/**
142107
* Gets the worker's service, building the environment and starting the
143108
* worker process first if needed.
@@ -146,7 +111,7 @@ public synchronized Environment environment() throws BuildException {
146111
*/
147112
public synchronized Service service() throws BuildException {
148113
if (service != null && service.isAlive()) return service;
149-
final Service s = launcher.apply(environment());
114+
final Service s = launcher.apply(environment.get());
150115
s.debug(msg -> {
151116
final Consumer<String> listener = debug;
152117
if (listener != null) listener.accept(msg);
@@ -161,7 +126,7 @@ public synchronized Service service() throws BuildException {
161126
return service;
162127
}
163128

164-
/** Creates a task to run on the resident worker. */
129+
/** Creates a task to run on this worker. */
165130
public Service.Task task(final String script,
166131
final Map<String, Object> inputs) throws BuildException
167132
{

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

Lines changed: 32 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -33,8 +33,11 @@
3333
import java.util.HashMap;
3434
import java.util.List;
3535
import java.util.Map;
36+
import java.util.function.Function;
3637
import java.util.function.Supplier;
3738

39+
import org.apposed.appose.Builder;
40+
import org.apposed.appose.Environment;
3841
import org.scijava.app.StatusService;
3942
import org.scijava.log.LogService;
4043
import org.scijava.plugin.Parameter;
@@ -45,9 +48,15 @@
4548
import org.scijava.task.TaskService;
4649

4750
/**
48-
* Keeps one {@link ResidentWorker} per Appose environment for the lifetime of
49-
* the application context, showing environment builds as SciJava tasks, and
50-
* ending every worker when the context is disposed.
51+
* Keeps Appose environments and their {@link ResidentWorker}s for the
52+
* lifetime of the application context, showing environment builds as SciJava
53+
* tasks, and ending every worker when the context is disposed.
54+
* <p>
55+
* Each environment currently gets exactly one worker. Several workers per
56+
* environment would let runs proceed in parallel, but each worker holds its
57+
* own copy of whatever a script loads (GPU memory included), and state kept
58+
* via {@code task.export(...)} would live in only one of them.
59+
* </p>
5160
*/
5261
@Plugin(type = Service.class)
5362
public class ResidentWorkerService extends AbstractService implements
@@ -63,35 +72,38 @@ public class ResidentWorkerService extends AbstractService implements
6372
@Parameter(required = false)
6473
private LogService log;
6574

66-
/** Workers by environment name, each with the configuration it was made for. */
67-
private final Map<String, Entry> workers = new HashMap<>();
75+
/** Environments by name, each with its configuration and worker. */
76+
private final Map<String, Entry> entries = new HashMap<>();
6877

6978
/**
70-
* Gets the resident worker for the named environment, creating it if
71-
* needed.
79+
* Gets the resident worker for the named environment, creating the
80+
* environment and worker if needed.
7281
* <p>
73-
* If the environment's configuration has changed since its worker was
74-
* created, as indicated by {@code config}, the old worker is released and
82+
* If the environment's configuration has changed since it was created, as
83+
* indicated by {@code config}, its old worker is released and both are
7584
* replaced, so that the next run uses the updated environment.
7685
* </p>
7786
*
7887
* @param name Name of the environment.
7988
* @param config The environment's configuration, e.g. its file contents.
80-
* @param factory Creates a worker for the environment.
89+
* @param builder Creates a builder for the environment.
90+
* @param launcher Creates a worker's service from the built environment.
8191
*/
8292
public ResidentWorker worker(final String name, final String config,
83-
final Supplier<ResidentWorker> factory)
93+
final Supplier<Builder<?>> builder,
94+
final Function<Environment, org.apposed.appose.Service> launcher)
8495
{
8596
final ResidentWorker stale;
8697
final ResidentWorker worker;
87-
synchronized (workers) {
88-
final Entry entry = workers.get(name);
98+
synchronized (entries) {
99+
final Entry entry = entries.get(name);
89100
if (entry != null && entry.config.equals(config)) return entry.worker;
90101
stale = entry == null ? null : entry.worker;
91-
worker = factory.get();
92-
worker.addBuildListener(SciJavaTasks.buildListener(taskService,
102+
final LazyEnvironment env = new LazyEnvironment(name, builder.get());
103+
env.addBuildListener(SciJavaTasks.buildListener(taskService,
93104
statusService, log));
94-
workers.put(name, new Entry(config, worker));
105+
worker = new ResidentWorker(env, launcher);
106+
entries.put(name, new Entry(config, worker));
95107
}
96108
if (stale != null) stale.release();
97109
return worker;
@@ -100,8 +112,8 @@ public ResidentWorker worker(final String name, final String config,
100112
/** Gets the current resident workers. */
101113
public List<ResidentWorker> workers() {
102114
final List<ResidentWorker> list = new ArrayList<>();
103-
synchronized (workers) {
104-
for (final Entry entry : workers.values()) list.add(entry.worker);
115+
synchronized (entries) {
116+
for (final Entry entry : entries.values()) list.add(entry.worker);
105117
}
106118
return list;
107119
}
@@ -119,8 +131,8 @@ public void releaseAll() {
119131
@Override
120132
public void dispose() {
121133
releaseAll();
122-
synchronized (workers) {
123-
workers.clear();
134+
synchronized (entries) {
135+
entries.clear();
124136
}
125137
}
126138

0 commit comments

Comments
 (0)