mirror of
https://github.com/baketnk/frame-yap.git
synced 2026-10-04 22:00:03 +02:00
feat(worker): dispatch verified backends and preserve request-level failures
This commit is contained in:
1 parent
ebaae02b97
commit
9832ec1fd4
8 files changed
+549
-18
No files matched your search
@@ -0,0 +1,146 @@
|
||||
"""Manifest-driven local worker dispatcher. One verified model, one child, no shell/network.
|
||||
|
||||
The child speaks frameyap-worker-v1; this process checks framing and correlation
|
||||
before forwarding Y/T/R/E. The native parent owns the private clip and deadlines.
|
||||
"""
|
||||
import argparse
|
||||
import os
|
||||
import signal
|
||||
from pathlib import Path
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
try:
|
||||
from .model_files import load_backends, check_model, ManifestError
|
||||
from .worker import read_frame, send_frame, private_dir
|
||||
except ImportError:
|
||||
from model_files import load_backends, check_model, ManifestError
|
||||
from worker import read_frame, send_frame, private_dir
|
||||
|
||||
|
||||
def launcher_command(backend, root, python, model_dir, clip_dir, threads):
|
||||
"""Return argv, never a shell string. Manifest paths are relative to release root."""
|
||||
relative = backend.launcher["path"]
|
||||
root = Path(root).resolve(strict=True)
|
||||
target = root / relative
|
||||
if not target.is_file() or target.is_symlink() or not target.resolve().is_relative_to(root):
|
||||
raise ValueError("unsafe or missing backend launcher")
|
||||
replacements = {"{model_dir}": str(model_dir), "{clip_dir}": str(clip_dir),
|
||||
"{threads}": str(threads)}
|
||||
args = [replacements.get(a, a) for a in backend.launcher["arguments"]]
|
||||
if backend.launcher["type"] == "python":
|
||||
return [python, str(target), *args]
|
||||
if not os.access(target, os.X_OK):
|
||||
raise ValueError("backend executable is not executable")
|
||||
return [str(target), *args]
|
||||
|
||||
|
||||
def serve(backend, root, python, model_dir, clip_dir, threads, advanced_debug=False,
|
||||
input_fd=0, output_fd=1):
|
||||
private_dir(clip_dir)
|
||||
if check_model(backend, model_dir)["state"] != "installed_verified":
|
||||
send_frame(output_fd, b"F", b"M")
|
||||
return 1
|
||||
try:
|
||||
argv = launcher_command(backend, root, python, model_dir, clip_dir, threads)
|
||||
except (OSError, ValueError):
|
||||
send_frame(output_fd, b"F", b"I")
|
||||
return 1
|
||||
env = os.environ.copy()
|
||||
env.update(HF_HUB_OFFLINE="1", TRANSFORMERS_OFFLINE="1", HF_DATASETS_OFFLINE="1",
|
||||
TOKENIZERS_PARALLELISM="false", CUDA_VISIBLE_DEVICES="",
|
||||
OMP_NUM_THREADS=str(threads), MKL_NUM_THREADS=str(threads),
|
||||
OPENBLAS_NUM_THREADS=str(threads), PYTHONDONTWRITEBYTECODE="1")
|
||||
# Protocol opt-in only for our known Redux worker. Arbitrary backend
|
||||
# launchers need not accept --advanced-debug or know this environment key.
|
||||
env["FRAMEYAP_ADVANCED_DEBUG"] = "1" if (advanced_debug and backend.id == "redux" and
|
||||
backend.launcher["type"] == "python" and backend.launcher["path"] == "python/frameyap/worker.py") else "0"
|
||||
# A child cannot inherit input-delivery authority. stdout is exclusively frames.
|
||||
try:
|
||||
child = subprocess.Popen(argv, stdin=subprocess.PIPE, stdout=subprocess.PIPE,
|
||||
stderr=None if advanced_debug else subprocess.DEVNULL,
|
||||
env=env, close_fds=True)
|
||||
except OSError:
|
||||
send_frame(output_fd, b"F", b"I")
|
||||
return 1
|
||||
def stop_child(_signal, _frame):
|
||||
# Native shutdown may SIGKILL this dispatcher after a short grace period.
|
||||
# Reap the direct model child before exiting, never leave inference running.
|
||||
child.kill()
|
||||
child.wait()
|
||||
raise SystemExit(1)
|
||||
|
||||
old_term = signal.signal(signal.SIGTERM, stop_child)
|
||||
try:
|
||||
ready = read_frame(child.stdout.fileno())
|
||||
if ready is None or ready[:1] == b"F":
|
||||
send_frame(output_fd, b"F", ready[1:2] if ready and ready[1:2] in (b"M", b"I") else b"D")
|
||||
return 1
|
||||
if ready != b"Y":
|
||||
send_frame(output_fd, b"F", b"D")
|
||||
return 1
|
||||
send_frame(output_fd, b"Y")
|
||||
while True:
|
||||
request = read_frame(input_fd)
|
||||
if request is None:
|
||||
return 0
|
||||
if len(request) != 9 or request[:1] != b"T":
|
||||
raise ValueError("invalid worker request")
|
||||
send_frame(child.stdin.fileno(), b"T", request[1:])
|
||||
reply = read_frame(child.stdout.fileno())
|
||||
if (reply is None or len(reply) < 9 or len(reply) > 4105 or
|
||||
reply[:1] not in (b"R", b"E") or reply[1:9] != request[1:]):
|
||||
raise ValueError("invalid backend response")
|
||||
send_frame(output_fd, reply[:1], reply[1:])
|
||||
finally:
|
||||
# An EOF/invalid protocol must not block cleanup if a child ignores TERM.
|
||||
# Native also kills the dedicated owned worker group after its grace period.
|
||||
try:
|
||||
child.stdin.close()
|
||||
if child.poll() is None:
|
||||
child.terminate()
|
||||
try:
|
||||
child.wait(timeout=0.2)
|
||||
except subprocess.TimeoutExpired:
|
||||
child.kill()
|
||||
child.wait()
|
||||
else:
|
||||
child.wait()
|
||||
finally:
|
||||
child.stdout.close()
|
||||
signal.signal(signal.SIGTERM, old_term)
|
||||
|
||||
|
||||
def main(argv=None):
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("--backend", required=True)
|
||||
parser.add_argument("--manifest-dir", required=True, type=Path)
|
||||
parser.add_argument("--root", required=True, type=Path)
|
||||
parser.add_argument("--python", required=True)
|
||||
parser.add_argument("--model", required=True, type=Path)
|
||||
parser.add_argument("--clip-dir", required=True, type=Path)
|
||||
parser.add_argument("--threads", required=True, type=int)
|
||||
parser.add_argument("--advanced-debug", action="store_true")
|
||||
args = parser.parse_args(argv)
|
||||
if not 1 <= args.threads <= 64:
|
||||
parser.error("threads must be 1..64")
|
||||
# Never let native/model library stdout corrupt the framed channel.
|
||||
protocol = os.dup(1)
|
||||
with open(os.devnull, "wb") as null:
|
||||
os.dup2(2 if args.advanced_debug else null.fileno(), 1)
|
||||
if not args.advanced_debug:
|
||||
os.dup2(null.fileno(), 2)
|
||||
try:
|
||||
try:
|
||||
backend = load_backends(args.manifest_dir)[args.backend]
|
||||
except (ManifestError, KeyError):
|
||||
send_frame(protocol, b"F", b"M")
|
||||
return 1
|
||||
return serve(backend, args.root, args.python, args.model, args.clip_dir,
|
||||
args.threads, args.advanced_debug, output_fd=protocol)
|
||||
finally:
|
||||
os.close(protocol)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -182,6 +182,9 @@ def main(argv=None):
|
||||
parser.add_argument("--advanced-debug", action="store_true",
|
||||
help="log full exceptions/runtime output and transcripts to stderr; may contain private speech")
|
||||
args = parser.parse_args(argv)
|
||||
# Dispatcher uses this explicit protocol preference for the pinned Redux
|
||||
# launcher, rather than adding flags to arbitrary backend commands.
|
||||
args.advanced_debug = args.advanced_debug or os.environ.get("FRAMEYAP_ADVANCED_DEBUG") == "1"
|
||||
if not 1 <= args.threads <= 64:
|
||||
parser.error("threads must be 1..64")
|
||||
# Native libraries sometimes print directly to fd 1. Keep those bytes out of
|
||||
|
||||
@@ -0,0 +1,106 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Offline status emitter / explicit installer handoff for the native model chooser.
|
||||
|
||||
Status is percent-encoded tab-separated records to keep the C++ UI free of a
|
||||
second manifest/JSON/hash implementation. Installer execution is opt-in ONLY.
|
||||
"""
|
||||
import argparse
|
||||
import hashlib
|
||||
import os
|
||||
import re
|
||||
import stat
|
||||
from pathlib import Path
|
||||
import sys
|
||||
from urllib.parse import quote
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "python"))
|
||||
from frameyap.model_files import load_backends, check_model
|
||||
|
||||
|
||||
def emit(*fields):
|
||||
print("\t".join(quote(str(field), safe="") for field in fields), flush=True)
|
||||
|
||||
|
||||
_DIGEST = re.compile(r"[0-9a-f]{64}\Z", re.ASCII)
|
||||
|
||||
|
||||
def manifest_sha(manifest_dir, ident):
|
||||
"""Hash bounded raw ID.json bytes, never a reconstructed manifest."""
|
||||
path = manifest_dir / (ident + ".json")
|
||||
fd = os.open(path, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK)
|
||||
with os.fdopen(fd, "rb") as stream:
|
||||
before = os.fstat(stream.fileno())
|
||||
if not stat.S_ISREG(before.st_mode) or before.st_size > 65536:
|
||||
raise ValueError("unsafe backend manifest")
|
||||
raw = stream.read(65537)
|
||||
after = os.fstat(stream.fileno())
|
||||
if len(raw) > 65536 or (before.st_size, before.st_mtime_ns, before.st_ctime_ns) != (
|
||||
after.st_size, after.st_mtime_ns, after.st_ctime_ns):
|
||||
raise ValueError("backend manifest changed")
|
||||
return hashlib.sha256(raw).hexdigest()
|
||||
|
||||
|
||||
def main(argv=None):
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
operation = parser.add_mutually_exclusive_group(required=True)
|
||||
operation.add_argument("--status", action="store_true")
|
||||
operation.add_argument("--install", action="store_true")
|
||||
parser.add_argument("--manifest-dir", type=Path, required=True)
|
||||
parser.add_argument("--model-store", type=Path, required=True)
|
||||
parser.add_argument("--backend")
|
||||
parser.add_argument("--installer", type=Path)
|
||||
parser.add_argument("--legacy-model", type=Path)
|
||||
parser.add_argument("--expected-manifest-sha256")
|
||||
args = parser.parse_args(argv)
|
||||
if not args.model_store.is_absolute() or not args.manifest_dir.is_absolute():
|
||||
parser.error("paths must be absolute")
|
||||
if args.status:
|
||||
if args.expected_manifest_sha256:
|
||||
parser.error("expected manifest applies only to install")
|
||||
try:
|
||||
# Bracket parsing and model checks with the same raw-byte snapshots.
|
||||
# Never emit a consent record for metadata changed during status.
|
||||
before = {p.stem: manifest_sha(args.manifest_dir, p.stem)
|
||||
for p in args.manifest_dir.iterdir() if p.suffix == ".json"}
|
||||
records = []
|
||||
for ident, backend in load_backends(args.manifest_dir).items():
|
||||
path = args.model_store / ident
|
||||
status = check_model(backend, path)
|
||||
if ident == "redux" and args.legacy_model and status["state"] == "not_installed":
|
||||
legacy = check_model(backend, args.legacy_model)
|
||||
if legacy["state"] == "installed_verified":
|
||||
status, path = legacy, args.legacy_model
|
||||
records.append(("ST", ident, backend.display_name, status["state"],
|
||||
status["reason"] or "", backend.total_bytes, backend.source,
|
||||
backend.license_id, backend.license_text, backend.attribution, path, before[ident]))
|
||||
if before != {p.stem: manifest_sha(args.manifest_dir, p.stem)
|
||||
for p in args.manifest_dir.iterdir() if p.suffix == ".json"}:
|
||||
raise ValueError("backend metadata changed during status")
|
||||
for record in records:
|
||||
emit(*record)
|
||||
emit("DONE")
|
||||
return 0
|
||||
except (ValueError, OSError) as error:
|
||||
emit("ERROR", type(error).__name__)
|
||||
return 1
|
||||
# No implicit installation: requires explicit --install AND backend and caller
|
||||
# consent. The exec makes the installer the owned child (no hidden grandchild).
|
||||
if (not args.backend or not args.installer or not args.expected_manifest_sha256 or
|
||||
not _DIGEST.fullmatch(args.expected_manifest_sha256)):
|
||||
parser.error("backend, installer and 64-character consent fingerprint required")
|
||||
try:
|
||||
backends = load_backends(args.manifest_dir)
|
||||
if args.backend not in backends or manifest_sha(args.manifest_dir, args.backend) != args.expected_manifest_sha256:
|
||||
parser.error("selected backend manifest changed since consent")
|
||||
except (OSError, ValueError) as error:
|
||||
parser.error(f"selected backend manifest unavailable: {type(error).__name__}")
|
||||
if not args.installer.is_file() or args.installer.is_symlink():
|
||||
parser.error("installer missing or unsafe")
|
||||
dest = args.model_store / args.backend
|
||||
os.execv("/bin/sh", ["sh", str(args.installer), "--install-model", "--backend", args.backend,
|
||||
"--model-dir", str(dest), "--expected-manifest-sha256",
|
||||
args.expected_manifest_sha256, "--yes", "--json"])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
+45
-17
@@ -250,17 +250,18 @@ void Worker::stop() {
|
||||
close_fd(s.to_child);
|
||||
close_fd(s.from_child);
|
||||
if (s.pid > 0) {
|
||||
::kill(s.pid, SIGTERM); // Only our direct child, never a process group.
|
||||
::kill(s.pid, SIGTERM); // Dispatcher handles and reaps its owned child.
|
||||
auto until = Clock::now() + std::chrono::milliseconds(500);
|
||||
int status = 0;
|
||||
while (::waitpid(s.pid, &status, WNOHANG) == 0 && Clock::now() < until) {
|
||||
try { s.drain_debug(); } catch (...) {} // stop/destructor are best effort
|
||||
::usleep(10000);
|
||||
}
|
||||
if (::waitpid(s.pid, &status, WNOHANG) == 0) {
|
||||
::kill(s.pid, SIGKILL);
|
||||
// Spawned in a dedicated group: if dispatcher was stuck or died before
|
||||
// reaping, stop its remaining model process too. Never signal a session.
|
||||
::kill(-s.pid, SIGKILL);
|
||||
if (::waitpid(s.pid, &status, WNOHANG) == 0)
|
||||
while (::waitpid(s.pid, &status, 0) < 0 && errno == EINTR) {}
|
||||
}
|
||||
s.pid = -1;
|
||||
}
|
||||
try { s.drain_debug(); } catch (...) {}
|
||||
@@ -276,10 +277,14 @@ void Worker::stop() {
|
||||
}
|
||||
|
||||
void Worker::start(const std::string& python, const std::string& script,
|
||||
const std::string& model, int threads, bool advanced_debug) {
|
||||
const std::string& model, int threads, bool advanced_debug,
|
||||
const std::string& backend, const std::string& manifest_dir,
|
||||
const std::string& root) {
|
||||
if (state_->pid > 0) throw std::logic_error("worker already started");
|
||||
if (python.empty() || script.empty() || model.empty() || threads < 1 || threads > 64)
|
||||
throw std::invalid_argument("python, script, local model and 1..64 threads required");
|
||||
if (!backend.empty() && (manifest_dir.empty() || root.empty()))
|
||||
throw std::invalid_argument("backend needs manifest directory and release root");
|
||||
check_runtime(::getenv("XDG_RUNTIME_DIR"));
|
||||
try {
|
||||
// The model's internal files are checked by the child before loading, never fetched.
|
||||
@@ -302,16 +307,27 @@ void Worker::start(const std::string& python, const std::string& script,
|
||||
std::vector<std::string> environment;
|
||||
for (char** e = environ; *e; ++e) {
|
||||
std::string_view entry(*e);
|
||||
if (entry.starts_with("PYTHONDONTWRITEBYTECODE=")) continue;
|
||||
if (entry.starts_with("PYTHONDONTWRITEBYTECODE=") ||
|
||||
entry.starts_with("FRAMEYAP_ADVANCED_DEBUG=")) continue;
|
||||
environment.emplace_back(*e);
|
||||
}
|
||||
environment.emplace_back("PYTHONDONTWRITEBYTECODE=1");
|
||||
// Explicit protocol opt-in; dispatcher forwards it to known Redux without
|
||||
// appending unsupported flags to third-party backend launchers.
|
||||
environment.emplace_back(std::string("FRAMEYAP_ADVANCED_DEBUG=") + (advanced_debug ? "1" : "0"));
|
||||
std::vector<char*> envp;
|
||||
for (auto& value : environment) envp.push_back(value.data());
|
||||
envp.push_back(nullptr);
|
||||
const char* args[] = {python.c_str(), script.c_str(), "--model", model.c_str(),
|
||||
"--threads", thread_arg.c_str(), "--clip-dir", state_->dir.c_str(),
|
||||
advanced_debug ? "--advanced-debug" : nullptr, nullptr};
|
||||
std::vector<std::string> arguments{python, script, "--model", model, "--threads", thread_arg,
|
||||
"--clip-dir", state_->dir};
|
||||
if (!backend.empty()) {
|
||||
arguments.insert(arguments.end(), {"--backend", backend, "--manifest-dir", manifest_dir,
|
||||
"--root", root, "--python", python});
|
||||
}
|
||||
if (advanced_debug) arguments.emplace_back("--advanced-debug");
|
||||
std::vector<char*> args;
|
||||
for (auto& argument : arguments) args.push_back(argument.data());
|
||||
args.push_back(nullptr);
|
||||
posix_spawn_file_actions_t actions;
|
||||
int error = posix_spawn_file_actions_init(&actions);
|
||||
if (error) {
|
||||
@@ -325,14 +341,23 @@ void Worker::start(const std::string& python, const std::string& script,
|
||||
if (!error && !advanced_debug) error = posix_spawn_file_actions_addopen(&actions, STDERR_FILENO, "/dev/null", O_WRONLY, 0);
|
||||
for (int fd : std::array<int, 6>{in[0], in[1], out[0], out[1], debug[0], debug[1]})
|
||||
if (!error && fd > STDERR_FILENO) error = posix_spawn_file_actions_addclose(&actions, fd);
|
||||
posix_spawnattr_t attr;
|
||||
bool attr_ready = false;
|
||||
if (!error) {
|
||||
error = posix_spawnattr_init(&attr);
|
||||
attr_ready = !error;
|
||||
}
|
||||
if (!error) error = posix_spawnattr_setflags(&attr, POSIX_SPAWN_SETPGROUP);
|
||||
if (!error) error = posix_spawnattr_setpgroup(&attr, 0);
|
||||
pid_t pid = -1;
|
||||
if (!error) error = posix_spawnp(&pid, python.c_str(), &actions, nullptr,
|
||||
const_cast<char* const*>(args), envp.data());
|
||||
if (!error) error = posix_spawnp(&pid, python.c_str(), &actions, &attr,
|
||||
args.data(), envp.data());
|
||||
if (attr_ready) posix_spawnattr_destroy(&attr);
|
||||
posix_spawn_file_actions_destroy(&actions);
|
||||
if (error) {
|
||||
close_fd(in[0]); close_fd(in[1]); close_fd(out[0]); close_fd(out[1]);
|
||||
close_fd(debug[0]); close_fd(debug[1]);
|
||||
throw std::runtime_error("cannot launch configured Python worker");
|
||||
throw std::runtime_error("cannot launch configured local worker");
|
||||
}
|
||||
close_fd(in[0]); close_fd(out[1]); close_fd(debug[1]);
|
||||
state_->pid = pid; state_->to_child = in[1]; state_->from_child = out[0];
|
||||
@@ -419,18 +444,21 @@ std::optional<WorkerReply> Worker::poll() {
|
||||
}
|
||||
if (type == 'F' && !s.loaded) {
|
||||
if (size == 2 && s.input[5] == 'M')
|
||||
throw std::runtime_error("Missing/mismatched pinned local model weights or private clip directory; check --model");
|
||||
throw std::runtime_error("Missing/mismatched local model files or private clip directory; check --model");
|
||||
if (size == 2 && s.input[5] == 'I')
|
||||
throw std::runtime_error("Authorized Python lacks compatible CPU moondream/torch dependencies; check --python");
|
||||
throw std::runtime_error("Local Redux model failed to load; check authorized CPU runtime and weights");
|
||||
throw std::runtime_error("Configured local worker runtime unavailable; check --python and backend launcher");
|
||||
throw std::runtime_error("Local backend failed to load; check runtime and verified model files");
|
||||
}
|
||||
if ((type != 'R' && type != 'E') || size < 9 || !s.loaded || !s.pending ||
|
||||
get64(s.input.data() + 5) != *s.pending || size - 9 > max_text)
|
||||
throw std::runtime_error("invalid or stale worker reply");
|
||||
WorkerReply reply{*s.pending, {}, {}};
|
||||
std::string text(reinterpret_cast<const char*>(s.input.data() + 13), size - 9);
|
||||
if (!valid_utf8(text)) throw std::runtime_error("invalid worker UTF-8 reply");
|
||||
if (type == 'R') reply.text = std::move(text);
|
||||
// A malformed but correctly framed, correlated transcript is a
|
||||
// request failure, not a broken worker stream. Preserve the loaded
|
||||
// model and allow a fresh recording. Never expose the bad bytes.
|
||||
if (!valid_utf8(text)) reply.error = "transcription failed";
|
||||
else if (type == 'R') reply.text = std::move(text);
|
||||
else reply.error = safe_request_error(text);
|
||||
s.pending.reset(); s.input.clear();
|
||||
::unlink((s.dir + "/clip.raw").c_str());
|
||||
|
||||
+3
-1
@@ -23,7 +23,9 @@ public:
|
||||
Worker(const Worker&) = delete;
|
||||
Worker& operator=(const Worker&) = delete;
|
||||
void start(const std::string& python, const std::string& script,
|
||||
const std::string& model, int threads = 2, bool advanced_debug = false);
|
||||
const std::string& model, int threads = 2, bool advanced_debug = false,
|
||||
const std::string& backend = {}, const std::string& manifest_dir = {},
|
||||
const std::string& root = {});
|
||||
bool ready() const;
|
||||
void submit(uint64_t id, const std::vector<float>& pcm);
|
||||
std::optional<WorkerReply> poll();
|
||||
|
||||
@@ -0,0 +1,210 @@
|
||||
"""Hardware-free generic dispatcher exercise with a second, tiny executable backend."""
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import signal
|
||||
import time
|
||||
from pathlib import Path
|
||||
import struct
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import unittest
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
sys.path.insert(0, str(ROOT / "python"))
|
||||
from frameyap.model_files import load_backends
|
||||
|
||||
|
||||
FAKE = '''#!/usr/bin/env python3
|
||||
import os, struct, sys
|
||||
read = sys.stdin.buffer.read
|
||||
write = sys.stdout.buffer.write
|
||||
assert len(sys.argv) in (3, 4), 'generic backend received unsupported flags'
|
||||
if os.environ.get('FRAMEYAP_ADVANCED_DEBUG') != '0':
|
||||
raise AssertionError('generic launcher must not receive Redux debug opt-in')
|
||||
if os.environ.get('FAKE_PID_PATH'):
|
||||
import signal
|
||||
signal.signal(signal.SIGTERM, signal.SIG_IGN)
|
||||
open(os.environ['FAKE_PID_PATH'], 'w').write(str(os.getpid()))
|
||||
write(struct.pack('<I', 1) + b'Y'); sys.stdout.buffer.flush()
|
||||
while True:
|
||||
header = read(4)
|
||||
if not header: break
|
||||
if os.environ.get('FAKE_PID_PATH'):
|
||||
import time
|
||||
time.sleep(60)
|
||||
size, = struct.unpack('<I', header)
|
||||
payload = read(size)
|
||||
assert payload[:1] == b'T' and len(payload) == 9
|
||||
assert os.path.isfile(sys.argv[2] + '/clip.raw')
|
||||
response = b'R' + payload[1:] + b'fixture transcript'
|
||||
write(struct.pack('<I', len(response)) + response); sys.stdout.buffer.flush()
|
||||
'''
|
||||
|
||||
|
||||
def frame(payload):
|
||||
return struct.pack('<I', len(payload)) + payload
|
||||
|
||||
|
||||
def read_frame(stream):
|
||||
size, = struct.unpack('<I', stream.read(4))
|
||||
return stream.read(size)
|
||||
|
||||
|
||||
class BackendDispatchTest(unittest.TestCase):
|
||||
def test_second_executable_without_cpp_or_model_dependencies(self):
|
||||
with tempfile.TemporaryDirectory() as temp:
|
||||
root = Path(temp)
|
||||
(root / 'assets/backends').mkdir(parents=True)
|
||||
(root / 'python/frameyap').mkdir(parents=True)
|
||||
(root / 'bin').mkdir()
|
||||
(root / 'model/fake').mkdir(parents=True)
|
||||
clip = root / 'clip'
|
||||
clip.mkdir(mode=0o700)
|
||||
(clip / 'clip.raw').write_bytes(b'\0' * (3200 * 4))
|
||||
model = root / 'model/fake'
|
||||
(model / 'weights.dat').write_bytes(b'fixture')
|
||||
fake = root / 'bin/fake'
|
||||
fake.write_text(FAKE)
|
||||
fake.chmod(0o700)
|
||||
manifest = {
|
||||
'schema': 1, 'id': 'fake', 'display_name': 'Fake backend',
|
||||
'launcher': {'type': 'executable', 'path': 'bin/fake',
|
||||
'arguments': ['{model_dir}', '{clip_dir}', '{threads}'],
|
||||
'protocol': 'frameyap-worker-v1'},
|
||||
'model': {'source': 'https://example.org/fake', 'revision': 'test',
|
||||
'files': [{'path': 'weights.dat', 'size': 7,
|
||||
'sha256': hashlib.sha256(b'fixture').hexdigest()}]},
|
||||
'attribution': 'Fixture only', 'license': {'id': 'MIT', 'text': 'Test-only fixture'},
|
||||
'requirements': {'cpu': 'None', 'gpu': 'None'}}
|
||||
(root / 'assets/backends/fake.json').write_text(json.dumps(manifest))
|
||||
self.assertEqual(set(load_backends(root / 'assets/backends')), {'fake'})
|
||||
argv = [sys.executable, str(ROOT / 'python/frameyap/backend_worker.py'),
|
||||
'--backend', 'fake', '--root', str(root), '--python', sys.executable,
|
||||
'--manifest-dir', str(root / 'assets/backends'), '--model', str(model),
|
||||
'--clip-dir', str(clip), '--threads', '2']
|
||||
env = {**os.environ, 'PYTHONDONTWRITEBYTECODE': '1'}
|
||||
with subprocess.Popen(argv, stdin=subprocess.PIPE, stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE, env=env) as child:
|
||||
self.assertEqual(read_frame(child.stdout), b'Y')
|
||||
child.stdin.write(frame(b'T' + (7).to_bytes(8, 'little')))
|
||||
child.stdin.flush()
|
||||
self.assertEqual(read_frame(child.stdout), b'R' + (7).to_bytes(8, 'little') + b'fixture transcript')
|
||||
child.stdin.close()
|
||||
self.assertEqual(child.wait(timeout=5), 0, child.stderr.read())
|
||||
# Debug stays enabled in dispatcher, but arbitrary backend argv/env
|
||||
# receives no Redux-specific flag or transcript logging opt-in.
|
||||
with subprocess.Popen(argv + ['--advanced-debug'], stdin=subprocess.PIPE,
|
||||
stdout=subprocess.PIPE, stderr=subprocess.PIPE, env=env) as child:
|
||||
self.assertEqual(read_frame(child.stdout), b'Y')
|
||||
child.stdin.close()
|
||||
self.assertEqual(child.wait(timeout=5), 0, child.stderr.read())
|
||||
(model / 'weights.dat').write_bytes(b'corrupt')
|
||||
with subprocess.Popen(argv, stdin=subprocess.PIPE, stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE, env=env) as child:
|
||||
self.assertEqual(read_frame(child.stdout), b'FM')
|
||||
child.stdin.close()
|
||||
self.assertEqual(child.wait(timeout=5), 1)
|
||||
|
||||
def test_status_consent_digest_rejected_before_installer_on_same_id_change(self):
|
||||
with tempfile.TemporaryDirectory() as temp:
|
||||
root = Path(temp)
|
||||
manifests = root / 'backends'
|
||||
manifests.mkdir()
|
||||
store = root / 'models'
|
||||
store.mkdir()
|
||||
manifest = {
|
||||
'schema': 1, 'id': 'fake', 'display_name': 'Fake',
|
||||
'launcher': {'type': 'executable', 'path': 'bin/fake',
|
||||
'arguments': ['{model_dir}', '{clip_dir}'], 'protocol': 'frameyap-worker-v1'},
|
||||
'model': {'source': 'https://example.org/original', 'revision': 'v1',
|
||||
'files': [{'path': 'fake', 'size': 1, 'sha256': '0' * 64}]},
|
||||
'attribution': 'Test', 'license': {'id': 'MIT', 'text': 'Fixture'},
|
||||
'requirements': {'cpu': 'None', 'gpu': 'None'}}
|
||||
source = manifests / 'fake.json'
|
||||
source.write_text(json.dumps(manifest))
|
||||
old = hashlib.sha256(source.read_bytes()).hexdigest()
|
||||
service = ROOT / 'scripts/backend-service.py'
|
||||
base = [sys.executable, str(service), '--manifest-dir', str(manifests),
|
||||
'--model-store', str(store)]
|
||||
status = subprocess.run(base + ['--status'], capture_output=True, text=True, check=True)
|
||||
self.assertIn(old, status.stdout)
|
||||
self.assertIn('DONE', status.stdout)
|
||||
marker = root / 'installer-ran'
|
||||
installer = root / 'installer.sh'
|
||||
installer.write_text('#!/bin/sh\nprintf "%s\\n" "$@" > ' + str(marker) + '\n')
|
||||
install = base + ['--install', '--backend', 'fake', '--installer', str(installer),
|
||||
'--expected-manifest-sha256', old]
|
||||
manifest['model']['source'] = 'https://example.org/changed'
|
||||
source.write_text(json.dumps(manifest))
|
||||
wrong = subprocess.run(install, capture_output=True, text=True)
|
||||
self.assertNotEqual(wrong.returncode, 0)
|
||||
self.assertFalse(marker.exists(), wrong.stdout)
|
||||
self.assertNotIn('DONE', wrong.stdout)
|
||||
# Reusing a backend ID cannot silently re-authorize new sources.
|
||||
new = hashlib.sha256(source.read_bytes()).hexdigest()
|
||||
self.assertNotEqual(new, old)
|
||||
self.assertEqual(subprocess.run(install[:-1] + [new], capture_output=True).returncode, 0)
|
||||
self.assertTrue(marker.exists())
|
||||
passed = marker.read_text().splitlines()
|
||||
self.assertEqual(passed[passed.index('--expected-manifest-sha256') + 1], new)
|
||||
|
||||
def test_dispatcher_reaps_child_ignoring_term_on_eof_and_native_style_stop(self):
|
||||
with tempfile.TemporaryDirectory() as temp:
|
||||
root = Path(temp)
|
||||
(root / 'assets/backends').mkdir(parents=True)
|
||||
(root / 'bin').mkdir()
|
||||
model = root / 'model'
|
||||
model.mkdir()
|
||||
(model / 'weights.dat').write_bytes(b'fixture')
|
||||
clip = root / 'clip'
|
||||
clip.mkdir(mode=0o700)
|
||||
(clip / 'clip.raw').write_bytes(b'\0' * 12800)
|
||||
fake = root / 'bin/fake'
|
||||
fake.write_text(FAKE)
|
||||
fake.chmod(0o700)
|
||||
manifest = {'schema': 1, 'id': 'fake', 'display_name': 'Fake',
|
||||
'launcher': {'type': 'executable', 'path': 'bin/fake',
|
||||
'arguments': ['{model_dir}', '{clip_dir}'], 'protocol': 'frameyap-worker-v1'},
|
||||
'model': {'source': 'https://example.org/fake', 'revision': 'test',
|
||||
'files': [{'path': 'weights.dat', 'size': 7,
|
||||
'sha256': hashlib.sha256(b'fixture').hexdigest()}]},
|
||||
'attribution': 'Fixture', 'license': {'id': 'MIT', 'text': 'Fixture'},
|
||||
'requirements': {'cpu': 'None', 'gpu': 'None'}}
|
||||
(root / 'assets/backends/fake.json').write_text(json.dumps(manifest))
|
||||
args = [sys.executable, str(ROOT / 'python/frameyap/backend_worker.py'),
|
||||
'--backend', 'fake', '--root', str(root), '--python', sys.executable,
|
||||
'--manifest-dir', str(root / 'assets/backends'), '--model', str(model),
|
||||
'--clip-dir', str(clip), '--threads', '2']
|
||||
for mode in ('eof', 'term'):
|
||||
pidfile = root / ('pid-' + mode)
|
||||
with subprocess.Popen(args, stdin=subprocess.PIPE, stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE, start_new_session=True,
|
||||
env={**os.environ, 'FAKE_PID_PATH': str(pidfile)}) as child:
|
||||
self.assertEqual(read_frame(child.stdout), b'Y')
|
||||
pid = int(pidfile.read_text())
|
||||
if mode == 'eof':
|
||||
child.stdin.close()
|
||||
else:
|
||||
child.stdin.write(frame(b'T' + (1).to_bytes(8, 'little')))
|
||||
child.stdin.flush()
|
||||
child.terminate()
|
||||
try:
|
||||
code = child.wait(timeout=3)
|
||||
self.assertIn(code, (0, 1), child.stderr.read())
|
||||
for _ in range(100):
|
||||
if not Path('/proc/' + str(pid)).exists():
|
||||
break
|
||||
time.sleep(0.01)
|
||||
self.assertFalse(Path('/proc/' + str(pid)).exists(), 'owned model still running')
|
||||
finally:
|
||||
# Same dedicated group as native. Kill only this test's processes.
|
||||
try:
|
||||
os.killpg(child.pid, signal.SIGKILL)
|
||||
except ProcessLookupError:
|
||||
pass
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
unittest.main()
|
||||
@@ -84,6 +84,20 @@ def fake_child():
|
||||
|
||||
|
||||
class WorkerTests(unittest.TestCase):
|
||||
def test_explicit_dispatcher_debug_environment_reaches_redux(self):
|
||||
# No model imports: the missing model exits after emitting F/M. This
|
||||
# verifies stderr consent while stdout remains an exact protocol frame.
|
||||
with tempfile.TemporaryDirectory() as temp:
|
||||
os.chmod(temp, 0o700)
|
||||
args = [sys.executable, str(Path(worker.__file__)), '--model', str(Path(temp) / 'missing'),
|
||||
'--clip-dir', temp, '--threads', '2']
|
||||
for enabled in ('0', '1'):
|
||||
result = subprocess.run(args, capture_output=True,
|
||||
env={**os.environ, 'FRAMEYAP_ADVANCED_DEBUG': enabled}, timeout=5)
|
||||
self.assertEqual(result.returncode, 1)
|
||||
self.assertEqual(result.stdout, struct.pack('<I', 2) + b'FM')
|
||||
self.assertEqual(b'advanced debugging ON' in result.stderr, enabled == '1')
|
||||
|
||||
def test_framing_and_bounds(self):
|
||||
r, w = os.pipe()
|
||||
try:
|
||||
|
||||
@@ -189,6 +189,10 @@ while True:
|
||||
assert len(request) == 13 and request[4:5] == b'T'
|
||||
if a.model == 'unsafe':
|
||||
value = b'E' + request[5:] + b'transcription failed [inference: RuntimeError] PRIVATE SPEECH'
|
||||
elif a.model == 'bad-utf8' and request[5] == 47:
|
||||
value = b'R' + request[5:] + b'bad\xff'
|
||||
elif a.model == 'control' and request[5] == 47:
|
||||
value = b'R' + request[5:] + b'bad\x1b[31m'
|
||||
elif a.model == 'safe':
|
||||
value = b'E' + request[5:] + b'transcription failed [inference: RuntimeError]'
|
||||
else:
|
||||
@@ -201,6 +205,24 @@ while True:
|
||||
const auto logs = state / "frameyap";
|
||||
const auto current = logs / "worker-debug.log";
|
||||
const auto previous = logs / "worker-debug.previous.log";
|
||||
for (const char* mode : {"bad-utf8", "control"}) {
|
||||
worker.start(argv[1], script, mode, 2);
|
||||
until([&] { worker.poll(); return worker.ready(); });
|
||||
worker.submit(47, clip);
|
||||
until([&] { reply = worker.poll(); return reply.has_value(); });
|
||||
assert(reply->id == 47 && worker.ready());
|
||||
if (std::string(mode) == "bad-utf8") {
|
||||
assert(reply->error == "transcription failed" && reply->text.empty());
|
||||
} else {
|
||||
// Framing permits controls; Controller rejects the transcript
|
||||
// as a request-local failure without stopping this worker.
|
||||
assert(reply->text == "bad\x1b[31m");
|
||||
}
|
||||
worker.submit(48, clip);
|
||||
until([&] { reply = worker.poll(); return reply.has_value(); });
|
||||
assert(reply->text == "ok" && worker.ready());
|
||||
worker.stop();
|
||||
}
|
||||
worker.start(argv[1], script, "safe", 2); // OFF must not create files.
|
||||
until([&] { worker.poll(); return worker.ready(); });
|
||||
worker.submit(300, clip);
|
||||
|
||||
Reference in new issue
Block a user