forked from retoor/devplacepy
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
32314fc6d6 |
File diff suppressed because one or more lines are too long
@@ -7,12 +7,12 @@ from devplacepy.cli._shared import _audit_cli
|
|||||||
def _remove_zip_artifacts(job):
|
def _remove_zip_artifacts(job):
|
||||||
import shutil
|
import shutil
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from devplacepy.services.jobs.zip_service import STAGING_DIR
|
from devplacepy.config import ZIP_STAGING_DIR
|
||||||
|
|
||||||
local_path = (job.get("result") or {}).get("local_path")
|
local_path = (job.get("result") or {}).get("local_path")
|
||||||
if local_path:
|
if local_path:
|
||||||
Path(local_path).unlink(missing_ok=True)
|
Path(local_path).unlink(missing_ok=True)
|
||||||
shutil.rmtree(STAGING_DIR / job["uid"], ignore_errors=True)
|
shutil.rmtree(ZIP_STAGING_DIR / job["uid"], ignore_errors=True)
|
||||||
|
|
||||||
|
|
||||||
def cmd_zips_prune(args):
|
def cmd_zips_prune(args):
|
||||||
|
|||||||
@@ -685,7 +685,6 @@ def import_from_dir(project_uid: str, src_dir, user: dict, *, skip_names=None) -
|
|||||||
|
|
||||||
|
|
||||||
def export_to_dir(project_uid: str, subpath: str, dest_dir) -> int:
|
def export_to_dir(project_uid: str, subpath: str, dest_dir) -> int:
|
||||||
_guard_writable(project_uid)
|
|
||||||
dest = Path(dest_dir).resolve()
|
dest = Path(dest_dir).resolve()
|
||||||
dest.mkdir(parents=True, exist_ok=True)
|
dest.mkdir(parents=True, exist_ok=True)
|
||||||
if subpath:
|
if subpath:
|
||||||
|
|||||||
@@ -2,10 +2,8 @@
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import socket
|
|
||||||
import subprocess
|
import subprocess
|
||||||
import sys
|
import sys
|
||||||
import time
|
|
||||||
|
|
||||||
from devplacepy.config import XMLRPC_BIND, XMLRPC_PORT
|
from devplacepy.config import XMLRPC_BIND, XMLRPC_PORT
|
||||||
from devplacepy.services.base import BaseService
|
from devplacepy.services.base import BaseService
|
||||||
@@ -27,48 +25,20 @@ class XmlrpcService(BaseService):
|
|||||||
def __init__(self) -> None:
|
def __init__(self) -> None:
|
||||||
super().__init__("xmlrpc", interval_seconds=XMLRPC_INTERVAL_SECONDS)
|
super().__init__("xmlrpc", interval_seconds=XMLRPC_INTERVAL_SECONDS)
|
||||||
self._process: subprocess.Popen | None = None
|
self._process: subprocess.Popen | None = None
|
||||||
self._last_stderr: str | None = None
|
|
||||||
|
|
||||||
def _alive(self) -> bool:
|
def _alive(self) -> bool:
|
||||||
return self._process is not None and self._process.poll() is None
|
return self._process is not None and self._process.poll() is None
|
||||||
|
|
||||||
def _spawn(self) -> None:
|
def _spawn(self) -> None:
|
||||||
self._last_stderr = None
|
|
||||||
self._process = subprocess.Popen(
|
self._process = subprocess.Popen(
|
||||||
[sys.executable, "-m", SERVER_MODULE],
|
[sys.executable, "-m", SERVER_MODULE],
|
||||||
stdout=subprocess.DEVNULL,
|
stdout=subprocess.DEVNULL,
|
||||||
stderr=subprocess.PIPE,
|
stderr=subprocess.DEVNULL,
|
||||||
)
|
)
|
||||||
self.log(
|
self.log(
|
||||||
f"Forking XML-RPC server spawning (pid {self._process.pid}) on "
|
f"Forking XML-RPC server started (pid {self._process.pid}) on "
|
||||||
f"{XMLRPC_BIND}:{XMLRPC_PORT}"
|
f"{XMLRPC_BIND}:{XMLRPC_PORT}"
|
||||||
)
|
)
|
||||||
time.sleep(0.5)
|
|
||||||
if self._process.poll() is not None:
|
|
||||||
stderr_data = self._process.communicate()[1]
|
|
||||||
if stderr_data:
|
|
||||||
self._last_stderr = stderr_data.decode("utf-8", errors="replace")
|
|
||||||
self.log(
|
|
||||||
f"Forking XML-RPC server died after spawn (pid {self._process.pid}, "
|
|
||||||
f"exit code {self._process.returncode})"
|
|
||||||
)
|
|
||||||
if self._last_stderr:
|
|
||||||
for line in self._last_stderr.strip().split("\n"):
|
|
||||||
self.log(f" stderr: {line}")
|
|
||||||
return
|
|
||||||
try:
|
|
||||||
sock = socket.create_connection((XMLRPC_BIND, XMLRPC_PORT), timeout=1)
|
|
||||||
sock.close()
|
|
||||||
except (OSError, socket.timeout):
|
|
||||||
self.log(
|
|
||||||
f"Forking XML-RPC server not yet ready (pid {self._process.pid}) "
|
|
||||||
f"on {XMLRPC_BIND}:{XMLRPC_PORT}"
|
|
||||||
)
|
|
||||||
else:
|
|
||||||
self.log(
|
|
||||||
f"Forking XML-RPC server started (pid {self._process.pid}) on "
|
|
||||||
f"{XMLRPC_BIND}:{XMLRPC_PORT}"
|
|
||||||
)
|
|
||||||
|
|
||||||
def _terminate(self) -> None:
|
def _terminate(self) -> None:
|
||||||
if not self._alive():
|
if not self._alive():
|
||||||
@@ -95,13 +65,6 @@ class XmlrpcService(BaseService):
|
|||||||
self.log(f"XML-RPC server healthy (pid {self._process.pid})")
|
self.log(f"XML-RPC server healthy (pid {self._process.pid})")
|
||||||
return
|
return
|
||||||
self.log("XML-RPC server not running, starting it")
|
self.log("XML-RPC server not running, starting it")
|
||||||
if self._process is not None:
|
|
||||||
self.log(
|
|
||||||
f"Previous process exited with code {self._process.returncode}"
|
|
||||||
)
|
|
||||||
if self._last_stderr:
|
|
||||||
for line in self._last_stderr.strip().split("\n"):
|
|
||||||
self.log(f" last stderr: {line}")
|
|
||||||
self._spawn()
|
self._spawn()
|
||||||
|
|
||||||
def collect_metrics(self) -> dict:
|
def collect_metrics(self) -> dict:
|
||||||
|
|||||||
@@ -3,7 +3,6 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
import sys
|
|
||||||
from socketserver import ForkingMixIn
|
from socketserver import ForkingMixIn
|
||||||
from xmlrpc.server import SimpleXMLRPCDispatcher, SimpleXMLRPCRequestHandler, SimpleXMLRPCServer
|
from xmlrpc.server import SimpleXMLRPCDispatcher, SimpleXMLRPCRequestHandler, SimpleXMLRPCServer
|
||||||
|
|
||||||
@@ -117,9 +116,6 @@ def main() -> None:
|
|||||||
server.serve_forever()
|
server.serve_forever()
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt:
|
||||||
logger.info("XML-RPC server interrupted")
|
logger.info("XML-RPC server interrupted")
|
||||||
except Exception:
|
|
||||||
logging.exception("XML-RPC server crashed with unhandled exception")
|
|
||||||
sys.exit(1)
|
|
||||||
finally:
|
finally:
|
||||||
server.server_close()
|
server.server_close()
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -16,7 +16,7 @@ def _init_db_fork_jobs():
|
|||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
def fork_env(tmp_path, monkeypatch):
|
def fork_env(tmp_path, monkeypatch):
|
||||||
monkeypatch.setattr(
|
monkeypatch.setattr(
|
||||||
"devplacepy.services.jobs.fork_service.STAGING_DIR", tmp_path / "staging"
|
"devplacepy.services.jobs.fork_service.FORK_STAGING_DIR", tmp_path / "staging"
|
||||||
)
|
)
|
||||||
monkeypatch.setattr("devplacepy.project_files.PROJECT_FILES_DIR", tmp_path / "pf")
|
monkeypatch.setattr("devplacepy.project_files.PROJECT_FILES_DIR", tmp_path / "pf")
|
||||||
yield tmp_path
|
yield tmp_path
|
||||||
|
|||||||
Reference in New Issue
Block a user