Add safe transcription errors and opt-in private debug logging

This commit is contained in:
baketnk committed 2026-09-24 15:07:53 -04:00
1 parent 29f3729df7
commit 6592e1af1c
21 files changed
+720 -65

No files matched your search

+52 -9
View File
@@ -13,7 +13,8 @@ namespace {
struct Json {
std::string value;
std::map<std::string, Json> object;
bool is_string = false, is_object = false, is_number = false;
bool is_string = false, is_object = false, is_number = false, is_bool = false;
size_t start = 0, end = 0; // original value span for non-destructive config updates
};
struct Parser {
std::string_view s;
@@ -70,23 +71,25 @@ struct Parser {
ws();
if (pos == s.size()) fail();
Json result;
if (s[pos] == '"') { result.value = str(); result.is_string = true; return result; }
result.start = pos;
auto done = [&]() { result.end = pos; return result; };
if (s[pos] == '"') { result.value = str(); result.is_string = true; return done(); }
if (eat('{')) {
result.is_object = true;
if (eat('}')) return result;
if (eat('}')) return done();
do {
ws(); if (pos == s.size() || s[pos] != '"') fail();
auto key = str();
if (!eat(':')) fail();
auto [it, inserted] = result.object.emplace(std::move(key), parse(depth + 1));
if (!inserted) fail();
if (eat('}')) return result;
if (eat('}')) return done();
} while (eat(','));
fail();
}
if (eat('[')) {
if (eat(']')) return result;
do { parse(depth + 1); if (eat(']')) return result; } while (eat(','));
if (eat(']')) return done();
do { parse(depth + 1); if (eat(']')) return done(); } while (eat(','));
fail();
}
size_t start = pos;
@@ -105,11 +108,15 @@ struct Parser {
}
}
if (pos == start) fail();
if (s.substr(start, pos - start) == "true" || s.substr(start, pos - start) == "false") {
result.is_bool = true;
result.value = s.substr(start, pos - start);
}
if (s[start] == '-' || (s[start] >= '0' && s[start] <= '9')) {
result.is_number = true;
result.value = s.substr(start, pos - start);
}
return result;
return done();
}
};
Rgba color(const Json& json) {
@@ -153,14 +160,16 @@ std::string read_file(const std::filesystem::path& p) {
data.assign(buf, size_t(in.gcount()));
return data;
}
void write_file(const std::filesystem::path& path, const std::string& content) {
void write_file(const std::filesystem::path& path, const std::string& content, bool private_file = false) {
auto temporary = path.string() + ".tmp." + std::to_string(std::chrono::steady_clock::now().time_since_epoch().count());
try {
{
std::ofstream out(temporary, std::ios::binary | std::ios::trunc);
out << content;
if (!out) throw std::runtime_error("Could not write generated OpenVR binding: " + path.string());
if (!out) throw std::runtime_error("Could not write file: " + path.string());
}
if (private_file) std::filesystem::permissions(temporary, std::filesystem::perms::owner_read |
std::filesystem::perms::owner_write, std::filesystem::perm_options::replace);
std::filesystem::rename(temporary, path);
} catch (...) {
std::error_code ignored;
@@ -185,6 +194,9 @@ Config load_config(const std::filesystem::path& path) {
if (key == "font") {
if (!value.is_string) throw std::runtime_error("Config font must be a path string");
config.font = value.value;
} else if (key == "advanced_debug") {
if (!value.is_bool) throw std::runtime_error("Config advanced_debug must be a boolean");
config.advanced_debug = value.value == "true";
} else if (key == "input_priority") {
if (!value.is_string || (value.value != "normal" && value.value != "experimental"))
throw std::runtime_error("Config input_priority must be normal or experimental");
@@ -238,6 +250,37 @@ Config load_config(const std::filesystem::path& path) {
}
return config;
}
bool save_advanced_debug(const std::filesystem::path& path, bool enabled) noexcept {
try {
if (path.empty() || !path.is_absolute() || std::filesystem::is_symlink(path)) return false;
const bool existing = std::filesystem::exists(path);
if (existing && !std::filesystem::is_regular_file(path)) return false;
if (existing) load_config(path); // reject invalid/unknown settings rather than erase customizations
std::string bytes = existing ? read_file(path) : "{}\n";
Parser parser{bytes};
auto root = parser.parse(); parser.ws();
if (!root.is_object || parser.pos != bytes.size()) return false;
const std::string value = enabled ? "true" : "false";
auto it = root.object.find("advanced_debug");
if (it != root.object.end()) {
if (!it->second.is_bool) return false;
if (it->second.value == value) return true;
bytes.replace(it->second.start, it->second.end - it->second.start, value);
} else {
bytes.insert(root.end - 1, std::string(root.object.empty() ? "" : ",") +
"\"advanced_debug\":" + value);
}
// The native reader rejects files >=4097 bytes, even if the JSON is valid.
if (bytes.size() > 4096) return false;
const auto parent = path.parent_path();
if (std::filesystem::is_symlink(parent)) return false;
std::filesystem::create_directories(parent);
write_file(path, bytes, true);
return true;
} catch (...) {
return false;
}
}
std::string resolve_font(const std::string& assets, const std::string& requested) {
const auto bundled = std::filesystem::path(assets) / "fonts/Inconsolata-Regular.ttf";
const std::array<std::filesystem::path, 4> candidates{requested, bundled,
+4
View File
@@ -18,12 +18,16 @@ struct Config {
std::string font; // absolute TTF/OTF path; empty uses the bundled face
// Requests OpenVR's experimental global action priority; SteamVR must allow it too.
bool experimental_input_priority = false;
bool advanced_debug = false; // opt-in full diagnostic logging; never raw audio recording
WristPlacement wrist;
// OpenVR action name -> physical Frame controller input path; empty disables it.
std::map<std::string, std::string> buttons;
};
std::filesystem::path default_config_path();
Config load_config(const std::filesystem::path& path);
// Update only advanced_debug in an existing valid config; false on invalid/unwritable paths.
// Other user customizations and formatting are retained; creates a minimal config if absent.
bool save_advanced_debug(const std::filesystem::path& path, bool enabled) noexcept;
std::string resolve_font(const std::string& assets, const std::string& requested);
// No writes or OpenVR access when no custom button mappings are specified.
// When customized, build a generated manifest and bindings under XDG cache.
+15 -3
View File
@@ -74,7 +74,7 @@ struct Overlay::Impl {
Config config;
PanelSurface surface;
Panel panel;
bool save_failed = false;
bool save_failed = false, debug_save_failed = false;
bool world_ready = false, placed = false, has_texture = false, shown = false;
vr::HmdMatrix34_t world_transform{};
std::optional<Mount> applied_mount;
@@ -151,6 +151,7 @@ struct Overlay::Impl {
laser_change_failed = true;
}
surface.set_lasers_anytime(lasers_anytime);
surface.set_advanced_debug(config.advanced_debug);
overlay_check(overlay->SetOverlayFlag(handle, vr::VROverlayFlags_VisibleInDashboard, true), overlay, "VisibleInDashboard");
vr::HmdVector2_t mouse_scale{{float(W), float(H)}};
overlay_check(overlay->SetOverlayMouseScale(handle, &mouse_scale), overlay, "SetOverlayMouseScale");
@@ -190,7 +191,8 @@ struct Overlay::Impl {
}
if (effective == Mount::World && applied_mount && *applied_mount != Mount::World)
world_ready = false; // a fresh world fallback near the wearer, not an old room location
std::string note = laser_change_failed ? "SteamVR declined the laser mode change." :
std::string note = debug_save_failed ? "Debug preference not saved; using it only for this session." :
laser_change_failed ? "SteamVR declined the laser mode change." :
save_failed ? "Preference could not be saved; using it for this session." : "";
if (effective != mount) note = save_failed ? "Wrist untracked; world fallback. Preference not saved." :
"Wrist not tracked - using world space until it returns.";
@@ -371,8 +373,17 @@ struct Overlay::Impl {
last_pointer_event = "up button=" + std::to_string(event.data.mouse.button);
if (event.data.mouse.button == vr::VRMouseButton_Left) {
auto event_result = surface.pointer_up(event.data.mouse.cursorIndex, event.data.mouse.x, H - event.data.mouse.y);
if (event_result.action || event_result.mount || event_result.recenter || event_result.lasers_anytime || event_result.open_bindings) ++pointer_actions;
if (event_result.action || event_result.mount || event_result.recenter || event_result.lasers_anytime || event_result.open_bindings || event_result.advanced_debug) ++pointer_actions;
if (event_result.action) result.push_back(*event_result.action);
if (event_result.advanced_debug) {
config.advanced_debug = *event_result.advanced_debug;
debug_save_failed = persist_mount && !save_advanced_debug(default_config_path(), config.advanced_debug);
surface.set_advanced_debug(config.advanced_debug);
reset_input(result);
// Runtime observes the change before handling this batch,
// restarts its worker and invalidates pending work/actions.
return result;
}
if (event_result.open_bindings) {
// The editor changes input ownership. Invalidate held gestures
// and pointer presses; never turn the returning release into input.
@@ -456,6 +467,7 @@ Overlay::Overlay(const std::string& assets, const std::string& font, std::option
Overlay::~Overlay() = default;
std::vector<UiAction> Overlay::poll() { return impl_->poll(); }
void Overlay::draw(const Panel& panel) { impl_->draw(panel); }
bool Overlay::advanced_debug() const { return impl_->config.advanced_debug; }
std::string Overlay::controls_status() {
// Compare the same actions across modes, before and after our pose/role gate.
// IsInputAvailable and a successful UpdateActionState are not delivery proof.
+1
View File
@@ -24,6 +24,7 @@ public:
Overlay& operator=(const Overlay&) = delete;
std::vector<UiAction> poll();
void draw(const Panel& panel);
bool advanced_debug() const;
std::string controls_status(); // diagnostic only, no input delivery
std::string pointer_status() const; // diagnostic counters, no input delivery
private:
+24 -12
View File
@@ -26,10 +26,10 @@ struct Rect {
}
};
enum class Control { Review, Settings, Bindings, OpenBindings, Prev, Next, Record, Cancel, Insert, Enter, Quit,
World, Left, Right, Head, Recenter, LasersAnytime };
World, Left, Right, Head, Recenter, LasersAnytime, AdvancedDebug };
enum class Tab { Review, Settings, Bindings };
struct Button { Rect r; Control id; const char* label; };
constexpr std::array<Button, 17> buttons{{
constexpr std::array<Button, 18> buttons{{
{{32, 138, 180, 46}, Control::Review, "Review"},
{{226, 138, 180, 46}, Control::Settings, "Settings"},
{{420, 138, 180, 46}, Control::Bindings, "Bindings"},
@@ -47,6 +47,7 @@ constexpr std::array<Button, 17> buttons{{
{{514, 308, 454, 58}, Control::Right, "Right wrist"},
{{32, 394, 300, 50}, Control::Recenter, "Recenter in front"},
{{514, 394, 454, 50}, Control::LasersAnytime, "Lasers anytime"},
{{32, 452, 936, 48}, Control::AdvancedDebug, "Advanced debugging (full logs)"},
}};
std::optional<UiAction> action(Control c) {
switch (c) {
@@ -105,7 +106,7 @@ struct PanelSurface::Impl {
Theme theme;
Color background, card, ink, muted, cyan, pink;
Tab tab = Tab::Review;
bool dirty = true, lasers_anytime = false;
bool dirty = true, lasers_anytime = false, advanced_debug = false;
std::string placement_note, binding_note;
std::array<std::string, 6> bindings{};
std::array<int, 2> pressed{{-1, -1}};
@@ -244,7 +245,8 @@ struct PanelSurface::Impl {
}
bool visible(Control c) const {
if (c == Control::Prev || c == Control::Next) return tab == Tab::Review;
if (mounting(c) || c == Control::Recenter || c == Control::LasersAnytime) return tab == Tab::Settings;
if (mounting(c) || c == Control::Recenter || c == Control::LasersAnytime ||
c == Control::AdvancedDebug) return tab == Tab::Settings;
if (c == Control::OpenBindings) return tab == Tab::Bindings;
return true;
}
@@ -316,10 +318,9 @@ struct PanelSurface::Impl {
32, 539, 21, muted, 968);
} else {
text("MOUNT AND INTERACTION", 32, 224, 22, muted, 968);
text("Lasers anytime enables system-wide laser mode while this panel is visible.", 32, 472, 20, muted, 968);
text(placement_note.empty() ? "May affect games. Default off; changes saved on this device." : placement_note,
32, 509, 22, muted, 968);
text("If off, open the dashboard to click this setting again.", 32, 538, 20, muted, 968);
text("Full logs may contain speech/text/paths. No saved audio clips.", 32, 520, 20, pink, 968);
text(placement_note.empty() ? "Toggle restarts worker; cancels current work. Lasers may affect games." : placement_note,
32, 540, 20, muted, 968);
}
rect({32, 550, 936, 1}, mix(card, cyan, .17f));
for (size_t i = 0; i < buttons.size(); ++i) {
@@ -330,7 +331,8 @@ struct PanelSurface::Impl {
(b.id == Control::Settings && tab == Tab::Settings) ||
(b.id == Control::Bindings && tab == Tab::Bindings) ||
(mounting(b.id) && *mounting(b.id) == mount) ||
(b.id == Control::LasersAnytime && lasers_anytime);
(b.id == Control::LasersAnytime && lasers_anytime) ||
(b.id == Control::AdvancedDebug && advanced_debug);
const Color fill = !on ? mix(background, card, .40f) :
selected ? mix(card, cyan, .14f) : card;
const Color accent = b.id == Control::Record && panel.recording ? pink : cyan;
@@ -342,9 +344,11 @@ struct PanelSurface::Impl {
const auto label = b.id == Control::Record && panel.recording ? "Stop" : b.label;
text(label, b.r.x + 16, b.r.y + b.r.h / 2 + 9, 27, on ? ink : mix(background, muted, .48f), b.r.x + b.r.w - 8);
if (mounting(b.id) && selected) text("ON", b.r.x + b.r.w - 56, b.r.y + 38, 23, cyan, b.r.x + b.r.w - 12);
if (b.id == Control::LasersAnytime)
text(lasers_anytime ? "ON" : "OFF", b.r.x + b.r.w - 66, b.r.y + 34, 22,
lasers_anytime ? cyan : muted, b.r.x + b.r.w - 12);
if (b.id == Control::LasersAnytime || b.id == Control::AdvancedDebug) {
bool active = b.id == Control::LasersAnytime ? lasers_anytime : advanced_debug;
text(active ? "ON" : "OFF", b.r.x + b.r.w - 66, b.r.y + b.r.h / 2 + 9, 22,
active ? cyan : muted, b.r.x + b.r.w - 12);
}
}
dirty = false;
return true;
@@ -375,6 +379,7 @@ SurfaceEvent PanelSurface::pointer_up(unsigned cursor, float x, float y) {
else if (auto m = mounting(c)) { impl_->mount = *m; result.mount = *m; impl_->reset(); impl_->dirty = true; }
else if (c == Control::Recenter) { result.recenter = true; impl_->reset(); }
else if (c == Control::LasersAnytime) result.lasers_anytime = !impl_->lasers_anytime;
else if (c == Control::AdvancedDebug) result.advanced_debug = !impl_->advanced_debug;
else if (c == Control::OpenBindings) { result.open_bindings = true; impl_->reset(); }
else if (c == Control::Review || c == Control::Settings || c == Control::Bindings) {
impl_->tab = c == Control::Review ? Tab::Review : c == Control::Settings ? Tab::Settings : Tab::Bindings;
@@ -402,4 +407,11 @@ void PanelSurface::set_lasers_anytime(bool enabled) {
impl_->dirty = true;
}
}
void PanelSurface::set_advanced_debug(bool enabled) {
if (impl_->advanced_debug != enabled) {
impl_->advanced_debug = enabled;
impl_->reset();
impl_->dirty = true;
}
}
} // namespace frameyap
+2
View File
@@ -12,6 +12,7 @@ struct SurfaceEvent {
std::optional<UiAction> action;
std::optional<Mount> mount;
std::optional<bool> lasers_anytime;
std::optional<bool> advanced_debug;
bool recenter = false;
bool open_bindings = false;
};
@@ -32,6 +33,7 @@ public:
void reset_pointers();
void set_placement_note(std::string note);
void set_lasers_anytime(bool enabled);
void set_advanced_debug(bool enabled);
// PTT, Cancel, Insert, Enter, left-grip gesture, right-grip gesture.
void set_bindings(std::array<std::string, 6> labels);
void set_binding_note(std::string note);
+16 -2
View File
@@ -11,6 +11,7 @@
#include <cstdlib>
#include <filesystem>
#include <fcntl.h>
#include <iostream>
#include <stdexcept>
#include <sys/file.h>
#include <sys/stat.h>
@@ -54,6 +55,7 @@ int run(const Options& options) {
Session session;
std::string detail = "Review mode. Other apps may also hear your mic. Enter is explicit.";
bool quit = false;
bool advanced_debug = overlay.advanced_debug();
const DeliveryFactory acquire = [&]() -> std::unique_ptr<DeliveryLease> {
auto input = std::make_unique<NativeDeliveryLease>(options.socket);
if (interrupted) throw std::runtime_error("Input cancelled before delivery");
@@ -77,7 +79,7 @@ int run(const Options& options) {
};
auto warm = [&] {
audio.close(); worker.stop(); session = Session{};
worker.start(options.python, options.worker, options.model, options.threads);
worker.start(options.python, options.worker, options.model, options.threads, advanced_debug);
detail = "Loading local model; microphone closed. Record again when Ready.";
};
auto stop_record = [&] {
@@ -108,6 +110,9 @@ int run(const Options& options) {
// E is a request-local error. The child still owns its
// loaded model and can accept the next utterance.
session.fail(); detail = reply->error + "; model ready. Record to retry.";
// Worker::poll allows only fixed diagnostic labels here,
// never exception messages, audio, paths or recognized text.
std::cerr << "FrameYap worker: " << reply->error << '\n';
} else {
session.reply(reply->id, reply->text);
detail = session.text().empty() ? "No speech recognized; try again." : "Focus your destination, then Insert. Cancel discards.";
@@ -139,7 +144,16 @@ int run(const Options& options) {
session.state() != State::Warming && session.state() != State::Transcribing && session.state() != State::Review};
if (panel.recording) panel.status += " - " + std::to_string(audio.seconds()) + " / 20s";
overlay.draw(panel);
for (auto action : overlay.poll()) {
auto actions = overlay.poll();
if (advanced_debug != overlay.advanced_debug()) {
advanced_debug = overlay.advanced_debug();
// Consent changes take effect before any more work or delivery. The
// Settings warning makes the worker restart/cancellation explicit.
try { warm(); }
catch (const std::exception& e) { session.fail(); detail = e.what(); }
std::erase_if(actions, [](UiAction action) { return action != UiAction::Quit; });
}
for (auto action : actions) {
try {
switch (action) {
case UiAction::Quit: quit = true; break;
+162 -13
View File
@@ -1,5 +1,6 @@
#include "worker.hpp"
#include <algorithm>
#include <array>
#include <cerrno>
#include <chrono>
@@ -26,6 +27,8 @@ namespace {
using Clock = std::chrono::steady_clock;
constexpr size_t max_frame = 65536;
constexpr size_t max_text = 4096;
constexpr size_t max_debug_log = 4 * 1024 * 1024;
constexpr std::string_view truncation = "\n[worker diagnostics truncated; further output discarded]\n";
void close_fd(int& fd) { if (fd >= 0) { ::close(fd); fd = -1; } }
void put32(unsigned char* p, uint32_t n) {
@@ -62,6 +65,102 @@ bool valid_utf8(std::string_view text) {
}
return true;
}
std::string safe_request_error(std::string_view error) {
// Treat E payloads as untrusted too: older/alternate workers might include
// raw exception messages. Only fixed diagnostic labels reach UI or logs.
for (const auto* stage : {"audio", "inference", "response", "unknown"}) {
for (const auto* category : {"MemoryError", "ImportError", "OSError", "UnicodeError", "TypeError",
"ValueError", "KeyError", "IndexError", "RuntimeError", "Exception"}) {
auto allowed = std::string("transcription failed [") + stage + ": " + category + "]";
if (error == allowed) return allowed;
}
}
return "transcription failed"; // compatible legacy/unknown error, never raw text
}
// Walk from / using directory fds: no symlinks, even in intermediate components.
// Root-owned system ancestors (including sticky /tmp) are permitted, but a
// writable non-sticky ancestor could redirect private diagnostics elsewhere.
int state_directory() {
const char* base = ::getenv("XDG_STATE_HOME");
std::string path;
if (base && *base) path = base;
else {
const char* home = ::getenv("HOME");
if (!home || !*home) throw std::runtime_error("advanced debug needs HOME or XDG_STATE_HOME");
path = std::string(home) + "/.local/state";
}
if (path.empty() || path[0] != '/')
throw std::runtime_error("advanced debug state path must be absolute");
while (path.size() > 1 && path.back() == '/') path.pop_back();
path += path == "/" ? "frameyap" : "/frameyap";
int dir = ::open("/", O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (dir < 0) throw std::runtime_error("cannot open debug state root");
try {
size_t pos = 1;
while (pos < path.size()) {
auto end = path.find('/', pos);
if (end == std::string::npos) end = path.size();
auto part = path.substr(pos, end - pos);
if (part.empty() || part == "." || part == "..")
throw std::runtime_error("unsafe debug state path");
if (::mkdirat(dir, part.c_str(), 0700) && errno != EEXIST)
throw std::runtime_error("cannot create private debug state directory");
int next = ::openat(dir, part.c_str(), O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (next < 0) throw std::runtime_error("unsafe debug state directory");
struct stat st{};
bool final = end == path.size();
if (::fstat(next, &st) || !S_ISDIR(st.st_mode) ||
(final && (st.st_uid != ::geteuid() || (st.st_mode & 07777) != 0700)) ||
(!final && ((st.st_uid != ::geteuid() && st.st_uid != 0) ||
((st.st_mode & 0022) && !(st.st_uid == 0 && (st.st_mode & S_ISVTX)))))) {
::close(next);
throw std::runtime_error("unsafe debug state directory permissions");
}
::close(dir); dir = next;
pos = end + 1;
}
return dir;
} catch (...) { ::close(dir); throw; }
}
// Open existing files without following links; never rotate someone else's file,
// a hardlink or a permissive target. Both names are checked before any mutation.
void check_log_target(int dir, const char* name) {
int fd = ::openat(dir, name, O_RDONLY | O_NOFOLLOW | O_NONBLOCK | O_CLOEXEC);
if (fd < 0) {
if (errno == ENOENT) return;
throw std::runtime_error("unsafe existing worker debug log");
}
struct stat st{};
bool valid = ::fstat(fd, &st) == 0 && S_ISREG(st.st_mode) &&
st.st_uid == ::geteuid() && (st.st_mode & 07777) == 0600 &&
st.st_nlink == 1 && st.st_size <= static_cast<off_t>(max_debug_log);
::close(fd);
if (!valid) throw std::runtime_error("unsafe existing worker debug log");
}
int create_debug_log() {
int dir = state_directory();
try {
check_log_target(dir, "worker-debug.log");
check_log_target(dir, "worker-debug.previous.log");
if (::unlinkat(dir, "worker-debug.previous.log", 0) && errno != ENOENT)
throw std::runtime_error("cannot rotate worker debug log");
if (::renameat(dir, "worker-debug.log", dir, "worker-debug.previous.log") && errno != ENOENT)
throw std::runtime_error("cannot rotate worker debug log");
int fd = ::openat(dir, "worker-debug.log", O_WRONLY | O_CREAT | O_EXCL | O_NOFOLLOW | O_CLOEXEC, 0600);
if (fd < 0) throw std::runtime_error("cannot create private worker debug log");
::close(dir);
return fd;
} catch (...) { ::close(dir); throw; }
}
void log_bytes(int fd, const unsigned char* data, size_t size) {
while (size) {
ssize_t n = ::write(fd, data, size);
if (n < 0 && errno == EINTR) continue;
if (n <= 0) throw std::runtime_error("worker debug log write failed");
data += n; size -= static_cast<size_t>(n);
}
}
void check_runtime(const char* path) {
if (!path || path[0] != '/') throw std::runtime_error("XDG_RUNTIME_DIR must be an absolute private directory");
struct stat st{};
@@ -97,13 +196,45 @@ void write_all(int fd, const unsigned char* data, size_t size) {
struct Worker::State {
pid_t pid = -1;
int to_child = -1, from_child = -1;
int to_child = -1, from_child = -1, debug_pipe = -1, debug_file = -1;
size_t debug_written = 0;
bool debug_truncated = false;
std::string dir;
bool loaded = false;
std::optional<uint64_t> pending;
Clock::time_point deadline{};
std::vector<unsigned char> input;
std::chrono::milliseconds warmup_timeout{120000}, request_timeout{60000};
void drain_debug() {
if (debug_pipe < 0) return;
// Limit work per poll: even an endlessly noisy child cannot trap the UI.
size_t budget = 256 * 1024;
unsigned char block[4096];
while (budget) {
ssize_t n = ::read(debug_pipe, block, std::min(budget, sizeof block));
if (n > 0) {
budget -= static_cast<size_t>(n);
if (!debug_truncated) {
size_t left = max_debug_log - truncation.size() - debug_written;
size_t count = std::min(left, static_cast<size_t>(n));
log_bytes(debug_file, block, count);
debug_written += count;
if (count < static_cast<size_t>(n)) {
log_bytes(debug_file, reinterpret_cast<const unsigned char*>(truncation.data()), truncation.size());
debug_written += truncation.size();
debug_truncated = true;
}
}
continue;
}
if (n == 0) { close_fd(debug_pipe); return; }
if (errno == EINTR) continue;
if (errno != EAGAIN && errno != EWOULDBLOCK)
throw std::runtime_error("worker debug pipe read failed");
return;
}
}
};
Worker::Worker(std::chrono::milliseconds warmup, std::chrono::milliseconds request)
@@ -122,14 +253,20 @@ void Worker::stop() {
::kill(s.pid, SIGTERM); // Only our direct child, never a process group.
auto until = Clock::now() + std::chrono::milliseconds(500);
int status = 0;
while (::waitpid(s.pid, &status, WNOHANG) == 0 && Clock::now() < until)
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);
while (::waitpid(s.pid, &status, 0) < 0 && errno == EINTR) {}
}
s.pid = -1;
}
try { s.drain_debug(); } catch (...) {}
close_fd(s.debug_pipe);
close_fd(s.debug_file);
s.debug_written = 0; s.debug_truncated = false;
if (!s.dir.empty()) {
::unlink((s.dir + "/clip.raw").c_str());
::rmdir(s.dir.c_str());
@@ -139,7 +276,7 @@ void Worker::stop() {
}
void Worker::start(const std::string& python, const std::string& script,
const std::string& model, int threads) {
const std::string& model, int threads, bool advanced_debug) {
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");
@@ -151,9 +288,12 @@ void Worker::start(const std::string& python, const std::string& script,
if (!::mkdtemp(tmp.data())) throw std::runtime_error("cannot create private clip directory");
state_->dir = tmp.data();
::chmod(state_->dir.c_str(), 0700);
int in[2]{-1,-1}, out[2]{-1,-1};
if (::pipe2(in, O_CLOEXEC) || ::pipe2(out, O_CLOEXEC)) {
if (advanced_debug) state_->debug_file = create_debug_log();
int in[2]{-1,-1}, out[2]{-1,-1}, debug[2]{-1,-1};
if (::pipe2(in, O_CLOEXEC) || ::pipe2(out, O_CLOEXEC) ||
(advanced_debug && ::pipe2(debug, O_CLOEXEC))) {
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 create worker pipes");
}
std::string thread_arg = std::to_string(threads);
@@ -170,17 +310,20 @@ void Worker::start(const std::string& python, const std::string& script,
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(), nullptr};
"--threads", thread_arg.c_str(), "--clip-dir", state_->dir.c_str(),
advanced_debug ? "--advanced-debug" : nullptr, nullptr};
posix_spawn_file_actions_t actions;
int error = posix_spawn_file_actions_init(&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 initialize worker spawn");
}
error = posix_spawn_file_actions_adddup2(&actions, in[0], STDIN_FILENO);
if (!error) error = posix_spawn_file_actions_adddup2(&actions, out[1], STDOUT_FILENO);
if (!error) error = posix_spawn_file_actions_addopen(&actions, STDERR_FILENO, "/dev/null", O_WRONLY, 0);
for (int fd : std::array<int, 4>{in[0], in[1], out[0], out[1]})
if (!error && advanced_debug) error = posix_spawn_file_actions_adddup2(&actions, debug[1], STDERR_FILENO);
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);
pid_t pid = -1;
if (!error) error = posix_spawnp(&pid, python.c_str(), &actions, nullptr,
@@ -188,13 +331,18 @@ void Worker::start(const std::string& python, const std::string& script,
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");
}
close_fd(in[0]); close_fd(out[1]);
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];
int flags = ::fcntl(state_->from_child, F_GETFL);
if (flags < 0 || ::fcntl(state_->from_child, F_SETFL, flags | O_NONBLOCK))
throw std::runtime_error("cannot set nonblocking worker pipe");
state_->debug_pipe = debug[0];
for (int fd : {state_->from_child, state_->debug_pipe}) {
if (fd < 0) continue;
int flags = ::fcntl(fd, F_GETFL);
if (flags < 0 || ::fcntl(fd, F_SETFL, flags | O_NONBLOCK))
throw std::runtime_error("cannot set nonblocking worker pipe");
}
state_->deadline = Clock::now() + state_->warmup_timeout;
} catch (...) { stop(); throw; }
}
@@ -240,6 +388,7 @@ std::optional<WorkerReply> Worker::poll() {
auto& s = *state_;
if (s.pid <= 0) return std::nullopt;
try {
s.drain_debug();
if ((!s.loaded || s.pending) && Clock::now() > s.deadline)
throw std::runtime_error(s.loaded ? "worker transcription timed out" : "worker warmup timed out");
unsigned char block[4096];
@@ -280,7 +429,7 @@ std::optional<WorkerReply> Worker::poll() {
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);
else reply.error = std::move(text);
else reply.error = safe_request_error(text);
s.pending.reset(); s.input.clear();
::unlink((s.dir + "/clip.raw").c_str());
return reply;
+1 -1
View File
@@ -23,7 +23,7 @@ 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);
const std::string& model, int threads = 2, bool advanced_debug = false);
bool ready() const;
void submit(uint64_t id, const std::vector<float>& pcm);
std::optional<WorkerReply> poll();