mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-10 10:06:15 -08:00
feat(embedded-mpv): frame-copy helper process and shm frame reader (native layer)
iptvnator_mpv_helper: one-process-per-session libmpv host that renders offscreen at viewport size (headless CGL + async PBO ring, validated in spikes/mpv-frame-copy), publishes BGRA frames into a seqlock shm ring with resize generations, plays audio directly, and speaks a stdio protocol — tab-separated commands in, JSON events out. The snapshot event mirrors NativeEmbeddedMpvSessionSnapshot; status semantics (END_FILE reasons, eof-reached with keep-open, pause gated on loaded path, fatal-only status flips) are ported from embedded_mpv.mm. embedded_mpv_frame_reader.node: plain-C N-API reader the preload script uses to memcpy the newest complete frame into a V8 ArrayBuffer (Electron's memory cage forbids zero-copy). Stub exports off macOS. Both build as extra binding.gyp targets through build-embedded-mpv.js; the helper gets the same libmpv dependency-path rewrite + ad-hoc re-sign as the addon and is validated by the forbidden-link check. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
1 parent
b0be6d9b8b
commit
aabce726e3
7 files changed
+1577
No files matched your search
@@ -473,8 +473,16 @@ function main() {
|
||||
|
||||
if (targetPlatform === 'darwin') {
|
||||
patchAddonForBundledRuntime(outputFile, outputLibDir);
|
||||
// The frame-copy helper executable links libmpv too and sits next to
|
||||
// the same lib/ directory, so it gets the identical dependency-path
|
||||
// rewrite + ad-hoc re-sign.
|
||||
const frameHelperFile = path.join(outputDir, 'iptvnator_mpv_helper');
|
||||
if (fs.existsSync(frameHelperFile)) {
|
||||
patchAddonForBundledRuntime(frameHelperFile, outputLibDir);
|
||||
}
|
||||
const forbiddenLinkErrors = validateNoForbiddenRuntimeLinks([
|
||||
outputFile,
|
||||
...(fs.existsSync(frameHelperFile) ? [frameHelperFile] : []),
|
||||
...runtimeManifest.dylibs.map((dylib) =>
|
||||
path.join(outputLibDir, dylib)
|
||||
),
|
||||
|
||||
@@ -75,6 +75,50 @@
|
||||
}
|
||||
]
|
||||
]
|
||||
},
|
||||
{
|
||||
"target_name": "embedded_mpv_frame_reader",
|
||||
"sources": [
|
||||
"src/embedded_mpv_frame_reader.c"
|
||||
],
|
||||
"include_dirs": [
|
||||
"helper"
|
||||
]
|
||||
},
|
||||
{
|
||||
"target_name": "iptvnator_mpv_helper",
|
||||
"type": "none",
|
||||
"conditions": [
|
||||
[
|
||||
"OS==\"mac\"",
|
||||
{
|
||||
"type": "executable",
|
||||
"sources": [
|
||||
"helper/mpv_frame_helper.cpp"
|
||||
],
|
||||
"include_dirs": [
|
||||
"<!(node -p \"process.env.LIBMPV_INCLUDE_DIR || '/opt/homebrew/include'\")",
|
||||
"helper"
|
||||
],
|
||||
"cflags_cc": [
|
||||
"-std=c++17"
|
||||
],
|
||||
"xcode_settings": {
|
||||
"CLANG_CXX_LANGUAGE_STANDARD": "c++17",
|
||||
"GCC_ENABLE_CPP_EXCEPTIONS": "YES",
|
||||
"MACOSX_DEPLOYMENT_TARGET": "11.0",
|
||||
"OTHER_LDFLAGS": [
|
||||
"-Wl,-rpath,@executable_path/lib",
|
||||
"-Wl,-rpath,@loader_path/lib"
|
||||
]
|
||||
},
|
||||
"libraries": [
|
||||
"-framework OpenGL",
|
||||
"<!(node -e \"const path = require('path'); const dir = process.env.LIBMPV_LIBRARY_DIR || '/opt/homebrew/lib'; process.stdout.write(path.join(dir, 'libmpv.2.dylib'))\")"
|
||||
]
|
||||
}
|
||||
]
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -0,0 +1,176 @@
|
||||
/*
|
||||
* stdio line protocol for the frame-copy helper.
|
||||
*
|
||||
* Inbound commands (stdin): one command per line, tab-separated fields.
|
||||
* The first field is the command name; the rest are `key=value` pairs with
|
||||
* percent-encoded values (%XX for at least '%', TAB, CR, LF). This avoids a
|
||||
* JSON parser in native code — the TypeScript adapter owns serialization.
|
||||
*
|
||||
* Outbound events (stdout): one JSON object per line. Emitting JSON only
|
||||
* needs string escaping, which lives here too.
|
||||
*/
|
||||
#pragma once
|
||||
|
||||
#include <cstdint>
|
||||
#include <cstdio>
|
||||
#include <mutex>
|
||||
#include <sstream>
|
||||
#include <string>
|
||||
#include <unordered_map>
|
||||
#include <vector>
|
||||
|
||||
namespace frame_helper {
|
||||
|
||||
inline int hexValue(char c) {
|
||||
if (c >= '0' && c <= '9') return c - '0';
|
||||
if (c >= 'a' && c <= 'f') return c - 'a' + 10;
|
||||
if (c >= 'A' && c <= 'F') return c - 'A' + 10;
|
||||
return -1;
|
||||
}
|
||||
|
||||
inline std::string percentDecode(const std::string& value) {
|
||||
std::string out;
|
||||
out.reserve(value.size());
|
||||
for (size_t i = 0; i < value.size(); i++) {
|
||||
if (value[i] == '%' && i + 2 < value.size() &&
|
||||
hexValue(value[i + 1]) >= 0 && hexValue(value[i + 2]) >= 0) {
|
||||
out.push_back(
|
||||
(char)((hexValue(value[i + 1]) << 4) | hexValue(value[i + 2])));
|
||||
i += 2;
|
||||
continue;
|
||||
}
|
||||
out.push_back(value[i]);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
struct Command {
|
||||
std::string name;
|
||||
std::unordered_map<std::string, std::string> args;
|
||||
|
||||
std::string get(const std::string& key,
|
||||
const std::string& fallback = "") const {
|
||||
const auto it = args.find(key);
|
||||
return it == args.end() ? fallback : it->second;
|
||||
}
|
||||
|
||||
double getDouble(const std::string& key, double fallback) const {
|
||||
const auto it = args.find(key);
|
||||
if (it == args.end()) return fallback;
|
||||
try {
|
||||
return std::stod(it->second);
|
||||
} catch (...) {
|
||||
return fallback;
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
inline Command parseCommandLine(const std::string& line) {
|
||||
Command command;
|
||||
std::stringstream stream(line);
|
||||
std::string field;
|
||||
bool first = true;
|
||||
while (std::getline(stream, field, '\t')) {
|
||||
if (first) {
|
||||
command.name = field;
|
||||
first = false;
|
||||
continue;
|
||||
}
|
||||
const size_t eq = field.find('=');
|
||||
if (eq == std::string::npos) continue;
|
||||
command.args[field.substr(0, eq)] =
|
||||
percentDecode(field.substr(eq + 1));
|
||||
}
|
||||
return command;
|
||||
}
|
||||
|
||||
inline std::string jsonEscape(const std::string& value) {
|
||||
std::string out;
|
||||
out.reserve(value.size() + 8);
|
||||
for (const char c : value) {
|
||||
switch (c) {
|
||||
case '"': out += "\\\""; break;
|
||||
case '\\': out += "\\\\"; break;
|
||||
case '\n': out += "\\n"; break;
|
||||
case '\r': out += "\\r"; break;
|
||||
case '\t': out += "\\t"; break;
|
||||
default:
|
||||
if ((unsigned char)c < 0x20) {
|
||||
char buf[8];
|
||||
std::snprintf(buf, sizeof(buf), "\\u%04x", c);
|
||||
out += buf;
|
||||
} else {
|
||||
out.push_back(c);
|
||||
}
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/* Minimal JSON writer for flat/one-level-nested event objects. */
|
||||
class JsonWriter {
|
||||
public:
|
||||
JsonWriter() { buffer_ += '{'; }
|
||||
|
||||
JsonWriter& str(const std::string& key, const std::string& value) {
|
||||
writeKey(key);
|
||||
buffer_ += '"';
|
||||
buffer_ += jsonEscape(value);
|
||||
buffer_ += '"';
|
||||
return *this;
|
||||
}
|
||||
|
||||
JsonWriter& num(const std::string& key, double value) {
|
||||
writeKey(key);
|
||||
char buf[64];
|
||||
std::snprintf(buf, sizeof(buf), "%.6g", value);
|
||||
buffer_ += buf;
|
||||
return *this;
|
||||
}
|
||||
|
||||
JsonWriter& boolean(const std::string& key, bool value) {
|
||||
writeKey(key);
|
||||
buffer_ += value ? "true" : "false";
|
||||
return *this;
|
||||
}
|
||||
|
||||
JsonWriter& nullValue(const std::string& key) {
|
||||
writeKey(key);
|
||||
buffer_ += "null";
|
||||
return *this;
|
||||
}
|
||||
|
||||
JsonWriter& raw(const std::string& key, const std::string& rawJson) {
|
||||
writeKey(key);
|
||||
buffer_ += rawJson;
|
||||
return *this;
|
||||
}
|
||||
|
||||
std::string finish() {
|
||||
buffer_ += '}';
|
||||
return buffer_;
|
||||
}
|
||||
|
||||
private:
|
||||
void writeKey(const std::string& key) {
|
||||
if (!first_) buffer_ += ',';
|
||||
first_ = false;
|
||||
buffer_ += '"';
|
||||
buffer_ += jsonEscape(key);
|
||||
buffer_ += "\":";
|
||||
}
|
||||
|
||||
std::string buffer_;
|
||||
bool first_ = true;
|
||||
};
|
||||
|
||||
/* stdout is shared by every thread that emits events. */
|
||||
inline void emitLine(const std::string& jsonLine) {
|
||||
static std::mutex stdoutMutex;
|
||||
std::lock_guard<std::mutex> lock(stdoutMutex);
|
||||
std::fwrite(jsonLine.data(), 1, jsonLine.size(), stdout);
|
||||
std::fputc('\n', stdout);
|
||||
std::fflush(stdout);
|
||||
}
|
||||
|
||||
} // namespace frame_helper
|
||||
@@ -0,0 +1,403 @@
|
||||
/*
|
||||
* Offscreen render + readback pipeline for the frame-copy helper (macOS).
|
||||
*
|
||||
* Headless CGL context, mpv render API into an FBO sized to the viewport,
|
||||
* async PBO readback ring, BGRA frames published into the FrameShm ring.
|
||||
* Validated in spikes/mpv-frame-copy (4K60 sustained on M1 Pro).
|
||||
*
|
||||
* Threading: everything here runs on the render thread except
|
||||
* requestResize()/stop(), which only touch atomics/mutex-guarded state.
|
||||
*/
|
||||
#pragma once
|
||||
|
||||
#define GL_SILENCE_DEPRECATION
|
||||
#include <OpenGL/OpenGL.h>
|
||||
#include <OpenGL/gl3.h>
|
||||
|
||||
#include <mpv/client.h>
|
||||
#include <mpv/render_gl.h>
|
||||
|
||||
#include <dlfcn.h>
|
||||
#include <fcntl.h>
|
||||
#include <sys/mman.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <condition_variable>
|
||||
#include <cstring>
|
||||
#include <functional>
|
||||
#include <mutex>
|
||||
#include <string>
|
||||
|
||||
#include "frame_helper_io.h"
|
||||
#include "frame_shm.h"
|
||||
|
||||
namespace frame_helper {
|
||||
|
||||
inline uint64_t nowNs() { return clock_gettime_nsec_np(CLOCK_MONOTONIC_RAW); }
|
||||
|
||||
struct ShmRing {
|
||||
FrameShmHeader* header = nullptr;
|
||||
uint8_t* base = nullptr;
|
||||
size_t size = 0;
|
||||
std::string name;
|
||||
|
||||
bool create(const std::string& shmName, int width, int height,
|
||||
uint32_t generation) {
|
||||
destroy();
|
||||
shm_unlink(shmName.c_str());
|
||||
const int fd = shm_open(shmName.c_str(), O_CREAT | O_EXCL | O_RDWR, 0600);
|
||||
if (fd < 0) return false;
|
||||
const uint64_t frameBytes = (uint64_t)width * 4u * (uint64_t)height;
|
||||
const uint64_t dataOffset =
|
||||
(sizeof(FrameShmHeader) + FRAME_SHM_DATA_ALIGN - 1) &
|
||||
~(uint64_t)(FRAME_SHM_DATA_ALIGN - 1);
|
||||
const size_t total =
|
||||
(size_t)(dataOffset + FRAME_SHM_RING_SLOTS * frameBytes);
|
||||
if (ftruncate(fd, (off_t)total) != 0) {
|
||||
close(fd);
|
||||
shm_unlink(shmName.c_str());
|
||||
return false;
|
||||
}
|
||||
void* mapped =
|
||||
mmap(nullptr, total, PROT_READ | PROT_WRITE, MAP_SHARED, fd, 0);
|
||||
close(fd);
|
||||
if (mapped == MAP_FAILED) {
|
||||
shm_unlink(shmName.c_str());
|
||||
return false;
|
||||
}
|
||||
std::memset(mapped, 0, sizeof(FrameShmHeader));
|
||||
auto* hdr = static_cast<FrameShmHeader*>(mapped);
|
||||
hdr->version = FRAME_SHM_VERSION;
|
||||
hdr->width = (uint32_t)width;
|
||||
hdr->height = (uint32_t)height;
|
||||
hdr->stride = (uint32_t)width * 4u;
|
||||
hdr->generation = generation;
|
||||
hdr->frame_bytes = frameBytes;
|
||||
hdr->data_offset = dataOffset;
|
||||
std::atomic_thread_fence(std::memory_order_release);
|
||||
hdr->magic = FRAME_SHM_MAGIC;
|
||||
|
||||
header = hdr;
|
||||
base = static_cast<uint8_t*>(mapped);
|
||||
size = total;
|
||||
name = shmName;
|
||||
return true;
|
||||
}
|
||||
|
||||
uint8_t* slotData(uint64_t seq) const {
|
||||
return base + header->data_offset +
|
||||
(seq % FRAME_SHM_RING_SLOTS) * header->frame_bytes;
|
||||
}
|
||||
|
||||
void destroy() {
|
||||
if (base) {
|
||||
munmap(base, size);
|
||||
shm_unlink(name.c_str());
|
||||
}
|
||||
header = nullptr;
|
||||
base = nullptr;
|
||||
size = 0;
|
||||
name.clear();
|
||||
}
|
||||
};
|
||||
|
||||
class RenderPipeline {
|
||||
public:
|
||||
/* Called from the render thread after a resize created a new shm
|
||||
* generation, so the main protocol layer can announce it. */
|
||||
std::function<void(const std::string& name, int width, int height,
|
||||
uint32_t generation)>
|
||||
onGenerationChanged;
|
||||
|
||||
bool start(mpv_handle* mpv, const std::string& shmBaseName, int width,
|
||||
int height, std::string& errorOut);
|
||||
void requestResize(int width, int height);
|
||||
void stop();
|
||||
void notifyUpdate(); /* mpv render update callback -> wake render loop */
|
||||
void runLoop(); /* render thread body */
|
||||
|
||||
mpv_render_context* renderContext() { return renderContext_; }
|
||||
/* 0 = initializing, 1 = ready, -1 = failed */
|
||||
int initState() const { return initState_.load(); }
|
||||
|
||||
private:
|
||||
bool setupGl(std::string& errorOut);
|
||||
bool rebuildTargets(int width, int height);
|
||||
void publishPending();
|
||||
void renderFrame();
|
||||
|
||||
mpv_handle* mpv_ = nullptr;
|
||||
CGLContextObj cgl_ = nullptr;
|
||||
mpv_render_context* renderContext_ = nullptr;
|
||||
void* glDylib_ = nullptr;
|
||||
|
||||
std::string shmBaseName_;
|
||||
ShmRing ring_;
|
||||
uint32_t generation_ = 0;
|
||||
|
||||
int width_ = 0;
|
||||
int height_ = 0;
|
||||
GLuint texture_ = 0;
|
||||
GLuint fbo_ = 0;
|
||||
GLuint pbos_[FRAME_SHM_RING_SLOTS] = {0};
|
||||
int cursor_ = 0;
|
||||
int pendingPbo_ = -1;
|
||||
uint64_t nextSeq_ = 1;
|
||||
|
||||
std::mutex mutex_;
|
||||
std::condition_variable cv_;
|
||||
bool updatePending_ = false;
|
||||
bool stopRequested_ = false;
|
||||
int pendingWidth_ = 0;
|
||||
int pendingHeight_ = 0;
|
||||
std::atomic<int> initState_{0};
|
||||
};
|
||||
|
||||
inline void* frameHelperGetProcAddress(void* ctx, const char* name) {
|
||||
return dlsym(ctx, name);
|
||||
}
|
||||
|
||||
inline bool RenderPipeline::start(mpv_handle* mpv,
|
||||
const std::string& shmBaseName, int width,
|
||||
int height, std::string& errorOut) {
|
||||
mpv_ = mpv;
|
||||
shmBaseName_ = shmBaseName;
|
||||
width_ = width;
|
||||
height_ = height;
|
||||
|
||||
glDylib_ = dlopen(
|
||||
"/System/Library/Frameworks/OpenGL.framework/Versions/Current/OpenGL",
|
||||
RTLD_LAZY | RTLD_LOCAL);
|
||||
if (!glDylib_) {
|
||||
errorOut = "failed to open the OpenGL framework";
|
||||
return false;
|
||||
}
|
||||
|
||||
CGLPixelFormatAttribute attrs[] = {
|
||||
kCGLPFAOpenGLProfile, (CGLPixelFormatAttribute)kCGLOGLPVersion_3_2_Core,
|
||||
kCGLPFAAccelerated,
|
||||
kCGLPFAColorSize, (CGLPixelFormatAttribute)24,
|
||||
kCGLPFAAlphaSize, (CGLPixelFormatAttribute)8,
|
||||
(CGLPixelFormatAttribute)0,
|
||||
};
|
||||
CGLPixelFormatObj pixelFormat = nullptr;
|
||||
GLint matched = 0;
|
||||
if (CGLChoosePixelFormat(attrs, &pixelFormat, &matched) != kCGLNoError ||
|
||||
!pixelFormat) {
|
||||
errorOut = "no accelerated CGL pixel format";
|
||||
return false;
|
||||
}
|
||||
const CGLError contextError =
|
||||
CGLCreateContext(pixelFormat, nullptr, &cgl_);
|
||||
CGLDestroyPixelFormat(pixelFormat);
|
||||
if (contextError != kCGLNoError || !cgl_) {
|
||||
errorOut = "failed to create a headless CGL context";
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
inline bool RenderPipeline::setupGl(std::string& errorOut) {
|
||||
CGLSetCurrentContext(cgl_);
|
||||
|
||||
if (!rebuildTargets(width_, height_)) {
|
||||
errorOut = "framebuffer setup failed";
|
||||
return false;
|
||||
}
|
||||
|
||||
mpv_opengl_init_params glInit = {frameHelperGetProcAddress, glDylib_};
|
||||
mpv_render_param params[] = {
|
||||
{MPV_RENDER_PARAM_API_TYPE,
|
||||
const_cast<char*>(MPV_RENDER_API_TYPE_OPENGL)},
|
||||
{MPV_RENDER_PARAM_OPENGL_INIT_PARAMS, &glInit},
|
||||
{MPV_RENDER_PARAM_INVALID, nullptr},
|
||||
};
|
||||
const int result = mpv_render_context_create(&renderContext_, mpv_, params);
|
||||
if (result < 0) {
|
||||
errorOut = mpv_error_string(result);
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
inline bool RenderPipeline::rebuildTargets(int width, int height) {
|
||||
if (texture_) glDeleteTextures(1, &texture_);
|
||||
if (fbo_) glDeleteFramebuffers(1, &fbo_);
|
||||
if (pbos_[0]) glDeleteBuffers(FRAME_SHM_RING_SLOTS, pbos_);
|
||||
pendingPbo_ = -1;
|
||||
cursor_ = 0;
|
||||
|
||||
width_ = width;
|
||||
height_ = height;
|
||||
const GLsizeiptr frameBytes = (GLsizeiptr)width * 4 * height;
|
||||
|
||||
glGenTextures(1, &texture_);
|
||||
glBindTexture(GL_TEXTURE_2D, texture_);
|
||||
glTexImage2D(GL_TEXTURE_2D, 0, GL_RGBA8, width, height, 0, GL_RGBA,
|
||||
GL_UNSIGNED_BYTE, nullptr);
|
||||
glTexParameteri(GL_TEXTURE_2D, GL_TEXTURE_MIN_FILTER, GL_LINEAR);
|
||||
glGenFramebuffers(1, &fbo_);
|
||||
glBindFramebuffer(GL_FRAMEBUFFER, fbo_);
|
||||
glFramebufferTexture2D(GL_FRAMEBUFFER, GL_COLOR_ATTACHMENT0, GL_TEXTURE_2D,
|
||||
texture_, 0);
|
||||
glGenBuffers(FRAME_SHM_RING_SLOTS, pbos_);
|
||||
for (GLuint pbo : pbos_) {
|
||||
glBindBuffer(GL_PIXEL_PACK_BUFFER, pbo);
|
||||
glBufferData(GL_PIXEL_PACK_BUFFER, frameBytes, nullptr, GL_STREAM_READ);
|
||||
}
|
||||
glBindBuffer(GL_PIXEL_PACK_BUFFER, 0);
|
||||
|
||||
if (glCheckFramebufferStatus(GL_FRAMEBUFFER) != GL_FRAMEBUFFER_COMPLETE) {
|
||||
return false;
|
||||
}
|
||||
|
||||
generation_ += 1;
|
||||
const std::string shmName =
|
||||
shmBaseName_ + "-g" + std::to_string(generation_);
|
||||
if (!ring_.create(shmName, width, height, generation_)) {
|
||||
return false;
|
||||
}
|
||||
if (onGenerationChanged) {
|
||||
onGenerationChanged(shmName, width, height, generation_);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
inline void RenderPipeline::notifyUpdate() {
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(mutex_);
|
||||
updatePending_ = true;
|
||||
}
|
||||
cv_.notify_all();
|
||||
}
|
||||
|
||||
inline void RenderPipeline::requestResize(int width, int height) {
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(mutex_);
|
||||
pendingWidth_ = width;
|
||||
pendingHeight_ = height;
|
||||
updatePending_ = true;
|
||||
}
|
||||
cv_.notify_all();
|
||||
}
|
||||
|
||||
inline void RenderPipeline::stop() {
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(mutex_);
|
||||
stopRequested_ = true;
|
||||
}
|
||||
cv_.notify_all();
|
||||
}
|
||||
|
||||
inline void RenderPipeline::publishPending() {
|
||||
if (pendingPbo_ < 0 || !ring_.header) return;
|
||||
glBindBuffer(GL_PIXEL_PACK_BUFFER, pbos_[pendingPbo_]);
|
||||
const void* mapped = glMapBufferRange(GL_PIXEL_PACK_BUFFER, 0,
|
||||
(GLsizeiptr)ring_.header->frame_bytes,
|
||||
GL_MAP_READ_BIT);
|
||||
if (mapped) {
|
||||
const uint64_t seq = nextSeq_++;
|
||||
FrameShmSlot& slot =
|
||||
ring_.header->slots[seq % FRAME_SHM_RING_SLOTS];
|
||||
slot.seq.store(0, std::memory_order_release);
|
||||
std::memcpy(ring_.slotData(seq), mapped,
|
||||
(size_t)ring_.header->frame_bytes);
|
||||
slot.produce_time_ns = nowNs();
|
||||
slot.seq.store(seq, std::memory_order_release);
|
||||
ring_.header->latest_seq.store(seq, std::memory_order_release);
|
||||
glUnmapBuffer(GL_PIXEL_PACK_BUFFER);
|
||||
}
|
||||
glBindBuffer(GL_PIXEL_PACK_BUFFER, 0);
|
||||
pendingPbo_ = -1;
|
||||
}
|
||||
|
||||
inline void RenderPipeline::renderFrame() {
|
||||
mpv_opengl_fbo fbo = {(int)fbo_, width_, height_, 0};
|
||||
int flipY = 1;
|
||||
mpv_render_param params[] = {
|
||||
{MPV_RENDER_PARAM_OPENGL_FBO, &fbo},
|
||||
{MPV_RENDER_PARAM_FLIP_Y, &flipY},
|
||||
{MPV_RENDER_PARAM_INVALID, nullptr},
|
||||
};
|
||||
mpv_render_context_render(renderContext_, params);
|
||||
|
||||
glBindFramebuffer(GL_READ_FRAMEBUFFER, fbo_);
|
||||
glReadBuffer(GL_COLOR_ATTACHMENT0);
|
||||
glPixelStorei(GL_PACK_ALIGNMENT, 1);
|
||||
glPixelStorei(GL_PACK_ROW_LENGTH, 0);
|
||||
publishPending();
|
||||
glBindBuffer(GL_PIXEL_PACK_BUFFER, pbos_[cursor_]);
|
||||
glReadPixels(0, 0, width_, height_, GL_BGRA, GL_UNSIGNED_INT_8_8_8_8_REV,
|
||||
nullptr);
|
||||
glBindBuffer(GL_PIXEL_PACK_BUFFER, 0);
|
||||
pendingPbo_ = cursor_;
|
||||
cursor_ = (cursor_ + 1) % FRAME_SHM_RING_SLOTS;
|
||||
mpv_render_context_report_swap(renderContext_);
|
||||
}
|
||||
|
||||
inline void RenderPipeline::runLoop() {
|
||||
std::string glError;
|
||||
if (!setupGl(glError)) {
|
||||
emitLine(JsonWriter()
|
||||
.str("event", "fatal")
|
||||
.str("error", "render init failed: " + glError)
|
||||
.finish());
|
||||
initState_.store(-1);
|
||||
return;
|
||||
}
|
||||
initState_.store(1);
|
||||
|
||||
while (true) {
|
||||
int resizeWidth = 0;
|
||||
int resizeHeight = 0;
|
||||
{
|
||||
std::unique_lock<std::mutex> lock(mutex_);
|
||||
cv_.wait_for(lock, std::chrono::milliseconds(100), [&] {
|
||||
return updatePending_ || stopRequested_;
|
||||
});
|
||||
if (stopRequested_) break;
|
||||
updatePending_ = false;
|
||||
resizeWidth = pendingWidth_;
|
||||
resizeHeight = pendingHeight_;
|
||||
pendingWidth_ = 0;
|
||||
pendingHeight_ = 0;
|
||||
}
|
||||
|
||||
if (resizeWidth > 0 && resizeHeight > 0 &&
|
||||
(resizeWidth != width_ || resizeHeight != height_)) {
|
||||
if (!rebuildTargets(resizeWidth, resizeHeight)) {
|
||||
emitLine(JsonWriter()
|
||||
.str("event", "fatal")
|
||||
.str("error", "resize target rebuild failed")
|
||||
.finish());
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
const uint64_t flags = mpv_render_context_update(renderContext_);
|
||||
if (flags & MPV_RENDER_UPDATE_FRAME) {
|
||||
renderFrame();
|
||||
} else {
|
||||
/* Flush a stranded readback so pause keeps the last frame. */
|
||||
publishPending();
|
||||
}
|
||||
if (ring_.header) {
|
||||
ring_.header->heartbeat_ns.store(nowNs(),
|
||||
std::memory_order_relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
mpv_render_context_set_update_callback(renderContext_, nullptr, nullptr);
|
||||
mpv_render_context_free(renderContext_);
|
||||
renderContext_ = nullptr;
|
||||
if (texture_) glDeleteTextures(1, &texture_);
|
||||
if (fbo_) glDeleteFramebuffers(1, &fbo_);
|
||||
if (pbos_[0]) glDeleteBuffers(FRAME_SHM_RING_SLOTS, pbos_);
|
||||
CGLSetCurrentContext(nullptr);
|
||||
if (cgl_) CGLDestroyContext(cgl_);
|
||||
ring_.destroy();
|
||||
}
|
||||
|
||||
} // namespace frame_helper
|
||||
@@ -0,0 +1,45 @@
|
||||
/*
|
||||
* Shared-memory frame ring for the frame-copy embedded MPV helper.
|
||||
*
|
||||
* Same triple-buffer seqlock protocol as spikes/mpv-frame-copy (validated
|
||||
* there on M1 Pro up to 4K60), plus a `generation` field: every viewport
|
||||
* resize creates a fresh shm segment named `<base>-g<generation>` and the
|
||||
* reader re-attaches when the helper announces the new generation.
|
||||
*
|
||||
* Must compile as C11 (reader addon) and C++17 (helper).
|
||||
*/
|
||||
#pragma once
|
||||
|
||||
#include <stdint.h>
|
||||
|
||||
#ifdef __cplusplus
|
||||
#include <atomic>
|
||||
typedef std::atomic<uint64_t> frame_shm_atomic_u64;
|
||||
#else
|
||||
#include <stdatomic.h>
|
||||
typedef _Atomic uint64_t frame_shm_atomic_u64;
|
||||
#endif
|
||||
|
||||
#define FRAME_SHM_MAGIC 0x564d5046u /* 'FPMV' */
|
||||
#define FRAME_SHM_VERSION 1u
|
||||
#define FRAME_SHM_RING_SLOTS 3u
|
||||
#define FRAME_SHM_DATA_ALIGN 4096u
|
||||
|
||||
typedef struct {
|
||||
frame_shm_atomic_u64 seq; /* 0 while the slot is being (re)written */
|
||||
uint64_t produce_time_ns; /* CLOCK_MONOTONIC_RAW after copy completes */
|
||||
} FrameShmSlot;
|
||||
|
||||
typedef struct {
|
||||
uint32_t magic; /* written last during init; readers gate on it */
|
||||
uint32_t version;
|
||||
uint32_t width;
|
||||
uint32_t height;
|
||||
uint32_t stride; /* bytes per row, tightly packed BGRA */
|
||||
uint32_t generation;
|
||||
uint64_t frame_bytes;
|
||||
uint64_t data_offset;
|
||||
frame_shm_atomic_u64 latest_seq; /* newest complete frame, 0 = none */
|
||||
frame_shm_atomic_u64 heartbeat_ns;
|
||||
FrameShmSlot slots[FRAME_SHM_RING_SLOTS];
|
||||
} FrameShmHeader;
|
||||
@@ -0,0 +1,681 @@
|
||||
/*
|
||||
* iptvnator-mpv-helper — frame-copy embedded MPV helper process (macOS).
|
||||
*
|
||||
* One process = one playback session. Owns libmpv end to end: decodes,
|
||||
* renders offscreen at viewport size, publishes BGRA frames into a shared
|
||||
* memory ring (frame_shm.h), and plays audio directly through the OS.
|
||||
*
|
||||
* Control protocol: tab-separated commands on stdin, JSON events on stdout
|
||||
* (see frame_helper_io.h). The `snapshot` event mirrors the TypeScript
|
||||
* NativeEmbeddedMpvSessionSnapshot shape so the Electron-side adapter can
|
||||
* cache it verbatim. Status semantics are ported from embedded_mpv.mm:
|
||||
* END_FILE reason mapping, eof-reached => ended (keep-open), pause flag
|
||||
* gated on a loaded path, only fatal/load errors flip the status.
|
||||
*
|
||||
* Usage:
|
||||
* iptvnator_mpv_helper --shm-base /impv-<id> --width 1280 --height 720
|
||||
* [--volume 0..1] [--hwdec auto]
|
||||
*/
|
||||
#include <mpv/client.h>
|
||||
|
||||
#include <algorithm>
|
||||
#include <atomic>
|
||||
#include <clocale>
|
||||
#include <cmath>
|
||||
#include <csignal>
|
||||
#include <cstdio>
|
||||
#include <cstring>
|
||||
#include <ctime>
|
||||
#include <iostream>
|
||||
#include <mutex>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
#include <vector>
|
||||
|
||||
#include "frame_helper_io.h"
|
||||
#include "frame_helper_render.h"
|
||||
|
||||
namespace {
|
||||
|
||||
using frame_helper::Command;
|
||||
using frame_helper::emitLine;
|
||||
using frame_helper::JsonWriter;
|
||||
using frame_helper::RenderPipeline;
|
||||
|
||||
struct TrackInfo {
|
||||
int64_t id = -1;
|
||||
std::string title;
|
||||
std::string language;
|
||||
bool defaultTrack = false;
|
||||
bool forced = false;
|
||||
};
|
||||
|
||||
struct SnapshotState {
|
||||
std::string status = "idle";
|
||||
double positionSeconds = 0;
|
||||
double durationSeconds = -1; /* <0 => null */
|
||||
double volume = 1; /* 0..1 */
|
||||
std::string streamUrl;
|
||||
std::vector<TrackInfo> audioTracks;
|
||||
int64_t selectedAudioTrackId = -1;
|
||||
std::vector<TrackInfo> subtitleTracks;
|
||||
int64_t selectedSubtitleTrackId = -1;
|
||||
double playbackSpeed = 1;
|
||||
std::string aspectOverride = "no";
|
||||
bool recordingActive = false;
|
||||
std::string recordingTargetPath;
|
||||
std::string recordingStartedAt;
|
||||
std::string recordingError;
|
||||
std::string error;
|
||||
};
|
||||
|
||||
struct HelperState {
|
||||
mpv_handle* mpv = nullptr;
|
||||
RenderPipeline pipeline;
|
||||
|
||||
std::mutex mutex;
|
||||
SnapshotState snapshot;
|
||||
bool paused = false;
|
||||
bool loadedPath = false;
|
||||
bool dirty = true;
|
||||
std::string lastEmittedStatus;
|
||||
uint64_t lastEmitNs = 0;
|
||||
|
||||
std::atomic<bool> running{true};
|
||||
};
|
||||
|
||||
HelperState g_state;
|
||||
|
||||
constexpr uint64_t SNAPSHOT_EMIT_INTERVAL_NS = 250ull * 1000 * 1000;
|
||||
|
||||
std::string isoTimestampNow() {
|
||||
char buffer[32];
|
||||
const time_t now = time(nullptr);
|
||||
struct tm utc;
|
||||
gmtime_r(&now, &utc);
|
||||
strftime(buffer, sizeof(buffer), "%Y-%m-%dT%H:%M:%SZ", &utc);
|
||||
return buffer;
|
||||
}
|
||||
|
||||
std::string tracksJson(const std::vector<TrackInfo>& tracks,
|
||||
int64_t selectedId) {
|
||||
std::string out = "[";
|
||||
bool first = true;
|
||||
for (const TrackInfo& track : tracks) {
|
||||
JsonWriter writer;
|
||||
writer.num("id", (double)track.id);
|
||||
if (!track.title.empty()) writer.str("title", track.title);
|
||||
if (!track.language.empty()) writer.str("language", track.language);
|
||||
writer.boolean("selected", track.id == selectedId);
|
||||
writer.boolean("defaultTrack", track.defaultTrack);
|
||||
writer.boolean("forced", track.forced);
|
||||
if (!first) out += ',';
|
||||
first = false;
|
||||
out += writer.finish();
|
||||
}
|
||||
out += ']';
|
||||
return out;
|
||||
}
|
||||
|
||||
/* Compose the snapshot event. Caller holds g_state.mutex. */
|
||||
std::string composeSnapshotLocked() {
|
||||
const SnapshotState& s = g_state.snapshot;
|
||||
JsonWriter writer;
|
||||
writer.str("event", "snapshot");
|
||||
writer.str("status", s.status);
|
||||
writer.num("positionSeconds", std::max(0.0, s.positionSeconds));
|
||||
if (s.durationSeconds >= 0) {
|
||||
writer.num("durationSeconds", s.durationSeconds);
|
||||
} else {
|
||||
writer.nullValue("durationSeconds");
|
||||
}
|
||||
writer.num("volume", std::clamp(s.volume, 0.0, 1.5));
|
||||
writer.str("streamUrl", s.streamUrl);
|
||||
writer.raw("audioTracks", tracksJson(s.audioTracks, s.selectedAudioTrackId));
|
||||
if (s.selectedAudioTrackId >= 0) {
|
||||
writer.num("selectedAudioTrackId", (double)s.selectedAudioTrackId);
|
||||
} else {
|
||||
writer.nullValue("selectedAudioTrackId");
|
||||
}
|
||||
writer.raw("subtitleTracks",
|
||||
tracksJson(s.subtitleTracks, s.selectedSubtitleTrackId));
|
||||
if (s.selectedSubtitleTrackId >= 0) {
|
||||
writer.num("selectedSubtitleTrackId",
|
||||
(double)s.selectedSubtitleTrackId);
|
||||
} else {
|
||||
writer.nullValue("selectedSubtitleTrackId");
|
||||
}
|
||||
writer.num("playbackSpeed", s.playbackSpeed);
|
||||
writer.str("aspectOverride", s.aspectOverride);
|
||||
JsonWriter recording;
|
||||
recording.boolean("active", s.recordingActive);
|
||||
if (!s.recordingTargetPath.empty())
|
||||
recording.str("targetPath", s.recordingTargetPath);
|
||||
if (!s.recordingStartedAt.empty())
|
||||
recording.str("startedAt", s.recordingStartedAt);
|
||||
if (!s.recordingError.empty()) recording.str("error", s.recordingError);
|
||||
writer.raw("recording", recording.finish());
|
||||
if (!s.error.empty()) writer.str("error", s.error);
|
||||
return writer.finish();
|
||||
}
|
||||
|
||||
/* Emit the snapshot when dirty: immediately on status changes, otherwise at
|
||||
* most every SNAPSHOT_EMIT_INTERVAL_NS. Called from the mpv event thread. */
|
||||
void maybeEmitSnapshot(bool force = false) {
|
||||
std::string line;
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(g_state.mutex);
|
||||
if (!g_state.dirty && !force) return;
|
||||
const uint64_t now = frame_helper::nowNs();
|
||||
const bool statusChanged =
|
||||
g_state.snapshot.status != g_state.lastEmittedStatus;
|
||||
if (!force && !statusChanged &&
|
||||
now - g_state.lastEmitNs < SNAPSHOT_EMIT_INTERVAL_NS) {
|
||||
return;
|
||||
}
|
||||
g_state.dirty = false;
|
||||
g_state.lastEmitNs = now;
|
||||
g_state.lastEmittedStatus = g_state.snapshot.status;
|
||||
line = composeSnapshotLocked();
|
||||
}
|
||||
emitLine(line);
|
||||
}
|
||||
|
||||
void markDirtyLocked() { g_state.dirty = true; }
|
||||
|
||||
/* ---- track-list parsing (ported from embedded_mpv.mm) ------------------- */
|
||||
|
||||
void parseTrackNode(const mpv_node& trackNode, const char* wantedType,
|
||||
std::vector<TrackInfo>& out) {
|
||||
if (trackNode.format != MPV_FORMAT_NODE_MAP || !trackNode.u.list) return;
|
||||
TrackInfo track;
|
||||
bool matchesType = false;
|
||||
const mpv_node_list& map = *trackNode.u.list;
|
||||
for (int i = 0; i < map.num; i++) {
|
||||
const char* key = map.keys[i];
|
||||
const mpv_node& value = map.values[i];
|
||||
if (!key) continue;
|
||||
if (std::strcmp(key, "type") == 0 &&
|
||||
value.format == MPV_FORMAT_STRING && value.u.string) {
|
||||
matchesType = std::strcmp(value.u.string, wantedType) == 0;
|
||||
} else if (std::strcmp(key, "id") == 0 &&
|
||||
value.format == MPV_FORMAT_INT64) {
|
||||
track.id = value.u.int64;
|
||||
} else if (std::strcmp(key, "title") == 0 &&
|
||||
value.format == MPV_FORMAT_STRING && value.u.string) {
|
||||
track.title = value.u.string;
|
||||
} else if (std::strcmp(key, "lang") == 0 &&
|
||||
value.format == MPV_FORMAT_STRING && value.u.string) {
|
||||
track.language = value.u.string;
|
||||
} else if (std::strcmp(key, "default") == 0 &&
|
||||
value.format == MPV_FORMAT_FLAG) {
|
||||
track.defaultTrack = value.u.flag != 0;
|
||||
} else if (std::strcmp(key, "forced") == 0 &&
|
||||
value.format == MPV_FORMAT_FLAG) {
|
||||
track.forced = value.u.flag != 0;
|
||||
}
|
||||
}
|
||||
if (matchesType && track.id >= 0) {
|
||||
out.push_back(std::move(track));
|
||||
}
|
||||
}
|
||||
|
||||
void updateTracksFromNode(const mpv_node& trackListNode) {
|
||||
if (trackListNode.format != MPV_FORMAT_NODE_ARRAY ||
|
||||
!trackListNode.u.list) {
|
||||
return;
|
||||
}
|
||||
std::vector<TrackInfo> audio;
|
||||
std::vector<TrackInfo> subs;
|
||||
for (int i = 0; i < trackListNode.u.list->num; i++) {
|
||||
parseTrackNode(trackListNode.u.list->values[i], "audio", audio);
|
||||
parseTrackNode(trackListNode.u.list->values[i], "sub", subs);
|
||||
}
|
||||
g_state.snapshot.audioTracks = std::move(audio);
|
||||
g_state.snapshot.subtitleTracks = std::move(subs);
|
||||
}
|
||||
|
||||
bool parseTrackSelection(const char* value, int64_t& out) {
|
||||
if (!value) return false;
|
||||
char* end = nullptr;
|
||||
const long long parsed = strtoll(value, &end, 10);
|
||||
if (end == value || (end && *end != '\0')) return false;
|
||||
out = parsed;
|
||||
return true;
|
||||
}
|
||||
|
||||
/* ---- mpv event loop (status semantics ported from embedded_mpv.mm) ------ */
|
||||
|
||||
void handlePropertyChange(const mpv_event_property& property) {
|
||||
const std::string name = property.name ? property.name : "";
|
||||
SnapshotState& s = g_state.snapshot;
|
||||
|
||||
if (name == "time-pos" && property.format == MPV_FORMAT_DOUBLE &&
|
||||
property.data) {
|
||||
s.positionSeconds = *static_cast<double*>(property.data);
|
||||
} else if (name == "duration" && property.format == MPV_FORMAT_DOUBLE &&
|
||||
property.data) {
|
||||
s.durationSeconds = *static_cast<double*>(property.data);
|
||||
} else if (name == "pause" && property.format == MPV_FORMAT_FLAG &&
|
||||
property.data) {
|
||||
g_state.paused = *static_cast<int*>(property.data) != 0;
|
||||
if (s.status != "loading" && s.status != "ended" &&
|
||||
s.status != "error" && g_state.loadedPath) {
|
||||
s.status = g_state.paused ? "paused" : "playing";
|
||||
}
|
||||
} else if (name == "eof-reached" && property.format == MPV_FORMAT_FLAG &&
|
||||
property.data) {
|
||||
if (*static_cast<int*>(property.data) != 0 && g_state.loadedPath) {
|
||||
s.status = "ended";
|
||||
}
|
||||
} else if (name == "volume" && property.format == MPV_FORMAT_DOUBLE &&
|
||||
property.data) {
|
||||
s.volume = *static_cast<double*>(property.data) / 100.0;
|
||||
} else if (name == "path" && property.format == MPV_FORMAT_STRING &&
|
||||
property.data) {
|
||||
const char* value = *static_cast<char**>(property.data);
|
||||
s.streamUrl = value ? value : "";
|
||||
} else if (name == "track-list" && property.format == MPV_FORMAT_NODE &&
|
||||
property.data) {
|
||||
updateTracksFromNode(*static_cast<mpv_node*>(property.data));
|
||||
} else if (name == "aid" && property.format == MPV_FORMAT_STRING &&
|
||||
property.data) {
|
||||
int64_t selected = -1;
|
||||
if (parseTrackSelection(*static_cast<char**>(property.data),
|
||||
selected)) {
|
||||
s.selectedAudioTrackId = selected;
|
||||
} else {
|
||||
s.selectedAudioTrackId = -1;
|
||||
}
|
||||
} else if (name == "sid" && property.format == MPV_FORMAT_STRING &&
|
||||
property.data) {
|
||||
int64_t selected = -1;
|
||||
if (parseTrackSelection(*static_cast<char**>(property.data),
|
||||
selected)) {
|
||||
s.selectedSubtitleTrackId = selected;
|
||||
} else {
|
||||
s.selectedSubtitleTrackId = -1;
|
||||
}
|
||||
} else if (name == "speed" && property.format == MPV_FORMAT_DOUBLE &&
|
||||
property.data) {
|
||||
s.playbackSpeed = *static_cast<double*>(property.data);
|
||||
} else if (name == "video-aspect-override" &&
|
||||
property.format == MPV_FORMAT_STRING && property.data) {
|
||||
const char* value = *static_cast<char**>(property.data);
|
||||
/* mpv reports the unset override as "-1.000000"; the renderer
|
||||
* contract uses "no" for that state. */
|
||||
s.aspectOverride =
|
||||
value && *value && std::atof(value) > 0 ? value : "no";
|
||||
} else {
|
||||
return;
|
||||
}
|
||||
markDirtyLocked();
|
||||
}
|
||||
|
||||
void runMpvEventLoop() {
|
||||
while (g_state.running.load()) {
|
||||
mpv_event* event = mpv_wait_event(g_state.mpv, 0.1);
|
||||
if (!event) continue;
|
||||
if (event->event_id == MPV_EVENT_NONE) {
|
||||
maybeEmitSnapshot();
|
||||
continue;
|
||||
}
|
||||
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(g_state.mutex);
|
||||
SnapshotState& s = g_state.snapshot;
|
||||
switch (event->event_id) {
|
||||
case MPV_EVENT_START_FILE:
|
||||
s.status = "loading";
|
||||
s.error.clear();
|
||||
s.audioTracks.clear();
|
||||
s.selectedAudioTrackId = -1;
|
||||
s.subtitleTracks.clear();
|
||||
s.selectedSubtitleTrackId = -1;
|
||||
g_state.loadedPath = false;
|
||||
markDirtyLocked();
|
||||
break;
|
||||
case MPV_EVENT_FILE_LOADED:
|
||||
g_state.loadedPath = true;
|
||||
s.status = g_state.paused ? "paused" : "playing";
|
||||
markDirtyLocked();
|
||||
break;
|
||||
case MPV_EVENT_END_FILE: {
|
||||
const auto* endFile =
|
||||
static_cast<mpv_event_end_file*>(event->data);
|
||||
if (endFile &&
|
||||
endFile->reason == MPV_END_FILE_REASON_ERROR) {
|
||||
s.status = "error";
|
||||
s.error = endFile->error < 0
|
||||
? mpv_error_string(endFile->error)
|
||||
: "Playback failed.";
|
||||
} else if (endFile &&
|
||||
endFile->reason == MPV_END_FILE_REASON_EOF &&
|
||||
g_state.running.load()) {
|
||||
s.status = "ended";
|
||||
} else if (endFile &&
|
||||
endFile->reason ==
|
||||
MPV_END_FILE_REASON_REDIRECT &&
|
||||
g_state.running.load()) {
|
||||
s.status = "loading";
|
||||
s.error.clear();
|
||||
g_state.loadedPath = false;
|
||||
} else if (g_state.running.load()) {
|
||||
s.status = "idle";
|
||||
}
|
||||
markDirtyLocked();
|
||||
break;
|
||||
}
|
||||
case MPV_EVENT_PROPERTY_CHANGE: {
|
||||
const auto* property =
|
||||
static_cast<mpv_event_property*>(event->data);
|
||||
if (property) handlePropertyChange(*property);
|
||||
break;
|
||||
}
|
||||
case MPV_EVENT_LOG_MESSAGE: {
|
||||
const auto* logMessage =
|
||||
static_cast<mpv_event_log_message*>(event->data);
|
||||
if (!logMessage || !logMessage->level ||
|
||||
!logMessage->text) {
|
||||
break;
|
||||
}
|
||||
const std::string level = logMessage->level;
|
||||
if (level == "error" || level == "fatal") {
|
||||
s.error = logMessage->text;
|
||||
markDirtyLocked();
|
||||
}
|
||||
if (level == "fatal") {
|
||||
s.status = "error";
|
||||
}
|
||||
emitLine(JsonWriter()
|
||||
.str("event", "log")
|
||||
.str("level", level)
|
||||
.str("prefix", logMessage->prefix
|
||||
? logMessage->prefix
|
||||
: "mpv")
|
||||
.str("text", logMessage->text)
|
||||
.finish());
|
||||
break;
|
||||
}
|
||||
case MPV_EVENT_SHUTDOWN:
|
||||
g_state.running.store(false);
|
||||
s.status = "closed";
|
||||
markDirtyLocked();
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
maybeEmitSnapshot();
|
||||
}
|
||||
maybeEmitSnapshot(true);
|
||||
}
|
||||
|
||||
/* ---- command handling ---------------------------------------------------- */
|
||||
|
||||
void setPropertyString(const char* name, const std::string& value) {
|
||||
mpv_set_property_string(g_state.mpv, name, value.c_str());
|
||||
}
|
||||
|
||||
void handleLoadCommand(const Command& command) {
|
||||
const std::string url = command.get("url");
|
||||
if (url.empty()) return;
|
||||
|
||||
/* Loading a replacement stream stops any active recording first,
|
||||
* mirroring the .mm behavior. */
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(g_state.mutex);
|
||||
if (g_state.snapshot.recordingActive) {
|
||||
setPropertyString("stream-record", "");
|
||||
g_state.snapshot.recordingActive = false;
|
||||
g_state.snapshot.recordingStartedAt.clear();
|
||||
}
|
||||
g_state.snapshot.streamUrl = url;
|
||||
g_state.snapshot.error.clear();
|
||||
g_state.snapshot.status = "loading";
|
||||
g_state.dirty = true;
|
||||
}
|
||||
|
||||
/* Generic pass-through: every `opt.<name>` argument becomes a loadfile
|
||||
* option, so the adapter controls start/user-agent/referrer/headers
|
||||
* without protocol changes. */
|
||||
std::string optionString;
|
||||
for (const auto& [key, value] : command.args) {
|
||||
if (key.rfind("opt.", 0) != 0 || value.empty()) continue;
|
||||
if (!optionString.empty()) optionString += ',';
|
||||
std::string escaped = value;
|
||||
/* loadfile options string uses %len%value quoting to stay safe */
|
||||
optionString += key.substr(4) + "=%" +
|
||||
std::to_string(escaped.size()) + "%" + escaped;
|
||||
}
|
||||
|
||||
const char* args[] = {"loadfile", url.c_str(), "replace",
|
||||
optionString.empty() ? nullptr : optionString.c_str(),
|
||||
nullptr};
|
||||
/* mpv >= 0.38 expects an index argument before options; use the
|
||||
* options-map-free form: loadfile <url> replace [options] */
|
||||
int result;
|
||||
if (optionString.empty()) {
|
||||
const char* plain[] = {"loadfile", url.c_str(), "replace", nullptr};
|
||||
result = mpv_command(g_state.mpv, plain);
|
||||
} else {
|
||||
const char* withOpts[] = {"loadfile", url.c_str(), "replace", "-1",
|
||||
optionString.c_str(), nullptr};
|
||||
result = mpv_command(g_state.mpv, withOpts);
|
||||
if (result == MPV_ERROR_INVALID_PARAMETER) {
|
||||
/* Older mpv without the index argument. */
|
||||
result = mpv_command(g_state.mpv, args);
|
||||
}
|
||||
}
|
||||
if (result < 0) {
|
||||
std::lock_guard<std::mutex> lock(g_state.mutex);
|
||||
g_state.snapshot.status = "error";
|
||||
g_state.snapshot.error = mpv_error_string(result);
|
||||
g_state.dirty = true;
|
||||
}
|
||||
}
|
||||
|
||||
void handleRecordCommand(const Command& command) {
|
||||
const std::string path = command.get("path");
|
||||
const int result = mpv_set_property_string(g_state.mpv, "stream-record",
|
||||
path.c_str());
|
||||
std::lock_guard<std::mutex> lock(g_state.mutex);
|
||||
if (result < 0) {
|
||||
g_state.snapshot.recordingError = mpv_error_string(result);
|
||||
} else if (path.empty()) {
|
||||
g_state.snapshot.recordingActive = false;
|
||||
g_state.snapshot.recordingStartedAt.clear();
|
||||
g_state.snapshot.recordingError.clear();
|
||||
} else {
|
||||
g_state.snapshot.recordingActive = true;
|
||||
g_state.snapshot.recordingTargetPath = path;
|
||||
g_state.snapshot.recordingStartedAt = isoTimestampNow();
|
||||
g_state.snapshot.recordingError.clear();
|
||||
}
|
||||
g_state.dirty = true;
|
||||
}
|
||||
|
||||
void handleCommand(const Command& command) {
|
||||
if (command.name == "load") {
|
||||
handleLoadCommand(command);
|
||||
} else if (command.name == "pause") {
|
||||
const int flag = command.get("value") == "1" ? 1 : 0;
|
||||
mpv_set_property(g_state.mpv, "pause", MPV_FORMAT_FLAG,
|
||||
const_cast<int*>(&flag));
|
||||
} else if (command.name == "seek") {
|
||||
const double seconds = command.getDouble("seconds", 0);
|
||||
const std::string value = std::to_string(seconds);
|
||||
const char* args[] = {"seek", value.c_str(), "absolute", nullptr};
|
||||
mpv_command(g_state.mpv, args);
|
||||
} else if (command.name == "volume") {
|
||||
double percent =
|
||||
std::clamp(command.getDouble("value", 1) * 100.0, 0.0, 100.0);
|
||||
mpv_set_property(g_state.mpv, "volume", MPV_FORMAT_DOUBLE, &percent);
|
||||
} else if (command.name == "aid" || command.name == "sid") {
|
||||
const std::string value = command.get("value");
|
||||
setPropertyString(command.name.c_str(),
|
||||
value == "-1" ? "no" : value);
|
||||
} else if (command.name == "speed") {
|
||||
double speed = std::clamp(command.getDouble("value", 1), 0.25, 4.0);
|
||||
mpv_set_property(g_state.mpv, "speed", MPV_FORMAT_DOUBLE, &speed);
|
||||
} else if (command.name == "aspect") {
|
||||
setPropertyString("video-aspect-override", command.get("value", "no"));
|
||||
} else if (command.name == "record") {
|
||||
handleRecordCommand(command);
|
||||
} else if (command.name == "size") {
|
||||
const int width = (int)command.getDouble("width", 0);
|
||||
const int height = (int)command.getDouble("height", 0);
|
||||
if (width >= 16 && height >= 16) {
|
||||
g_state.pipeline.requestResize(width, height);
|
||||
}
|
||||
} else if (command.name == "quit") {
|
||||
g_state.running.store(false);
|
||||
mpv_wakeup(g_state.mpv);
|
||||
}
|
||||
}
|
||||
|
||||
void runStdinLoop() {
|
||||
std::string line;
|
||||
while (g_state.running.load() && std::getline(std::cin, line)) {
|
||||
if (line.empty()) continue;
|
||||
handleCommand(frame_helper::parseCommandLine(line));
|
||||
}
|
||||
/* stdin EOF => the parent process died or closed us: shut down. */
|
||||
g_state.running.store(false);
|
||||
mpv_wakeup(g_state.mpv);
|
||||
}
|
||||
|
||||
void onRenderUpdate(void* pipeline) {
|
||||
static_cast<RenderPipeline*>(pipeline)->notifyUpdate();
|
||||
}
|
||||
|
||||
struct HelperArgs {
|
||||
std::string shmBase = "/impv";
|
||||
std::string hwdec = "auto";
|
||||
int width = 1280;
|
||||
int height = 720;
|
||||
double volume = 1;
|
||||
};
|
||||
|
||||
HelperArgs parseArgs(int argc, char** argv) {
|
||||
HelperArgs args;
|
||||
for (int i = 1; i < argc; i++) {
|
||||
const std::string arg = argv[i];
|
||||
auto next = [&]() -> std::string {
|
||||
return i + 1 < argc ? argv[++i] : "";
|
||||
};
|
||||
if (arg == "--shm-base") args.shmBase = next();
|
||||
else if (arg == "--width") args.width = std::atoi(next().c_str());
|
||||
else if (arg == "--height") args.height = std::atoi(next().c_str());
|
||||
else if (arg == "--volume") args.volume = std::atof(next().c_str());
|
||||
else if (arg == "--hwdec") args.hwdec = next();
|
||||
}
|
||||
args.width = std::max(16, args.width);
|
||||
args.height = std::max(16, args.height);
|
||||
return args;
|
||||
}
|
||||
|
||||
} // namespace
|
||||
|
||||
int main(int argc, char** argv) {
|
||||
std::setlocale(LC_NUMERIC, "C");
|
||||
signal(SIGPIPE, SIG_IGN);
|
||||
const HelperArgs args = parseArgs(argc, argv);
|
||||
|
||||
g_state.mpv = mpv_create();
|
||||
if (!g_state.mpv) {
|
||||
emitLine(JsonWriter()
|
||||
.str("event", "fatal")
|
||||
.str("error", "mpv_create failed")
|
||||
.finish());
|
||||
return 1;
|
||||
}
|
||||
|
||||
mpv_set_option_string(g_state.mpv, "vo", "libmpv");
|
||||
mpv_set_option_string(g_state.mpv, "hwdec", args.hwdec.c_str());
|
||||
mpv_set_option_string(g_state.mpv, "keep-open", "yes");
|
||||
mpv_set_option_string(g_state.mpv, "idle", "yes");
|
||||
mpv_set_option_string(g_state.mpv, "input-default-bindings", "no");
|
||||
mpv_set_option_string(g_state.mpv, "osc", "no");
|
||||
if (mpv_initialize(g_state.mpv) < 0) {
|
||||
emitLine(JsonWriter()
|
||||
.str("event", "fatal")
|
||||
.str("error", "mpv_initialize failed")
|
||||
.finish());
|
||||
return 1;
|
||||
}
|
||||
mpv_request_log_messages(g_state.mpv, "warn");
|
||||
|
||||
double initialVolumePercent =
|
||||
std::clamp(args.volume, 0.0, 1.0) * 100.0;
|
||||
mpv_set_property(g_state.mpv, "volume", MPV_FORMAT_DOUBLE,
|
||||
&initialVolumePercent);
|
||||
|
||||
mpv_observe_property(g_state.mpv, 1, "time-pos", MPV_FORMAT_DOUBLE);
|
||||
mpv_observe_property(g_state.mpv, 2, "duration", MPV_FORMAT_DOUBLE);
|
||||
mpv_observe_property(g_state.mpv, 3, "pause", MPV_FORMAT_FLAG);
|
||||
mpv_observe_property(g_state.mpv, 4, "volume", MPV_FORMAT_DOUBLE);
|
||||
mpv_observe_property(g_state.mpv, 5, "path", MPV_FORMAT_STRING);
|
||||
mpv_observe_property(g_state.mpv, 6, "track-list", MPV_FORMAT_NODE);
|
||||
mpv_observe_property(g_state.mpv, 7, "aid", MPV_FORMAT_STRING);
|
||||
mpv_observe_property(g_state.mpv, 8, "sid", MPV_FORMAT_STRING);
|
||||
mpv_observe_property(g_state.mpv, 9, "speed", MPV_FORMAT_DOUBLE);
|
||||
mpv_observe_property(g_state.mpv, 10, "video-aspect-override",
|
||||
MPV_FORMAT_STRING);
|
||||
mpv_observe_property(g_state.mpv, 11, "eof-reached", MPV_FORMAT_FLAG);
|
||||
|
||||
g_state.pipeline.onGenerationChanged = [](const std::string& name,
|
||||
int width, int height,
|
||||
uint32_t generation) {
|
||||
emitLine(JsonWriter()
|
||||
.str("event", "shm")
|
||||
.str("name", name)
|
||||
.num("width", width)
|
||||
.num("height", height)
|
||||
.num("generation", generation)
|
||||
.finish());
|
||||
};
|
||||
|
||||
std::string startError;
|
||||
if (!g_state.pipeline.start(g_state.mpv, args.shmBase, args.width,
|
||||
args.height, startError)) {
|
||||
emitLine(JsonWriter()
|
||||
.str("event", "fatal")
|
||||
.str("error", startError)
|
||||
.finish());
|
||||
return 1;
|
||||
}
|
||||
|
||||
std::thread renderThread([] { g_state.pipeline.runLoop(); });
|
||||
/* The render context is created inside the render thread; hook the
|
||||
* update callback once it exists (or bail if GL init failed). */
|
||||
while (g_state.pipeline.initState() == 0) {
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(5));
|
||||
}
|
||||
if (g_state.pipeline.initState() < 0) {
|
||||
renderThread.join();
|
||||
mpv_terminate_destroy(g_state.mpv);
|
||||
return 1;
|
||||
}
|
||||
mpv_render_context_set_update_callback(g_state.pipeline.renderContext(),
|
||||
onRenderUpdate, &g_state.pipeline);
|
||||
|
||||
emitLine(JsonWriter()
|
||||
.str("event", "hello")
|
||||
.num("pid", (double)getpid())
|
||||
.num("protocolVersion", 1)
|
||||
.finish());
|
||||
|
||||
std::thread stdinThread(runStdinLoop);
|
||||
runMpvEventLoop();
|
||||
|
||||
g_state.pipeline.stop();
|
||||
renderThread.join();
|
||||
mpv_terminate_destroy(g_state.mpv);
|
||||
if (stdinThread.joinable()) {
|
||||
/* stdin thread exits on EOF; detach if the pipe is still open. */
|
||||
stdinThread.detach();
|
||||
}
|
||||
emitLine(JsonWriter().str("event", "bye").finish());
|
||||
return 0;
|
||||
}
|
||||
@@ -0,0 +1,220 @@
|
||||
/*
|
||||
* embedded_mpv_frame_reader — N-API reader for the frame-copy shm ring.
|
||||
*
|
||||
* Loaded by the Electron preload script; copies the newest complete BGRA
|
||||
* frame from the helper's shared-memory ring (native/helper/frame_shm.h)
|
||||
* into a caller-provided ArrayBuffer. A memcpy is mandatory: Electron's V8
|
||||
* memory cage forbids external ArrayBuffers over foreign memory.
|
||||
*
|
||||
* Plain C N-API (no node-addon-api) so it stays ABI-stable and trivial.
|
||||
* macOS-only for now — other platforms export an empty object; the
|
||||
* TypeScript side gates on platform before requiring it.
|
||||
*/
|
||||
#define NAPI_VERSION 8
|
||||
#include <node_api.h>
|
||||
|
||||
#ifdef __APPLE__
|
||||
|
||||
#include <fcntl.h>
|
||||
#include <stdatomic.h>
|
||||
#include <stdint.h>
|
||||
#include <string.h>
|
||||
#include <sys/mman.h>
|
||||
#include <sys/stat.h>
|
||||
#include <time.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include "../helper/frame_shm.h"
|
||||
|
||||
static FrameShmHeader* g_header = NULL;
|
||||
static uint8_t* g_base = NULL;
|
||||
static size_t g_size = 0;
|
||||
|
||||
static uint64_t now_ns(void) {
|
||||
return clock_gettime_nsec_np(CLOCK_MONOTONIC_RAW);
|
||||
}
|
||||
|
||||
static void unmap_current(void) {
|
||||
if (g_base) {
|
||||
munmap(g_base, g_size);
|
||||
}
|
||||
g_base = NULL;
|
||||
g_header = NULL;
|
||||
g_size = 0;
|
||||
}
|
||||
|
||||
static napi_value throw_error(napi_env env, const char* message) {
|
||||
napi_throw_error(env, NULL, message);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
static void set_named_double(napi_env env, napi_value target, const char* key,
|
||||
double value) {
|
||||
napi_value wrapped;
|
||||
napi_create_double(env, value, &wrapped);
|
||||
napi_set_named_property(env, target, key, wrapped);
|
||||
}
|
||||
|
||||
/* open(name) -> { width, height, stride, frameBytes, generation } */
|
||||
static napi_value Open(napi_env env, napi_callback_info info) {
|
||||
size_t argc = 1;
|
||||
napi_value argv[1];
|
||||
napi_get_cb_info(env, info, &argc, argv, NULL, NULL);
|
||||
if (argc < 1) return throw_error(env, "open(name) requires the shm name");
|
||||
|
||||
char name[256];
|
||||
size_t written = 0;
|
||||
if (napi_get_value_string_utf8(env, argv[0], name, sizeof(name),
|
||||
&written) != napi_ok) {
|
||||
return throw_error(env, "shm name must be a string");
|
||||
}
|
||||
|
||||
const int fd = shm_open(name, O_RDONLY, 0);
|
||||
if (fd < 0) return throw_error(env, "shm_open failed (helper gone?)");
|
||||
struct stat segment;
|
||||
if (fstat(fd, &segment) != 0 ||
|
||||
segment.st_size < (off_t)sizeof(FrameShmHeader)) {
|
||||
close(fd);
|
||||
return throw_error(env, "shm segment too small");
|
||||
}
|
||||
void* mapped =
|
||||
mmap(NULL, (size_t)segment.st_size, PROT_READ, MAP_SHARED, fd, 0);
|
||||
close(fd);
|
||||
if (mapped == MAP_FAILED) return throw_error(env, "mmap failed");
|
||||
|
||||
FrameShmHeader* header = (FrameShmHeader*)mapped;
|
||||
if (header->magic != FRAME_SHM_MAGIC ||
|
||||
header->version != FRAME_SHM_VERSION) {
|
||||
munmap(mapped, (size_t)segment.st_size);
|
||||
return throw_error(env, "frame ring not initialized yet");
|
||||
}
|
||||
|
||||
unmap_current();
|
||||
g_base = (uint8_t*)mapped;
|
||||
g_size = (size_t)segment.st_size;
|
||||
g_header = header;
|
||||
|
||||
napi_value result;
|
||||
napi_create_object(env, &result);
|
||||
set_named_double(env, result, "width", header->width);
|
||||
set_named_double(env, result, "height", header->height);
|
||||
set_named_double(env, result, "stride", header->stride);
|
||||
set_named_double(env, result, "frameBytes", (double)header->frame_bytes);
|
||||
set_named_double(env, result, "generation", header->generation);
|
||||
return result;
|
||||
}
|
||||
|
||||
/* latestSeq() -> number (0 = no frame yet / not attached) */
|
||||
static napi_value LatestSeq(napi_env env, napi_callback_info info) {
|
||||
napi_value result;
|
||||
const uint64_t seq =
|
||||
g_header
|
||||
? atomic_load_explicit(&g_header->latest_seq, memory_order_acquire)
|
||||
: 0;
|
||||
napi_create_double(env, (double)seq, &result);
|
||||
return result;
|
||||
}
|
||||
|
||||
/* copyLatest(arrayBuffer) -> { seq, ageMs, torn } | null */
|
||||
static napi_value CopyLatest(napi_env env, napi_callback_info info) {
|
||||
size_t argc = 1;
|
||||
napi_value argv[1];
|
||||
napi_get_cb_info(env, info, &argc, argv, NULL, NULL);
|
||||
if (!g_header) return throw_error(env, "call open() first");
|
||||
if (argc < 1) return throw_error(env, "copyLatest(buffer) needs a buffer");
|
||||
|
||||
void* destination = NULL;
|
||||
size_t destination_size = 0;
|
||||
if (napi_get_arraybuffer_info(env, argv[0], &destination,
|
||||
&destination_size) != napi_ok) {
|
||||
return throw_error(env, "argument must be an ArrayBuffer");
|
||||
}
|
||||
if (destination_size < g_header->frame_bytes) {
|
||||
return throw_error(env, "buffer smaller than one frame");
|
||||
}
|
||||
|
||||
napi_value null_value;
|
||||
napi_get_null(env, &null_value);
|
||||
|
||||
const uint64_t seq =
|
||||
atomic_load_explicit(&g_header->latest_seq, memory_order_acquire);
|
||||
if (seq == 0) return null_value;
|
||||
|
||||
FrameShmSlot* slot = &g_header->slots[seq % FRAME_SHM_RING_SLOTS];
|
||||
if (atomic_load_explicit(&slot->seq, memory_order_acquire) != seq) {
|
||||
return null_value; /* writer is racing this slot; next tick wins */
|
||||
}
|
||||
|
||||
const uint8_t* source = g_base + g_header->data_offset +
|
||||
(seq % FRAME_SHM_RING_SLOTS) *
|
||||
g_header->frame_bytes;
|
||||
memcpy(destination, source, (size_t)g_header->frame_bytes);
|
||||
const uint64_t copied_at = now_ns();
|
||||
const int torn =
|
||||
atomic_load_explicit(&slot->seq, memory_order_acquire) != seq;
|
||||
|
||||
napi_value result;
|
||||
napi_create_object(env, &result);
|
||||
set_named_double(env, result, "seq", (double)seq);
|
||||
set_named_double(env, result, "ageMs",
|
||||
(double)(copied_at - slot->produce_time_ns) / 1e6);
|
||||
napi_value torn_value;
|
||||
napi_get_boolean(env, torn, &torn_value);
|
||||
napi_set_named_property(env, result, "torn", torn_value);
|
||||
return result;
|
||||
}
|
||||
|
||||
/* producerAliveMs() -> ms since the helper's last heartbeat (-1 unknown) */
|
||||
static napi_value ProducerAliveMs(napi_env env, napi_callback_info info) {
|
||||
napi_value result;
|
||||
double ms = -1;
|
||||
if (g_header) {
|
||||
const uint64_t heartbeat =
|
||||
atomic_load_explicit(&g_header->heartbeat_ns, memory_order_relaxed);
|
||||
if (heartbeat) ms = (double)(now_ns() - heartbeat) / 1e6;
|
||||
}
|
||||
napi_create_double(env, ms, &result);
|
||||
return result;
|
||||
}
|
||||
|
||||
/* close() -> undefined */
|
||||
static napi_value Close(napi_env env, napi_callback_info info) {
|
||||
unmap_current();
|
||||
napi_value undefined;
|
||||
napi_get_undefined(env, &undefined);
|
||||
return undefined;
|
||||
}
|
||||
|
||||
static napi_value Init(napi_env env, napi_value exports) {
|
||||
const struct {
|
||||
const char* name;
|
||||
napi_callback fn;
|
||||
} exported[] = {
|
||||
{"open", Open},
|
||||
{"latestSeq", LatestSeq},
|
||||
{"copyLatest", CopyLatest},
|
||||
{"producerAliveMs", ProducerAliveMs},
|
||||
{"close", Close},
|
||||
};
|
||||
for (size_t i = 0; i < sizeof(exported) / sizeof(exported[0]); i++) {
|
||||
napi_value fn;
|
||||
napi_create_function(env, exported[i].name, NAPI_AUTO_LENGTH,
|
||||
exported[i].fn, NULL, &fn);
|
||||
napi_set_named_property(env, exports, exported[i].name, fn);
|
||||
}
|
||||
return exports;
|
||||
}
|
||||
|
||||
#else /* !__APPLE__ */
|
||||
|
||||
static napi_value Init(napi_env env, napi_value exports) {
|
||||
/* Frame-copy is macOS-only for now; loading this module elsewhere
|
||||
* yields an empty exports object and the TS side treats it as
|
||||
* unsupported. */
|
||||
(void)env;
|
||||
return exports;
|
||||
}
|
||||
|
||||
#endif
|
||||
|
||||
NAPI_MODULE(embedded_mpv_frame_reader, Init)
|
||||
Reference in new issue
Block a user