Two problems with the prior prepare_submit_job flow:
1. The agent could pass a script_path that only existed in its own
context (LLM-generated code not yet on disk) or a path on the host
filesystem that's invisible inside the container. The new check
surfaces this as a 400 with a clear remediation hint instead of
letting spark-submit fail later with an opaque FileNotFoundError -> 500.
2. The container-isolation issue: any path the agent gives is interpreted
inside the container. The two ways a file can legitimately exist there
are (a) generate_job_file(code=...) just wrote it to
SPARK_EXECUTOR_JOBS_DIR (the default ./data/jobs/ is the only
gitignored dir that survives restarts), or (b) a host dir was
mounted via -v. The error message spells both out so an agent can
self-correct.
Implementation:
- submit.py: new _check_script_path() that raises ValueError (-> 400)
when the path is missing, empty, or a directory. Called in both
prepare_submit_job and confirm_submit_job (defense in depth).
- prepare_submit_job checks the connection FIRST (KeyError -> 404)
before the script (ValueError -> 400), so an agent with both problems
sees the more fundamental 'unknown connection' error first.
- server.py / requests.py: route description and Pydantic field
description spell out the generate_job_file pattern so an LLM
reading the tool schema learns the right next step.
8 new tests; 8 existing tests adjusted to create real files (they used
synthetic /tmp/*.py paths that don't exist).
66 lines
2.4 KiB
Python
66 lines
2.4 KiB
Python
# coding=utf-8
|
|
import io
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
|
|
from spark_executor.core import connection_store, pending_store
|
|
from spark_executor.core.connection_store import ConnectionStore
|
|
from spark_executor.core.pending_store import PendingStore
|
|
from spark_executor.models import Connection
|
|
from spark_executor.tools import connections, submit
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _fresh_stores(tmp_path: Path, monkeypatch):
|
|
monkeypatch.setattr(connection_store, "DEFAULT_DATA_DIR", str(tmp_path))
|
|
monkeypatch.setattr(connection_store, "store", ConnectionStore())
|
|
monkeypatch.setattr(pending_store, "DEFAULT_DATA_DIR", str(tmp_path))
|
|
monkeypatch.setattr(pending_store, "store", PendingStore())
|
|
connections.store = connection_store.store
|
|
submit.conn_store = connection_store.store
|
|
submit.pending_store = pending_store.store
|
|
submit.job_store = submit.job_store.__class__() # fresh in-memory job store
|
|
|
|
|
|
@pytest.fixture
|
|
def log_capture():
|
|
"""Attach an in-memory sink to loguru so tests can assert on emitted lines."""
|
|
from common.logging import logger
|
|
buf = io.StringIO()
|
|
handler_id = logger.add(buf, level="DEBUG", format="{level}|{message}")
|
|
yield buf
|
|
logger.remove(handler_id)
|
|
|
|
|
|
def test_save_connection_emits_info_log(log_capture):
|
|
connections.save_connection(name="prod", master="yarn")
|
|
text = log_capture.getvalue()
|
|
assert "INFO" in text
|
|
assert "save_connection enter" in text
|
|
assert "DEBUG" in text
|
|
assert "connection saved" in text
|
|
|
|
|
|
def test_prepare_submit_job_emits_debug_and_info(log_capture, tmp_path):
|
|
connections.save_connection(name="prod", master="yarn", deploy_mode="cluster")
|
|
script = tmp_path / "demo.py"
|
|
script.write_text("print('hi')\n")
|
|
log_capture.truncate(0); log_capture.seek(0)
|
|
submit.prepare_submit_job(connection="prod", script_path=str(script), queue="research")
|
|
text = log_capture.getvalue()
|
|
assert "DEBUG|prepare_submit_job enter" in text
|
|
assert "INFO|prepare_submit_job ok" in text
|
|
assert f"script_path={script}" in text
|
|
assert "queue=research" in text
|
|
|
|
|
|
def test_get_unknown_pending_job_emits_debug(log_capture):
|
|
log_capture.truncate(0); log_capture.seek(0)
|
|
import pytest as _pytest
|
|
with _pytest.raises(KeyError):
|
|
submit.get_pending_job("p_doesnotexist")
|
|
text = log_capture.getvalue()
|
|
assert "DEBUG|get_pending_job enter" in text
|
|
assert "p_doesnotexist" in text
|