The worker script is handed to the interpreter by path, so app/flow was sys.path[0] for every node: `import queue` got the engine's. It now drops its own directory before anything else imports, and runs with the deployment's credentials scrubbed out of its environment. Also: reload builds off the event loop, the pool wakes what is blocked on it when it stops, a refused metrics flush is kept for the next one rather than dropped, and the cascade events are paired through failures. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017MeiWk3Yq12n2pTvnQWYvt
178 lines
5.8 KiB
Python
178 lines
5.8 KiB
Python
"""Python nodes run in a worker process, and stay there when things go wrong."""
|
|
|
|
import sys
|
|
import threading
|
|
import time
|
|
from collections.abc import Iterator
|
|
|
|
import pytest
|
|
|
|
from app.flow.workers import NodeCancelled, NodeTimeout, PythonWorkerPool
|
|
|
|
|
|
@pytest.fixture
|
|
def pool() -> Iterator[PythonWorkerPool]:
|
|
# One worker: a respawn is then provably the same slot coming back.
|
|
worker_pool = PythonWorkerPool(python=sys.executable, size=1)
|
|
worker_pool.start()
|
|
yield worker_pool
|
|
worker_pool.stop()
|
|
|
|
|
|
def run(pool: PythonWorkerPool, code: str, node: str = "demo", **kwargs):
|
|
return pool.run(
|
|
"demo", node, code, kwargs, {"factor": 2}, f"demo.{node}", timeout=5
|
|
)
|
|
|
|
|
|
def test_a_node_returns_its_value_and_what_it_printed(pool, capsys):
|
|
result = run(
|
|
pool,
|
|
"def process(value, params):\n"
|
|
" print('seen', value)\n"
|
|
" return {'out': value * params['factor']}\n",
|
|
value=21,
|
|
)
|
|
assert result == {"out": 42}
|
|
# The proxy writes them to stdout, which is where the engine's tee is.
|
|
assert "seen 21" in capsys.readouterr().out
|
|
|
|
|
|
def test_a_failure_keeps_its_class_and_points_at_the_node(pool):
|
|
with pytest.raises(Exception) as caught:
|
|
run(pool, "def process(params):\n raise ValueError('bad input')\n")
|
|
|
|
# The engine renders a node error as "<class>: <message>", so both have to
|
|
# survive the trip.
|
|
assert type(caught.value).__name__ == "ValueError"
|
|
assert str(caught.value) == "bad input"
|
|
assert "<node demo." in caught.value.remote_traceback
|
|
assert "ValueError: bad input" in caught.value.remote_traceback
|
|
|
|
|
|
def test_a_node_that_kills_its_worker_is_an_ordinary_error(pool):
|
|
with pytest.raises(Exception, match="worker died"):
|
|
run(pool, "import os\n\n\ndef process(params):\n os._exit(1)\n")
|
|
|
|
assert run(pool, "def process(params):\n return {'out': 1}\n") == {"out": 1}
|
|
|
|
|
|
def test_a_node_that_runs_too_long_is_killed_and_the_pool_recovers(pool):
|
|
started = time.monotonic()
|
|
with pytest.raises(NodeTimeout):
|
|
pool.run(
|
|
"demo",
|
|
"slow",
|
|
"import time\n\n\ndef process(params):\n time.sleep(30)\n",
|
|
{},
|
|
{},
|
|
"demo.slow",
|
|
timeout=1,
|
|
)
|
|
assert time.monotonic() - started < 10
|
|
|
|
# The killed worker's slot is refilled on the next call.
|
|
assert run(pool, "def process(params):\n return {'out': 2}\n") == {"out": 2}
|
|
|
|
|
|
def test_a_running_node_can_be_cancelled(pool):
|
|
def stop_it() -> None:
|
|
for _ in range(100):
|
|
if pool.cancel("demo.slow"):
|
|
return
|
|
time.sleep(0.05)
|
|
|
|
stopper = threading.Thread(target=stop_it)
|
|
stopper.start()
|
|
try:
|
|
with pytest.raises(NodeCancelled):
|
|
pool.run(
|
|
"demo",
|
|
"slow",
|
|
"import time\n\n\ndef process(params):\n time.sleep(30)\n",
|
|
{},
|
|
{},
|
|
"demo.slow",
|
|
timeout=30,
|
|
)
|
|
finally:
|
|
stopper.join()
|
|
|
|
|
|
def test_a_result_that_is_not_json_is_refused(pool):
|
|
with pytest.raises(Exception, match="cannot be sent back as JSON"):
|
|
run(pool, "def process(params):\n return {'out': {1, 2}}\n")
|
|
|
|
|
|
def test_compiling_reports_where_the_source_is_wrong(pool):
|
|
assert pool.compile("demo", "broken", "def process(params)\n return {}\n")
|
|
assert pool.compile("demo", "fine", "def process(params):\n return {}\n") is None
|
|
|
|
|
|
def test_a_node_imports_the_standard_library_not_the_engines_own_modules(pool):
|
|
# The worker script lives in app/flow, which holds queue.py, secrets.py and
|
|
# more; the interpreter would put that directory first on sys.path.
|
|
result = run(
|
|
pool,
|
|
"import queue\nimport secrets\n\n\n"
|
|
"def process(params):\n"
|
|
" return {'out': [queue.Queue().qsize(), len(secrets.token_hex(4))]}\n",
|
|
)
|
|
assert result == {"out": [0, 8]}
|
|
|
|
|
|
def test_the_engines_secrets_are_not_in_a_workers_environment(pool, monkeypatch):
|
|
monkeypatch.setenv("SECRET_KEY", "not-for-nodes")
|
|
monkeypatch.setenv("POSTGRES_PASSWORD", "not-for-nodes")
|
|
monkeypatch.setenv("FLUKSIO_HARMLESS", "fine")
|
|
# A fresh process, so it is built from the environment set just now.
|
|
pool.respawn_all()
|
|
|
|
result = run(
|
|
pool,
|
|
"import os\n\n\n"
|
|
"def process(params):\n"
|
|
" return {'out': [k for k in ('SECRET_KEY', 'POSTGRES_PASSWORD',\n"
|
|
" 'FLUKSIO_HARMLESS') if k in os.environ]}\n",
|
|
)
|
|
assert result == {"out": ["FLUKSIO_HARMLESS"]}
|
|
|
|
|
|
def test_a_pool_can_stop_while_a_node_is_running(pool):
|
|
# One slot, taken by a node that will not finish on its own, and a second
|
|
# call queued behind it. The engine's node threads are not daemons, so a
|
|
# wait here is a shutdown that never completes.
|
|
outcomes: list[str] = []
|
|
|
|
def call(node: str) -> None:
|
|
try:
|
|
pool.run(
|
|
"demo",
|
|
node,
|
|
"import time\n\n\ndef process(params):\n time.sleep(60)\n",
|
|
{},
|
|
{},
|
|
f"demo.{node}",
|
|
timeout=60,
|
|
)
|
|
outcomes.append("returned")
|
|
except Exception as exc:
|
|
outcomes.append(type(exc).__name__)
|
|
|
|
busy = threading.Thread(target=call, args=("busy",))
|
|
busy.start()
|
|
# Let the first one take the slot, so the second is blocked acquiring it.
|
|
time.sleep(1)
|
|
waiting = threading.Thread(target=call, args=("waiting",))
|
|
waiting.start()
|
|
time.sleep(0.2)
|
|
|
|
pool.stop()
|
|
for thread in (busy, waiting):
|
|
thread.join(timeout=10)
|
|
assert not thread.is_alive()
|
|
assert len(outcomes) == 2
|
|
|
|
with pytest.raises(Exception, match="shutting down"):
|
|
run(pool, "def process(params):\n return {'out': 1}\n")
|