26 Commits

Author SHA1 Message Date
Johanness 4402fb2e12 fix decklink-out: use stream_frame_ counter for DeckLink scheduling instead of TAI-based grains_read_ 2026-05-30 23:26:51 +03:00
Johanness b434c61b86 decklink-out: non-blocking frame pool, skip late frames instead of realign, 5 preroll frames 2026-05-30 23:24:54 +03:00
Johanness bd1f060c73 decklink-out: reuse frames via ScheduledFrameCompleted callback pool 2026-05-30 23:23:15 +03:00
Johanness 26febc41d6 decklink-out: preroll frames before StartScheduledPlayback, detect output format 2026-05-30 23:21:13 +03:00
Johanness a3cd45cc49 decklink-out: add mode support check, limit CreateVideoFrame error spam 2026-05-30 23:18:02 +03:00
Johanness b4f766fb9d add HRESULT logging to CreateVideoFrame for debug 2026-05-30 23:15:00 +03:00
Johanness b5d9c7cd3f fix decklink-out: init grains_read from current time to avoid epoch realignment 2026-05-30 23:13:24 +03:00
Johanness c2e823beda fix realignment 2026-05-30 23:10:35 +03:00
Johanness 003279669e fix output res&frame rate 2026-05-30 23:07:15 +03:00
Johanness 2206fd71b3 some fixes 2026-05-30 22:12:47 +03:00
Johanness 8b5fb3dc36 fix skipped write_index 2026-05-28 23:55:16 +03:00
Johanness 34f72c22f0 hardcoded decklink mode fix 2026-05-28 23:43:26 +03:00
Johanness 96bef6cd32 fix: use relative symlink for extern/mxl (works on any checkout path) 2026-05-28 22:51:45 +03:00
Johanness d02223224f feat: add DeckLink input and output nodes
decklink-in:
- IDeckLinkInputCallback::VideoInputFrameArrived captures V210 frames
- Frame data stored with mutex, copied to MXL grain in process thread
- Uses mxlSleepUntil for TAI-time-based grain pacing
- Configurable: device_index (int), mode (1080i50/1080p50/etc)
- Auto-detect input format via bmdVideoInputEnableFormatDetection
- 1 output port: video_out (V210)

decklink-out:
- Reads V210 grains from MXL, schedules playback via DeckLink output
- Creates IDeckLinkMutableVideoFrame, copies grain data, schedules
- Uses ScheduledFrameCompleted callback for frame completion
- Configurable: device_index (int), mode (1080i50/1080p50/etc)
- 1 input port: video_in (V210)

Both nodes:
- Use DeckLinkAPIDispatch.cpp for CreateDeckLinkIteratorInstance
- Access pixel data via IDeckLinkVideoBuffer (latest SDK API)
- Graceful device open/close on writer/reader add/remove
- Build conditionally via DMF_BUILD_DECKLINK + DECKLINK_SDK_DIR
2026-05-28 22:23:00 +03:00
Johanness db41a91967 Merge feature/phase1-framework: Phase 1 complete
Phase 1: Framework + Passthrough — Two processes exchanging video via MXL, controlled by engine.

Implemented:
- Project scaffold: CMake + vcpkg, MXL SDK submodule
- libdmf-node: Node interface, PortDef, ControlServer (WS + HTTP /cmd),
  NodeRunner (CLI, MXL lifecycle, separate process thread)
- libdmf-engine: Graph model, FlowManager (UUID v4, NMOS V210 flow defs),
  ProcessManager (fork/exec), ApiServer (REST), NodeControlClient
- dmf-node-passthrough: V210 grain read→memcpy→write, TAI-time-based
  grain pacing with mxlSleepUntil, 2-grain read delay
- dmf-studio-engine: REST API with full graph CRUD + start/stop +
  connect-input/connect-output/disconnect-port/command endpoints
- Proper shutdown: release MXL readers/writers before destroy, engine
  kills child processes on exit
- Unit tests: 6 cases, 22 assertions, all passing
2026-05-26 23:04:45 +03:00
Johanness 4b50f49172 chore: phase 1 cleanup — proper shutdown, flow resource cleanup, test fix
- node_runner: track all readers/writers in flow_resources vector,
  release them via mxlReleaseFlowWriter/mxlReleaseFlowReader before
  destroying the MXL instance (fixes 'leaked flow writer' warning)
- node_runner: remove_writer/remove_reader commands now also release
  the MXL flow resources, not just reset the node's optional<>
- engine: call process_manager.stop_all(graph) on shutdown to kill
  child node processes (prevents orphaned passthrough processes)
- graph: add get_node_mut() for process_manager to update node state
- process_manager: implement stop_all(Graph&) that SIGTERMs all
  running node processes
- passthrough: reduce logging to first-grain and realign-once only
- tests: update flow format assertion to match NMOS
  (urn:x-nmos:format:video instead of video/v210)
2026-05-26 23:04:16 +03:00
Johanness 5b48d86c6e fix: separate grain processing thread from LWS event loop
The fundamental issue: mxlSleepUntil/mxlSleepForNs blocks the entire
thread, making LWS unresponsive and causing control commands to hang.
Meanwhile, LWS poll(1) adds latency that makes grain timing unreliable.

Solution: run node->process() in a dedicated thread that can freely use
mxlSleepUntil for precise grain timing, while the main thread runs
control_server->poll(10) for LWS event handling.

Passthrough now uses mxlSleepUntil(deadline) matching the mxl-gst-sink
Cursor pattern, with 2-grain read delay for buffering.
2026-05-26 22:55:57 +03:00
Johanness d677b654f0 fix: non-blocking grain timing — use mxlGetNsUntilIndex instead of mxlSleepUntil
mxlSleepUntil blocks the entire thread including LWS poll, causing the
control server to become unresponsive and grains to be missed. Instead,
check mxlGetNsUntilIndex and skip process() if the grain isn't due yet
(>2ms away). The node loop spins with poll(1) keeping LWS responsive,
and only attempts a grain read when it's nearly due.
2026-05-26 22:51:57 +03:00
Johanness 503023f3e5 fix: use mxlSleepUntil delivery deadline pattern for grain timing
Match the mxl-gst-sink pattern: sleep until the next grain's delivery
deadline using mxlSleepUntil, then read with a short 5ms timeout.
Uses 2-grain read delay (40ms at 50fps) for buffering headroom.

Previous approach of mxlSleepForNs + 20ms GetGrain timeout caused the
passthrough to fall behind: each iteration took ~25ms (poll+sleep+read),
missing grains and perpetually chasing headIndex via realign.
2026-05-26 22:46:33 +03:00
Johanness 2c712592dd fix: use TAI-time-based grain index alignment in passthrough
Instead of chasing headIndex from mxlFlowReaderGetRuntimeInfo (which
points to the NEXT grain to be written, causing perpetual TOO_EARLY),
the passthrough now uses mxlTimestampToIndex + mxlGetNsUntilIndex for
proper timing alignment, matching the pattern used by mxl-gst-sink.

Key changes:
- realign() computes read_index from current TAI time minus 1 grain delay
- process() uses mxlSleepForNs to wait until the target grain is due
- on_add_reader fetches grain_rate from mxlFlowConfigInfo
- Removed separate write_index_ (uses grain_info.index for writer)
2026-05-26 22:40:46 +03:00
Johanness 3e1477c56c fix: LWS HTTP connection leak, grain index alignment, and reader head tracking
- control_server: add Connection: close header + always return -1 after
  HTTP response to force-close connection. Without this, lws_service()
  blocks forever after the first POST /cmd, freezing the node process loop.
- passthrough: initialize read_index_ to runtime.headIndex when reader
  is added (prevents TOO_EARLY errors on first read)
- passthrough: handle MXL_ERR_OUT_OF_RANGE_TOO_EARLY in addition to
  TOO_LATE (both jump to head index)
- passthrough: use grain_info.index as writer index instead of separate
  write_index_ counter (MXL writers must use TAI-based grain indices)
- passthrough: reduce grain read timeout from 100ms to 20ms for tighter
  loop with LWS poll
- node_runner: add 'status' command handler that sends node status
  back via control server
2026-05-26 22:22:46 +03:00
Johanness 5b2d420e71 feat: add node HTTP control, engine-to-node communication, connect-input/output API
- Node control server now accepts HTTP POST /cmd for command dispatch
  (in addition to existing WebSocket control)
- Added NodeControlClient: engine sends commands to nodes via HTTP POST
- New REST endpoints:
  POST /api/graph/nodes/:id/connect-input  - connect node input to MXL flow
  POST /api/graph/nodes/:id/connect-output - create new MXL flow + connect output
  POST /api/graph/nodes/:id/disconnect-port - remove reader/writer
  POST /api/graph/nodes/:id/command - send raw command to node
- Flow IDs now use proper UUID v4 format (MXL requires standard UUIDs)
- Flow definitions use NMOS format (urn:x-nmos:format:video)
- Engine passes --mxl-domain to node processes
- User-provided node IDs (via 'id' field in POST body)
- End-to-end verified: testsrc → passthrough → new MXL output flow
2026-05-26 21:33:17 +03:00
Johanness 1d6f93679f fix: engine REST API - handle GET/DELETE requests immediately, use user-provided node IDs
- Use lws_http_get_uri_and_method() for proper HTTP method detection
  (GET, POST, PUT, DELETE all supported)
- Handle GET/DELETE in LWS_CALLBACK_HTTP without waiting for body
- Support user-provided node IDs via config.id field
- Fixes all REST API endpoints hanging on GET requests
2026-05-26 21:08:06 +03:00
Johanness 37f02f01d9 docs: add build and test instructions 2026-05-26 01:42:44 +03:00
Johanness d138478c1a fix: WebSocket control server - add HTTP mount for WS upgrade negotiation
LWS v4.x requires an HTTP mount with LWSMPRO_CALLBACK to properly
route WebSocket upgrade requests to the dmf-control protocol handler.
Without this, WS connections were rejected with 403 Forbidden.
2026-05-26 01:40:42 +03:00
Johanness 5b5ffa1308 feat: scaffold project framework with engine, node skeleton, and passthrough node
- CMake build system with vcpkg dependency management
- libdmf-node: Node interface, port definitions, WebSocket control
  server, node runner with MXL lifecycle management
- libdmf-engine: Graph model (nodes/edges CRUD, serialization),
  FlowManager (UUID + V210 NMOS flow definitions), ProcessManager
  (fork/exec node processes), ApiServer (REST API for graph control)
- Passthrough node: reads V210 grains from MXL, copies, writes to MXL
- Unit tests for Graph and FlowManager (6 cases, 21 assertions)
- MXL SDK as git submodule (symlinked from extern/mxl)
2026-05-26 01:32:50 +03:00
43 changed files with 3360 additions and 0 deletions
+12
View File
@@ -1 +1,13 @@
ref_arch.pdf
build/
.cache/
CMakeUserPresets.json
compile_commands.json
.vcpkg/
vcpkg_installed/
node_modules/
web/dist/
*.o
*.a
*.so
*.d
+40
View File
@@ -0,0 +1,40 @@
cmake_minimum_required(VERSION 3.24 FATAL_ERROR)
project(dmf-studio
VERSION 0.1.0
LANGUAGES CXX C
)
set(CMAKE_CXX_STANDARD 20)
set(CMAKE_CXX_STANDARD_REQUIRED ON)
set(CMAKE_EXPORT_COMPILE_COMMANDS ON)
list(APPEND CMAKE_MODULE_PATH "${CMAKE_CURRENT_SOURCE_DIR}/cmake")
option(DMF_BUILD_TESTS "Build tests" ON)
option(DMF_BUILD_MXL_TOOLS "Build MXL tools (testsrc, sink)" OFF)
option(DMF_BUILD_DECKLINK "Build DeckLink I/O nodes" ON)
set(DECKLINK_SDK_DIR "" CACHE PATH "Path to Blackmagic DeckLink SDK root")
if(DMF_BUILD_DECKLINK AND DECKLINK_SDK_DIR)
add_subdirectory(nodes/decklink-in)
add_subdirectory(nodes/decklink-out)
endif()
find_package(fmt CONFIG REQUIRED)
find_package(spdlog CONFIG REQUIRED)
find_package(nlohmann_json CONFIG REQUIRED)
find_package(Libwebsockets CONFIG REQUIRED)
find_package(Catch2 CONFIG QUIET)
set(BUILD_TESTS OFF CACHE BOOL "" FORCE)
add_subdirectory(extern/mxl)
add_subdirectory(libs/dmf-node)
add_subdirectory(libs/dmf-engine)
add_subdirectory(nodes/passthrough)
add_subdirectory(engine)
if(DMF_BUILD_TESTS AND Catch2_FOUND)
add_subdirectory(tests)
endif()
+187
View File
@@ -0,0 +1,187 @@
# DMF Studio
Node-based visual production platform for the Dynamic Media Facility architecture.
## Prerequisites
- CMake 3.24+
- C++20 compiler (GCC 12+, Clang 15+)
- vcpkg (at `~/vcpkg` or set `CMAKE_TOOLCHAIN_FILE`)
- GStreamer (for MXL test tools only)
## Build
```bash
# Configure (from project root)
cmake -B build \
-DCMAKE_BUILD_TYPE=Debug \
-DCMAKE_TOOLCHAIN_FILE=$HOME/vcpkg/scripts/buildsystems/vcpkg.cmake \
-DBUILD_SHARED_LIBS=OFF \
-DBUILD_DOCS=OFF \
-DBUILD_TESTS=OFF \
-DBUILD_TOOLS=OFF \
-DBUILD_UTILS=OFF
# Build all
cmake --build build -j$(nproc)
```
Binaries end up in:
- `build/engine/dmf-studio-engine`
- `build/nodes/passthrough/dmf-node-passthrough`
## Rebuild after changes
```bash
# Full rebuild
cmake --build build -j$(nproc)
# Rebuild only one target (faster)
cmake --build build -j$(nproc) --target dmf-node-passthrough
cmake --build build -j$(nproc) --target dmf-studio-engine
```
## Clean rebuild
```bash
rm -rf build
# Then re-run the configure + build steps above
```
## Run unit tests
```bash
build/tests/dmf-test-graph
```
## Test with MXL
### 1. Create MXL domain
```bash
mkdir -p /tmp/dmf-mxl
```
### 2. Start MXL test source (writes V210 video flow)
Create a flow config file:
```bash
cat > /tmp/v210_50p.json << 'EOF'
{
"id": "a0000001-0000-0000-0000-000000000001",
"description": "DMF Studio test video flow",
"format": "urn:x-nmos:format:video",
"label": "DMF Studio Test Video",
"tags": {
"urn:x-nmos:tag:grouphint/v1.0": ["dmf-studio:Video"]
},
"media_type": "video/v210",
"grain_rate": {"numerator": 50, "denominator": 1},
"frame_width": 1920,
"frame_height": 1080,
"interlace_mode": "progressive",
"colorspace": "BT709",
"components": [
{"name": "Y", "width": 1920, "height": 1080, "bit_depth": 10},
{"name": "Cb", "width": 960, "height": 1080, "bit_depth": 10},
{"name": "Cr", "width": 960, "height": 1080, "bit_depth": 10}
]
}
EOF
```
Start test source (needs GStreamer + MXL tools built separately):
```bash
mxl-gst-testsrc -d /tmp/dmf-mxl -v /tmp/v210_50p.json --pattern smpte &
```
Check active flows:
```bash
mxl-info --domain /tmp/dmf-mxl
```
### 3. Start passthrough node
```bash
build/nodes/passthrough/dmf-node-passthrough \
--node-id pass1 \
--control-port 9100 \
--mxl-domain /tmp/dmf-mxl
```
Options:
- `--node-id, -n` — Unique node instance ID (required)
- `--control-port, -p` — WebSocket control port (required)
- `--mxl-domain, -d` — MXL domain path (default: `/dev/shm/mxl`)
- `--config, -c` — Node config as JSON string
### 4. Start engine
```bash
export DMF_STUDIO_BIN_DIR=build/nodes/passthrough
build/engine/dmf-studio-engine --port 8080
```
### 5. Control via REST API
```bash
# Add nodes
curl -X POST http://localhost:8080/api/graph/nodes \
-H "Content-Type: application/json" \
-d '{"type":"passthrough","id":"pass1"}'
curl -X POST http://localhost:8080/api/graph/nodes \
-H "Content-Type: application/json" \
-d '{"type":"passthrough","id":"pass2"}'
# Connect nodes (creates MXL flow between them)
curl -X POST http://localhost:8080/api/graph/edges \
-H "Content-Type: application/json" \
-d '{"from_node":"pass1","from_port":"video_out","to_node":"pass2","to_port":"video_in"}'
# View graph
curl http://localhost:8080/api/graph
# Start all nodes
curl -X POST http://localhost:8080/api/graph/start
# Stop all nodes
curl -X POST http://localhost:8080/api/graph/stop
# Remove node
curl -X DELETE http://localhost:8080/api/graph/nodes/pass1
# Remove edge
curl -X DELETE http://localhost:8080/api/graph/edges/pass1:video_out->pass2:video_in
```
### 6. Control node directly via WebSocket
Connect to `ws://localhost:9100` with subprotocol `dmf-control`:
```bash
wscat -c ws://localhost:9100 -s dmf-control
```
Commands:
```json
{"cmd": "add_reader", "flow_id": "<uuid>", "port_id": "video_in"}
{"cmd": "add_writer", "flow_id": "<uuid>", "port_id": "video_out", "flow_def": {<NMOS flow JSON>}}
{"cmd": "remove_reader", "port_id": "video_in"}
{"cmd": "remove_writer", "port_id": "video_out"}
{"cmd": "configure", "params": {}}
{"cmd": "shutdown"}
```
## Project structure
```
libs/dmf-node/ — Node skeleton library (interface, WS control, MXL lifecycle)
libs/dmf-engine/ — Engine library (graph model, flow manager, process manager, REST API)
nodes/passthrough/ — Passthrough node (1 MXL in → 1 MXL out, memcpy)
engine/ — Engine binary
tests/ — Unit tests
```
+66
View File
@@ -0,0 +1,66 @@
#!/bin/bash
set -e
ENGINE_PORT=${ENGINE_PORT:-9000}
BASE_URL="http://127.0.0.1:${ENGINE_PORT}"
DEVICE_INDEX=${DEVICE_INDEX:-0}
MODE=${MODE:-1080i50}
log() { echo "=== $1 ==="; }
log "Adding decklink-in node (device=${DEVICE_INDEX}, mode=${MODE})"
curl -s -X POST "${BASE_URL}/api/graph/nodes" \
-H "Content-Type: application/json" \
-d "{\"type\":\"decklink-in\",\"id\":\"sdi_in\",\"config\":{\"device_index\":${DEVICE_INDEX},\"mode\":\"${MODE}\"}}" | python3 -m json.tool 2>/dev/null || echo ""
log "Adding passthrough node"
curl -s -X POST "${BASE_URL}/api/graph/nodes" \
-H "Content-Type: application/json" \
-d '{"type":"passthrough","id":"pass1"}' | python3 -m json.tool 2>/dev/null || echo ""
log "Adding decklink-out node (device=${DEVICE_INDEX}, mode=${MODE})"
curl -s -X POST "${BASE_URL}/api/graph/nodes" \
-H "Content-Type: application/json" \
-d "{\"type\":\"decklink-out\",\"id\":\"sdi_out\",\"config\":{\"device_index\":${DEVICE_INDEX},\"mode\":\"${MODE}\"}}" | python3 -m json.tool 2>/dev/null || echo ""
log "Connecting sdi_in → pass1"
curl -s -X POST "${BASE_URL}/api/graph/edges" \
-H "Content-Type: application/json" \
-d '{"from_node":"sdi_in","from_port":"video_out","to_node":"pass1","to_port":"video_in"}' | python3 -m json.tool 2>/dev/null || echo ""
log "Connecting pass1 → sdi_out"
curl -s -X POST "${BASE_URL}/api/graph/edges" \
-H "Content-Type: application/json" \
-d '{"from_node":"pass1","from_port":"video_out","to_node":"sdi_out","to_port":"video_in"}' | python3 -m json.tool 2>/dev/null || echo ""
log "Starting graph"
curl -s -X POST "${BASE_URL}/api/graph/start" | python3 -m json.tool 2>/dev/null || echo ""
sleep 3
log "Status: sdi_in"
curl -s -X POST "${BASE_URL}/api/graph/nodes/sdi_in/command" \
-H "Content-Type: application/json" \
-d '{"cmd":"status"}' | python3 -m json.tool 2>/dev/null || echo ""
log "Status: pass1"
curl -s -X POST "${BASE_URL}/api/graph/nodes/pass1/command" \
-H "Content-Type: application/json" \
-d '{"cmd":"status"}' | python3 -m json.tool 2>/dev/null || echo ""
log "Status: sdi_out"
curl -s -X POST "${BASE_URL}/api/graph/nodes/sdi_out/command" \
-H "Content-Type: application/json" \
-d '{"cmd":"status"}' | python3 -m json.tool 2>/dev/null || echo ""
log "Graph state"
curl -s "${BASE_URL}/api/graph" | python3 -m json.tool 2>/dev/null || echo ""
echo ""
echo "=== SDI pipeline running. Press Enter to stop. ==="
read
log "Stopping graph"
curl -s -X POST "${BASE_URL}/api/graph/stop" | python3 -m json.tool 2>/dev/null || echo ""
log "Done"
+5
View File
@@ -0,0 +1,5 @@
add_executable(dmf-studio-engine
src/main.cpp
)
target_link_libraries(dmf-studio-engine PRIVATE dmf-engine)
+57
View File
@@ -0,0 +1,57 @@
#include <dmf-engine/api_server.hpp>
#include <dmf-engine/graph.hpp>
#include <dmf-engine/flow_manager.hpp>
#include <dmf-engine/process_manager.hpp>
#include <dmf-engine/node_control_client.hpp>
#include <spdlog/spdlog.h>
#include <csignal>
#include <atomic>
static std::atomic<bool> g_running{true};
static void signal_handler(int /*signum*/) {
g_running = false;
}
int main(int argc, char* argv[]) {
uint16_t port = 8080;
std::string mxl_domain = "/dev/shm/mxl";
for (int i = 1; i < argc; ++i) {
std::string arg = argv[i];
if ((arg == "--port" || arg == "-p") && i + 1 < argc) {
port = static_cast<uint16_t>(std::stoi(argv[++i]));
} else if ((arg == "--mxl-domain" || arg == "-d") && i + 1 < argc) {
mxl_domain = argv[++i];
} else if (arg == "--help" || arg == "-h") {
spdlog::info("Usage: dmf-studio-engine [options]");
spdlog::info(" --port, -p API server port (default: 8080)");
spdlog::info(" --mxl-domain, -d MXL domain path (default: /dev/shm/mxl)");
return 0;
}
}
std::signal(SIGINT, signal_handler);
std::signal(SIGTERM, signal_handler);
spdlog::info("DMF Studio Engine starting on port {}", port);
dmf_engine::Graph graph;
dmf_engine::FlowManager flow_manager;
dmf_engine::ProcessManager process_manager;
dmf_engine::NodeControlClient control_client;
dmf_engine::ApiServer api_server(port, graph, flow_manager, process_manager, control_client, mxl_domain);
spdlog::info("DMF Studio Engine ready");
while (g_running) {
api_server.poll(100);
}
spdlog::info("DMF Studio Engine shutting down");
process_manager.stop_all(graph);
return 0;
}
Vendored Symlink
+1
View File
@@ -0,0 +1 @@
../mxl
+19
View File
@@ -0,0 +1,19 @@
add_library(dmf-engine STATIC
src/graph.cpp
src/flow_manager.cpp
src/process_manager.cpp
src/api_server.cpp
src/node_control_client.cpp
)
target_include_directories(dmf-engine PUBLIC
include
)
target_link_libraries(dmf-engine PUBLIC
dmf-node
nlohmann_json::nlohmann_json
spdlog::spdlog
fmt::fmt
websockets
)
@@ -0,0 +1,33 @@
#pragma once
#include <dmf-engine/graph.hpp>
#include <functional>
#include <memory>
#include <string>
struct lws_context;
namespace dmf_engine {
using RequestHandler = std::function<std::string(const std::string& method, const std::string& path, const std::string& body)>;
class ApiServer {
public:
ApiServer(uint16_t port, Graph& graph, class FlowManager& flow_manager, class ProcessManager& process_manager, class NodeControlClient& control_client, const std::string& mxl_domain = "/dev/shm/mxl");
~ApiServer();
ApiServer(const ApiServer&) = delete;
ApiServer& operator=(const ApiServer&) = delete;
void poll(int timeout_ms);
private:
void handle_rest(const std::string& method, const std::string& path, const std::string& body,
std::string& response_status, std::string& response_content_type, std::string& response_body);
struct Impl;
std::unique_ptr<Impl> impl_;
};
} // namespace dmf_engine
@@ -0,0 +1,22 @@
#pragma once
#include <dmf-engine/types.hpp>
#include <nlohmann/json.hpp>
#include <string>
namespace dmf_engine {
class FlowManager {
public:
FlowManager();
FlowId create_flow_id();
nlohmann::json create_v210_flow_def(const FlowId& flow_id, int width, int height, int fps_numerator, int fps_denominator) const;
private:
int flow_counter_ = 0;
};
} // namespace dmf_engine
@@ -0,0 +1,63 @@
#pragma once
#include <dmf-engine/types.hpp>
#include <nlohmann/json.hpp>
#include <optional>
#include <string>
#include <unordered_map>
#include <vector>
namespace dmf_engine {
enum class NodeState {
Stopped,
Starting,
Running,
Stopping,
};
struct GraphNode {
NodeId id;
std::string type;
nlohmann::json config;
NodeState state = NodeState::Stopped;
uint16_t control_port = 0;
int pid = 0;
};
struct GraphEdge {
EdgeId id;
NodeId from_node;
PortId from_port;
NodeId to_node;
PortId to_port;
FlowId flow_id;
nlohmann::json flow_def;
};
class Graph {
public:
NodeId add_node(const std::string& type, const nlohmann::json& config = {});
bool remove_node(const NodeId& node_id);
const GraphNode* get_node(const NodeId& node_id) const;
GraphNode* get_node_mut(const NodeId& node_id);
std::vector<GraphNode> get_nodes() const;
EdgeId add_edge(const NodeId& from_node, const PortId& from_port,
const NodeId& to_node, const PortId& to_port,
const FlowId& flow_id, const nlohmann::json& flow_def);
bool remove_edge(const EdgeId& edge_id);
std::vector<GraphEdge> get_edges() const;
std::vector<GraphEdge> get_edges_for_node(const NodeId& node_id) const;
nlohmann::json serialize() const;
private:
std::unordered_map<NodeId, GraphNode> nodes_;
std::unordered_map<EdgeId, GraphEdge> edges_;
int next_node_num_ = 0;
};
} // namespace dmf_engine
@@ -0,0 +1,22 @@
#pragma once
#include <string>
#include <unordered_map>
#include <cstdint>
namespace dmf_engine {
class NodeControlClient {
public:
bool send_command(uint16_t port, const std::string& json_cmd);
std::string send_command_with_response(uint16_t port, const std::string& json_cmd);
void register_node(const std::string& node_id, uint16_t port);
void unregister_node(const std::string& node_id);
uint16_t get_port(const std::string& node_id) const;
private:
std::unordered_map<std::string, uint16_t> node_ports_;
};
} // namespace dmf_engine
@@ -0,0 +1,23 @@
#pragma once
#include <dmf-engine/graph.hpp>
#include <string>
#include <unordered_map>
namespace dmf_engine {
class ProcessManager {
public:
bool start_node(GraphNode& node, const std::string& mxl_domain, uint16_t base_port);
bool stop_node(GraphNode& node);
void stop_all();
void stop_all(Graph& graph);
bool is_running(const NodeId& node_id) const;
private:
std::string find_node_binary(const std::string& node_type) const;
};
} // namespace dmf_engine
@@ -0,0 +1,13 @@
#pragma once
#include <string>
#include <cstdint>
namespace dmf_engine {
using NodeId = std::string;
using EdgeId = std::string;
using FlowId = std::string;
using PortId = std::string;
} // namespace dmf_engine
+448
View File
@@ -0,0 +1,448 @@
#include <dmf-engine/api_server.hpp>
#include <dmf-engine/graph.hpp>
#include <dmf-engine/flow_manager.hpp>
#include <dmf-engine/process_manager.hpp>
#include <dmf-engine/node_control_client.hpp>
#include <libwebsockets.h>
#include <nlohmann/json.hpp>
#include <spdlog/spdlog.h>
#include <cstring>
#include <string>
namespace dmf_engine {
struct ApiServerImpl {
uint16_t port = 0;
Graph* graph = nullptr;
FlowManager* flow_manager = nullptr;
ProcessManager* process_manager = nullptr;
NodeControlClient* control_client = nullptr;
std::string mxl_domain = "/dev/shm/mxl";
struct lws_context* context = nullptr;
};
struct ApiServer::Impl {
std::unique_ptr<ApiServerImpl> data;
};
struct HttpRequest {
std::string method;
std::string path;
std::string body;
};
static ApiServerImpl* g_impl = nullptr;
static int send_json_response(struct lws* wsi, const std::string& status_str,
const std::string& json_body) {
auto headers = "HTTP/1.1 " + status_str + "\r\n"
"Content-Type: application/json\r\n"
"Access-Control-Allow-Origin: *\r\n"
"Access-Control-Allow-Methods: GET, POST, DELETE, PUT, OPTIONS\r\n"
"Access-Control-Allow-Headers: Content-Type\r\n"
"Content-Length: " + std::to_string(json_body.size()) + "\r\n"
"\r\n";
std::vector<uint8_t> buf(LWS_PRE + headers.size() + json_body.size());
std::memcpy(buf.data() + LWS_PRE, headers.data(), headers.size());
std::memcpy(buf.data() + LWS_PRE + headers.size(), json_body.data(), json_body.size());
lws_write(wsi, buf.data() + LWS_PRE, headers.size() + json_body.size(), LWS_WRITE_HTTP);
if (lws_http_transaction_completed(wsi)) {
return -1;
}
return 0;
}
static void handle_request(const std::string& method, const std::string& path,
const std::string& body,
std::string& status_str, std::string& response_body) {
if (!g_impl || !g_impl->graph) {
status_str = "500 Internal Server Error";
response_body = nlohmann::json({{"error", "Server not initialized"}}).dump();
return;
}
auto ok = [&](const nlohmann::json& j) {
status_str = "200 OK";
response_body = j.dump();
};
auto created = [&](const nlohmann::json& j) {
status_str = "201 Created";
response_body = j.dump();
};
auto error_resp = [&](int code, const std::string& msg) {
status_str = std::to_string(code) + " Error";
response_body = nlohmann::json({{"error", msg}}).dump();
};
try {
nlohmann::json req_body = body.empty() ? nlohmann::json::object() : nlohmann::json::parse(body);
auto& graph = *g_impl->graph;
auto& fm = *g_impl->flow_manager;
auto& pm = *g_impl->process_manager;
auto& cc = *g_impl->control_client;
if (path == "/api/graph" && method == "GET") {
ok(graph.serialize());
} else if (path == "/api/graph/nodes" && method == "POST") {
if (!req_body.contains("type")) {
error_resp(400, "Missing 'type' field");
} else {
auto type = req_body["type"].get<std::string>();
auto config = req_body.value("config", nlohmann::json::object());
if (req_body.contains("id") && req_body["id"].is_string()) {
config["id"] = req_body["id"];
}
auto id = graph.add_node(type, config);
created({{"id", id}});
}
} else if (path.find("/api/graph/nodes/") == 0 && method == "DELETE") {
auto node_id = path.substr(std::string("/api/graph/nodes/").length());
cc.unregister_node(node_id);
if (graph.remove_node(node_id)) {
ok({{"deleted", node_id}});
} else {
error_resp(404, "Node not found: " + node_id);
}
} else if (path == "/api/graph/edges" && method == "POST") {
if (!req_body.contains("from_node") || !req_body.contains("to_node")) {
error_resp(400, "Missing from_node/to_node");
} else {
auto from_node = req_body["from_node"].get<std::string>();
auto from_port = req_body.value("from_port", "video_out");
auto to_node = req_body["to_node"].get<std::string>();
auto to_port = req_body.value("to_port", "video_in");
auto flow_id = fm.create_flow_id();
auto flow_def = fm.create_v210_flow_def(flow_id, 1920, 1080, 50, 1);
auto edge_id = graph.add_edge(from_node, from_port, to_node, to_port, flow_id, flow_def);
auto from_port_num = cc.get_port(from_node);
auto to_port_num = cc.get_port(to_node);
if (from_port_num > 0) {
nlohmann::json cmd;
cmd["cmd"] = "add_writer";
cmd["port_id"] = from_port;
cmd["flow_id"] = flow_id;
cmd["flow_def"] = flow_def;
cc.send_command(from_port_num, cmd.dump());
}
if (to_port_num > 0) {
nlohmann::json cmd;
cmd["cmd"] = "add_reader";
cmd["port_id"] = to_port;
cmd["flow_id"] = flow_id;
cc.send_command(to_port_num, cmd.dump());
}
created({{"id", edge_id}, {"flow_id", flow_id}});
}
} else if (path.find("/api/graph/edges/") == 0 && method == "DELETE") {
auto edge_id = path.substr(std::string("/api/graph/edges/").length());
auto edges = graph.get_edges();
const GraphEdge* edge = nullptr;
for (auto& e : edges) {
if (e.id == edge_id) { edge = &e; break; }
}
if (graph.remove_edge(edge_id)) {
if (edge) {
auto from_port_num = cc.get_port(edge->from_node);
auto to_port_num = cc.get_port(edge->to_node);
if (from_port_num > 0) {
nlohmann::json cmd;
cmd["cmd"] = "remove_writer";
cmd["port_id"] = edge->from_port;
cc.send_command(from_port_num, cmd.dump());
}
if (to_port_num > 0) {
nlohmann::json cmd;
cmd["cmd"] = "remove_reader";
cmd["port_id"] = edge->to_port;
cc.send_command(to_port_num, cmd.dump());
}
}
ok({{"deleted", edge_id}});
} else {
error_resp(404, "Edge not found: " + edge_id);
}
} else if (path == "/api/graph/start" && method == "POST") {
auto nodes = graph.get_nodes();
uint16_t port = 9100;
for (auto& node : nodes) {
pm.start_node(const_cast<GraphNode&>(node), g_impl->mxl_domain, port);
cc.register_node(node.id, port);
port++;
}
for (auto& edge : graph.get_edges()) {
auto from_port_num = cc.get_port(edge.from_node);
auto to_port_num = cc.get_port(edge.to_node);
if (from_port_num > 0) {
nlohmann::json cmd;
cmd["cmd"] = "add_writer";
cmd["port_id"] = edge.from_port;
cmd["flow_id"] = edge.flow_id;
cmd["flow_def"] = edge.flow_def;
cc.send_command(from_port_num, cmd.dump());
}
if (to_port_num > 0) {
nlohmann::json cmd;
cmd["cmd"] = "add_reader";
cmd["port_id"] = edge.to_port;
cmd["flow_id"] = edge.flow_id;
cc.send_command(to_port_num, cmd.dump());
}
}
ok({{"status", "started"}});
} else if (path == "/api/graph/stop" && method == "POST") {
auto nodes = graph.get_nodes();
for (auto& node : nodes) {
cc.unregister_node(node.id);
pm.stop_node(const_cast<GraphNode&>(node));
}
ok({{"status", "stopped"}});
} else if (path.find("/api/graph/nodes/") == 0 && path.find("/connect-input") != std::string::npos && method == "POST") {
auto prefix = std::string("/api/graph/nodes/");
auto suffix_start = path.find("/connect-input");
auto node_id = path.substr(prefix.length(), suffix_start - prefix.length());
if (!req_body.contains("flow_id") || !req_body.contains("port_id")) {
error_resp(400, "Missing flow_id/port_id");
return;
}
auto port_num = cc.get_port(node_id);
if (port_num == 0) {
error_resp(404, "Node not running or not found: " + node_id);
return;
}
nlohmann::json cmd;
cmd["cmd"] = "add_reader";
cmd["port_id"] = req_body["port_id"].get<std::string>();
cmd["flow_id"] = req_body["flow_id"].get<std::string>();
if (cc.send_command(port_num, cmd.dump())) {
ok({{"node_id", node_id}, {"connected_input", req_body["port_id"]}, {"flow_id", req_body["flow_id"]}});
} else {
error_resp(500, "Failed to send command to node");
}
} else if (path.find("/api/graph/nodes/") == 0 && path.find("/connect-output") != std::string::npos && method == "POST") {
auto prefix = std::string("/api/graph/nodes/");
auto suffix_start = path.find("/connect-output");
auto node_id = path.substr(prefix.length(), suffix_start - prefix.length());
auto port_id = req_body.value("port_id", "video_out");
auto port_num = cc.get_port(node_id);
if (port_num == 0) {
error_resp(404, "Node not running or not found: " + node_id);
return;
}
int fps_num = 50, fps_den = 1;
int width = 1920, height = 1080;
nlohmann::json status_cmd;
status_cmd["cmd"] = "status";
auto status_resp = cc.send_command_with_response(port_num, status_cmd.dump());
if (!status_resp.empty()) {
try {
auto sr = nlohmann::json::parse(status_resp);
if (sr.contains("data")) {
auto& d = sr["data"];
if (d.contains("grain_rate")) {
fps_num = d["grain_rate"].value("numerator", fps_num);
fps_den = d["grain_rate"].value("denominator", fps_den);
}
width = d.value("width", width);
height = d.value("height", height);
}
} catch (...) {}
}
auto flow_id = fm.create_flow_id();
auto flow_def = fm.create_v210_flow_def(flow_id, width, height, fps_num, fps_den);
nlohmann::json cmd;
cmd["cmd"] = "add_writer";
cmd["port_id"] = port_id;
cmd["flow_id"] = flow_id;
cmd["flow_def"] = flow_def;
if (cc.send_command(port_num, cmd.dump())) {
ok({{"node_id", node_id}, {"connected_output", port_id}, {"flow_id", flow_id}});
} else {
error_resp(500, "Failed to send command to node");
}
} else if (path.find("/api/graph/nodes/") == 0 && path.find("/disconnect-port") != std::string::npos && method == "POST") {
auto prefix = std::string("/api/graph/nodes/");
auto suffix_start = path.find("/disconnect-port");
auto node_id = path.substr(prefix.length(), suffix_start - prefix.length());
if (!req_body.contains("port_id")) {
error_resp(400, "Missing port_id");
return;
}
auto port_id = req_body["port_id"].get<std::string>();
auto port_num = cc.get_port(node_id);
if (port_num == 0) {
error_resp(404, "Node not running or not found: " + node_id);
return;
}
auto direction = req_body.value("direction", "input");
nlohmann::json cmd;
if (direction == "output") {
cmd["cmd"] = "remove_writer";
} else {
cmd["cmd"] = "remove_reader";
}
cmd["port_id"] = port_id;
if (cc.send_command(port_num, cmd.dump())) {
ok({{"node_id", node_id}, {"disconnected", port_id}});
} else {
error_resp(500, "Failed to send command to node");
}
} else if (path.find("/api/graph/nodes/") == 0 && path.find("/command") != std::string::npos && method == "POST") {
auto prefix = std::string("/api/graph/nodes/");
auto suffix_start = path.find("/command");
auto node_id = path.substr(prefix.length(), suffix_start - prefix.length());
auto port_num = cc.get_port(node_id);
if (port_num == 0) {
error_resp(404, "Node not running or not found: " + node_id);
return;
}
if (cc.send_command(port_num, req_body.dump())) {
ok({{"node_id", node_id}, {"sent", true}});
} else {
error_resp(500, "Failed to send command to node");
}
} else {
error_resp(404, "Not found: " + method + " " + path);
}
} catch (const nlohmann::json::exception& e) {
error_resp(400, std::string("JSON error: ") + e.what());
} catch (const std::exception& e) {
error_resp(500, e.what());
}
}
static int callback_http(struct lws* wsi, enum lws_callback_reasons reason,
void* user, void* in, size_t len) {
auto* req = static_cast<HttpRequest*>(user);
switch (reason) {
case LWS_CALLBACK_HTTP: {
new (req) HttpRequest();
char* uri_ptr = nullptr;
int uri_len = 0;
int method = lws_http_get_uri_and_method(wsi, &uri_ptr, &uri_len);
switch (method) {
case LWSHUMETH_GET: req->method = "GET"; break;
case LWSHUMETH_POST: req->method = "POST"; break;
case LWSHUMETH_PUT: req->method = "PUT"; break;
case LWSHUMETH_DELETE: req->method = "DELETE"; break;
default: req->method = "GET"; break;
}
if (uri_ptr && uri_len > 0) {
req->path.assign(uri_ptr, uri_len);
}
if (req->method == "GET" || req->method == "DELETE") {
std::string status_str, response_body;
handle_request(req->method, req->path, "", status_str, response_body);
return send_json_response(wsi, status_str, response_body);
}
int body_len = lws_hdr_total_length(wsi, WSI_TOKEN_HTTP_CONTENT_LENGTH);
if (body_len == 0) {
std::string status_str, response_body;
handle_request(req->method, req->path, "", status_str, response_body);
return send_json_response(wsi, status_str, response_body);
}
req->body.reserve(body_len);
break;
}
case LWS_CALLBACK_HTTP_BODY: {
req->body.append(static_cast<char*>(in), len);
break;
}
case LWS_CALLBACK_HTTP_BODY_COMPLETION: {
std::string status_str, response_body;
handle_request(req->method, req->path, req->body, status_str, response_body);
return send_json_response(wsi, status_str, response_body);
}
default:
break;
}
return 0;
}
static const struct lws_protocols protocols[] = {
{"http-api", callback_http, sizeof(HttpRequest), 4096},
{nullptr, nullptr, 0, 0},
};
ApiServer::ApiServer(uint16_t port, Graph& graph, FlowManager& flow_manager, ProcessManager& process_manager, NodeControlClient& control_client, const std::string& mxl_domain)
: impl_(std::make_unique<Impl>()) {
impl_->data = std::make_unique<ApiServerImpl>();
impl_->data->port = port;
impl_->data->graph = &graph;
impl_->data->flow_manager = &flow_manager;
impl_->data->process_manager = &process_manager;
impl_->data->control_client = &control_client;
impl_->data->mxl_domain = mxl_domain;
g_impl = impl_->data.get();
struct lws_context_creation_info info;
std::memset(&info, 0, sizeof(info));
info.port = port;
info.protocols = protocols;
info.gid = -1;
info.uid = -1;
impl_->data->context = lws_create_context(&info);
if (!impl_->data->context) {
spdlog::error("Failed to create HTTP context on port {}", port);
throw std::runtime_error("Failed to create HTTP context");
}
spdlog::info("API server: listening on port {}", port);
}
ApiServer::~ApiServer() {
if (impl_->data && impl_->data->context) {
lws_context_destroy(impl_->data->context);
}
if (g_impl == impl_->data.get()) {
g_impl = nullptr;
}
}
void ApiServer::poll(int timeout_ms) {
lws_service(impl_->data->context, timeout_ms);
}
} // namespace dmf_engine
+64
View File
@@ -0,0 +1,64 @@
#include <dmf-engine/flow_manager.hpp>
#include <nlohmann/json.hpp>
#include <spdlog/spdlog.h>
#include <random>
#include <sstream>
namespace dmf_engine {
FlowManager::FlowManager() = default;
FlowId FlowManager::create_flow_id() {
std::random_device rd;
std::mt19937 gen(rd());
std::uniform_int_distribution<uint32_t> dist(0, 0xFFFFFFFF);
uint32_t a = dist(gen);
uint16_t b = dist(gen) & 0xFFFF;
uint16_t c = (dist(gen) & 0x0FFF) | 0x4000;
uint16_t d = (dist(gen) & 0x3FFF) | 0x8000;
uint32_t e1 = dist(gen);
uint16_t e2 = dist(gen) & 0xFFFF;
std::stringstream ss;
ss << std::hex << std::setfill('0');
ss << std::setw(8) << a << "-";
ss << std::setw(4) << b << "-";
ss << std::setw(4) << c << "-";
ss << std::setw(4) << d << "-";
ss << std::setw(8) << e1 << std::setw(4) << e2;
auto id = ss.str();
spdlog::info("FlowManager: created flow ID: {}", id);
return id;
}
nlohmann::json FlowManager::create_v210_flow_def(const FlowId& flow_id, int width, int height, int fps_numerator, int fps_denominator) const {
auto grain_size = (width * 2) * height;
nlohmann::json flow_def;
flow_def["id"] = flow_id;
flow_def["label"] = "DMF Studio Flow " + flow_id;
flow_def["description"] = "Auto-generated V210 flow";
flow_def["format"] = "urn:x-nmos:format:video";
flow_def["grain_rate"] = {{"numerator", fps_numerator}, {"denominator", fps_denominator}};
flow_def["media_type"] = "video/v210";
flow_def["frame_width"] = width;
flow_def["frame_height"] = height;
flow_def["interlace_mode"] = "progressive";
flow_def["colorspace"] = "BT709";
flow_def["components"] = nlohmann::json::array({
{{"name", "Y"}, {"width", width}, {"height", height}, {"bit_depth", 10}},
{{"name", "Cb"}, {"width", width / 2}, {"height", height}, {"bit_depth", 10}},
{{"name", "Cr"}, {"width", width / 2}, {"height", height}, {"bit_depth", 10}},
});
flow_def["tags"] = nlohmann::json::object();
flow_def["tags"]["urn:x-nmos:tag:grouphint/v1.0"] = nlohmann::json::array({"DMF Studio:Video"});
return flow_def;
}
} // namespace dmf_engine
+145
View File
@@ -0,0 +1,145 @@
#include <dmf-engine/graph.hpp>
#include <nlohmann/json.hpp>
#include <spdlog/spdlog.h>
namespace dmf_engine {
NodeId Graph::add_node(const std::string& type, const nlohmann::json& config) {
std::string id;
if (config.contains("id") && config["id"].is_string()) {
id = config["id"].get<std::string>();
} else {
id = type + "_" + std::to_string(next_node_num_++);
}
GraphNode node;
node.id = id;
node.type = type;
node.config = config;
node.state = NodeState::Stopped;
nodes_[id] = std::move(node);
spdlog::info("Graph: added node '{}' type='{}'", id, type);
return id;
}
bool Graph::remove_node(const NodeId& node_id) {
auto it = nodes_.find(node_id);
if (it == nodes_.end()) {
spdlog::warn("Graph: node '{}' not found", node_id);
return false;
}
std::vector<EdgeId> edges_to_remove;
for (const auto& [eid, edge] : edges_) {
if (edge.from_node == node_id || edge.to_node == node_id) {
edges_to_remove.push_back(eid);
}
}
for (const auto& eid : edges_to_remove) {
edges_.erase(eid);
spdlog::info("Graph: removed edge '{}' (connected to removed node '{}')", eid, node_id);
}
nodes_.erase(it);
spdlog::info("Graph: removed node '{}'", node_id);
return true;
}
const GraphNode* Graph::get_node(const NodeId& node_id) const {
auto it = nodes_.find(node_id);
return it != nodes_.end() ? &it->second : nullptr;
}
GraphNode* Graph::get_node_mut(const NodeId& node_id) {
auto it = nodes_.find(node_id);
return it != nodes_.end() ? &it->second : nullptr;
}
std::vector<GraphNode> Graph::get_nodes() const {
std::vector<GraphNode> result;
for (const auto& [_, node] : nodes_) {
result.push_back(node);
}
return result;
}
EdgeId Graph::add_edge(const NodeId& from_node, const PortId& from_port,
const NodeId& to_node, const PortId& to_port,
const FlowId& flow_id, const nlohmann::json& flow_def) {
auto id = from_node + ":" + from_port + "->" + to_node + ":" + to_port;
GraphEdge edge;
edge.id = id;
edge.from_node = from_node;
edge.from_port = from_port;
edge.to_node = to_node;
edge.to_port = to_port;
edge.flow_id = flow_id;
edge.flow_def = flow_def;
edges_[id] = std::move(edge);
spdlog::info("Graph: added edge '{}' flow_id={}", id, flow_id);
return id;
}
bool Graph::remove_edge(const EdgeId& edge_id) {
auto it = edges_.find(edge_id);
if (it == edges_.end()) {
spdlog::warn("Graph: edge '{}' not found", edge_id);
return false;
}
edges_.erase(it);
spdlog::info("Graph: removed edge '{}'", edge_id);
return true;
}
std::vector<GraphEdge> Graph::get_edges() const {
std::vector<GraphEdge> result;
for (const auto& [_, edge] : edges_) {
result.push_back(edge);
}
return result;
}
std::vector<GraphEdge> Graph::get_edges_for_node(const NodeId& node_id) const {
std::vector<GraphEdge> result;
for (const auto& [_, edge] : edges_) {
if (edge.from_node == node_id || edge.to_node == node_id) {
result.push_back(edge);
}
}
return result;
}
nlohmann::json Graph::serialize() const {
auto j = nlohmann::json::object();
auto nodes_arr = nlohmann::json::array();
for (const auto& [id, node] : nodes_) {
nodes_arr.push_back({
{"id", node.id},
{"type", node.type},
{"config", node.config},
{"state", static_cast<int>(node.state)},
{"control_port", node.control_port},
{"pid", node.pid},
});
}
j["nodes"] = nodes_arr;
auto edges_arr = nlohmann::json::array();
for (const auto& [id, edge] : edges_) {
edges_arr.push_back({
{"id", edge.id},
{"from_node", edge.from_node},
{"from_port", edge.from_port},
{"to_node", edge.to_node},
{"to_port", edge.to_port},
{"flow_id", edge.flow_id},
});
}
j["edges"] = edges_arr;
return j;
}
} // namespace dmf_engine
+144
View File
@@ -0,0 +1,144 @@
#include <dmf-engine/node_control_client.hpp>
#include <spdlog/spdlog.h>
#include <arpa/inet.h>
#include <cstring>
#include <netinet/in.h>
#include <sys/socket.h>
#include <unistd.h>
#include <sstream>
namespace dmf_engine {
void NodeControlClient::register_node(const std::string& node_id, uint16_t port) {
node_ports_[node_id] = port;
spdlog::info("NodeControlClient: registered node '{}' on port {}", node_id, port);
}
void NodeControlClient::unregister_node(const std::string& node_id) {
node_ports_.erase(node_id);
}
uint16_t NodeControlClient::get_port(const std::string& node_id) const {
auto it = node_ports_.find(node_id);
return it != node_ports_.end() ? it->second : 0;
}
bool NodeControlClient::send_command(uint16_t port, const std::string& json_cmd) {
int fd = socket(AF_INET, SOCK_STREAM, 0);
if (fd < 0) {
spdlog::error("NodeControlClient: socket() failed: {}", strerror(errno));
return false;
}
struct timeval tv;
tv.tv_sec = 2;
tv.tv_usec = 0;
setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
setsockopt(fd, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
struct sockaddr_in addr;
std::memset(&addr, 0, sizeof(addr));
addr.sin_family = AF_INET;
addr.sin_port = htons(port);
inet_pton(AF_INET, "127.0.0.1", &addr.sin_addr);
if (connect(fd, reinterpret_cast<struct sockaddr*>(&addr), sizeof(addr)) < 0) {
spdlog::error("NodeControlClient: connect to port {} failed: {}", port, strerror(errno));
close(fd);
return false;
}
std::ostringstream req;
req << "POST /cmd HTTP/1.1\r\n"
<< "Host: 127.0.0.1:" << port << "\r\n"
<< "Content-Type: application/json\r\n"
<< "Content-Length: " << json_cmd.size() << "\r\n"
<< "Connection: close\r\n"
<< "\r\n"
<< json_cmd;
auto request = req.str();
auto sent = write(fd, request.data(), request.size());
if (sent != static_cast<ssize_t>(request.size())) {
spdlog::error("NodeControlClient: write failed on port {}", port);
close(fd);
return false;
}
char resp_buf[4096] = {};
auto n = read(fd, resp_buf, sizeof(resp_buf) - 1);
close(fd);
if (n <= 0) {
spdlog::error("NodeControlClient: read failed on port {}", port);
return false;
}
std::string resp(resp_buf, n);
bool ok = resp.find("200 OK") != std::string::npos || resp.find("201 Created") != std::string::npos;
if (!ok) {
spdlog::warn("NodeControlClient: command failed on port {}: {}", port, resp.substr(0, 100));
}
return ok;
}
std::string NodeControlClient::send_command_with_response(uint16_t port, const std::string& json_cmd) {
int fd = socket(AF_INET, SOCK_STREAM, 0);
if (fd < 0) {
spdlog::error("NodeControlClient: socket() failed: {}", strerror(errno));
return "";
}
struct timeval tv;
tv.tv_sec = 2;
tv.tv_usec = 0;
setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
setsockopt(fd, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
struct sockaddr_in addr;
std::memset(&addr, 0, sizeof(addr));
addr.sin_family = AF_INET;
addr.sin_port = htons(port);
inet_pton(AF_INET, "127.0.0.1", &addr.sin_addr);
if (connect(fd, reinterpret_cast<struct sockaddr*>(&addr), sizeof(addr)) < 0) {
spdlog::error("NodeControlClient: connect to port {} failed: {}", port, strerror(errno));
close(fd);
return "";
}
std::ostringstream req;
req << "POST /cmd HTTP/1.1\r\n"
<< "Host: 127.0.0.1:" << port << "\r\n"
<< "Content-Type: application/json\r\n"
<< "Content-Length: " << json_cmd.size() << "\r\n"
<< "Connection: close\r\n"
<< "\r\n"
<< json_cmd;
auto request = req.str();
auto sent = write(fd, request.data(), request.size());
if (sent != static_cast<ssize_t>(request.size())) {
spdlog::error("NodeControlClient: write failed on port {}", port);
close(fd);
return "";
}
std::string full_resp;
char resp_buf[4096];
while (true) {
auto n = read(fd, resp_buf, sizeof(resp_buf));
if (n <= 0) break;
full_resp.append(resp_buf, n);
}
close(fd);
auto body_start = full_resp.find("\r\n\r\n");
if (body_start == std::string::npos) return "";
return full_resp.substr(body_start + 4);
}
} // namespace dmf_engine
+113
View File
@@ -0,0 +1,113 @@
#include <dmf-engine/process_manager.hpp>
#include <spdlog/spdlog.h>
#include <cstdlib>
#include <filesystem>
#include <signal.h>
#include <sys/types.h>
#include <sys/wait.h>
#include <unistd.h>
namespace dmf_engine {
std::string ProcessManager::find_node_binary(const std::string& node_type) const {
std::string binary_name = "dmf-node-" + node_type;
if (const char* env_path = std::getenv("DMF_NODE_PATH")) {
auto candidate = std::filesystem::path(env_path) / binary_name;
if (std::filesystem::exists(candidate)) {
return candidate.string();
}
}
if (const char* self_dir_env = std::getenv("DMF_STUDIO_BIN_DIR")) {
auto candidate = std::filesystem::path(self_dir_env) / binary_name;
if (std::filesystem::exists(candidate)) {
return candidate.string();
}
}
return binary_name;
}
bool ProcessManager::start_node(GraphNode& node, const std::string& mxl_domain, uint16_t base_port) {
if (node.state == NodeState::Running) {
spdlog::warn("ProcessManager: node '{}' is already running", node.id);
return true;
}
auto binary = find_node_binary(node.type);
node.control_port = base_port;
pid_t pid = fork();
if (pid < 0) {
spdlog::error("ProcessManager: fork() failed for node '{}': {}", node.id, strerror(errno));
return false;
}
if (pid == 0) {
std::string port_str = std::to_string(node.control_port);
std::string config_str = node.config.dump();
execlp(binary.c_str(), binary.c_str(),
"--node-id", node.id.c_str(),
"--control-port", port_str.c_str(),
"--mxl-domain", mxl_domain.c_str(),
"--config", config_str.c_str(),
nullptr);
spdlog::error("ProcessManager: execlp failed for '{}': {}", binary, strerror(errno));
_exit(1);
}
node.pid = pid;
node.state = NodeState::Running;
spdlog::info("ProcessManager: started node '{}' pid={} port={}", node.id, pid, node.control_port);
return true;
}
bool ProcessManager::stop_node(GraphNode& node) {
if (node.state != NodeState::Running || node.pid <= 0) {
spdlog::warn("ProcessManager: node '{}' is not running", node.id);
return true;
}
if (kill(node.pid, SIGTERM) != 0) {
spdlog::error("ProcessManager: failed to send SIGTERM to node '{}': {}", node.id, strerror(errno));
return false;
}
int status = 0;
waitpid(node.pid, &status, 0);
spdlog::info("ProcessManager: stopped node '{}' pid={}", node.id, node.pid);
node.pid = 0;
node.state = NodeState::Stopped;
return true;
}
void ProcessManager::stop_all() {
}
void ProcessManager::stop_all(Graph& graph) {
for (auto& node : graph.get_nodes()) {
if (node.state == NodeState::Running && node.pid > 0) {
kill(node.pid, SIGTERM);
int status = 0;
waitpid(node.pid, &status, 0);
auto* mut_node = graph.get_node_mut(node.id);
if (mut_node) {
mut_node->pid = 0;
mut_node->state = NodeState::Stopped;
}
spdlog::info("ProcessManager: stopped node '{}' pid={}", node.id, node.pid);
}
}
}
bool ProcessManager::is_running(const NodeId& node_id) const {
return false;
}
} // namespace dmf_engine
+17
View File
@@ -0,0 +1,17 @@
add_library(dmf-node STATIC
src/control_server.cpp
src/node_runner.cpp
)
target_include_directories(dmf-node PUBLIC
include
${CMAKE_SOURCE_DIR}/extern/mxl/lib/include
)
target_link_libraries(dmf-node PUBLIC
mxl
nlohmann_json::nlohmann_json
spdlog::spdlog
fmt::fmt
websockets
)
@@ -0,0 +1,35 @@
#pragma once
#include <dmf-node/node.hpp>
#include <functional>
#include <string>
#include <unordered_map>
#include <vector>
struct lws;
namespace dmf_node {
using CommandHandler = std::function<nlohmann::json(const nlohmann::json& payload)>;
using StatusCallback = std::function<void(const nlohmann::json& event)>;
class ControlServer {
public:
ControlServer(uint16_t port, StatusCallback on_event);
~ControlServer();
ControlServer(const ControlServer&) = delete;
ControlServer& operator=(const ControlServer&) = delete;
void register_command(const std::string& cmd, CommandHandler handler);
void send_event(const nlohmann::json& event);
void poll(int timeout_ms);
private:
struct Impl;
std::unique_ptr<Impl> impl_;
};
} // namespace dmf_node
+32
View File
@@ -0,0 +1,32 @@
#pragma once
#include <dmf-node/port.hpp>
#include <dmf-node/types.hpp>
#include <mxl/flow.h>
#include <mxl/mxl.h>
#include <nlohmann/json.hpp>
#include <memory>
#include <vector>
namespace dmf_node {
class Node {
public:
virtual ~Node() = default;
virtual std::string type() const = 0;
virtual std::vector<PortDef> input_ports() const = 0;
virtual std::vector<PortDef> output_ports() const = 0;
virtual void configure(const nlohmann::json& params) = 0;
virtual void on_add_writer(const std::string& port_id, mxlFlowWriter writer) = 0;
virtual void on_add_reader(const std::string& port_id, mxlFlowReader reader) = 0;
virtual void on_remove_writer(const std::string& port_id) = 0;
virtual void on_remove_reader(const std::string& port_id) = 0;
virtual void process() = 0;
virtual nlohmann::json status() const = 0;
};
} // namespace dmf_node
@@ -0,0 +1,39 @@
#pragma once
#include <dmf-node/node.hpp>
#include <dmf-node/control_server.hpp>
#include <mxl/mxl.h>
#include <atomic>
#include <memory>
#include <string>
namespace dmf_node {
class NodeRunner {
public:
template <typename NodeType>
static int run(int argc, char* argv[]) {
NodeRunner runner;
if (!runner.parse_args(argc, argv)) {
return 1;
}
auto node = std::make_unique<NodeType>();
return runner.exec(std::move(node));
}
private:
bool parse_args(int argc, char* argv[]);
int exec(std::unique_ptr<Node> node);
std::string node_id_;
uint16_t control_port_ = 0;
std::string mxl_domain_ = "/dev/shm/mxl";
std::string config_str_;
mxlInstance mxl_instance_ = nullptr;
std::atomic<bool> running_{false};
};
} // namespace dmf_node
+25
View File
@@ -0,0 +1,25 @@
#pragma once
#include <cstdint>
#include <string>
namespace dmf_node {
enum class MediaType : uint8_t {
VideoV210,
AudioFloat32,
AncData,
};
enum class PortDirection : uint8_t {
Input,
Output,
};
struct PortDef {
std::string id;
PortDirection direction;
MediaType media_type;
};
} // namespace dmf_node
+11
View File
@@ -0,0 +1,11 @@
#pragma once
#include <string>
namespace dmf_node {
using NodeId = std::string;
using PortId = std::string;
using FlowId = std::string;
} // namespace dmf_node
+265
View File
@@ -0,0 +1,265 @@
#include <dmf-node/control_server.hpp>
#include <dmf-node/node.hpp>
#include <libwebsockets.h>
#include <nlohmann/json.hpp>
#include <spdlog/spdlog.h>
#include <cstring>
#include <mutex>
#include <string>
#include <unordered_map>
#include <vector>
namespace dmf_node {
struct ControlServerData {
StatusCallback on_event;
std::unordered_map<std::string, CommandHandler> commands;
std::mutex send_mutex;
std::vector<std::string> send_queue;
struct lws* client_wsi = nullptr;
};
struct ControlServer::Impl {
uint16_t port;
std::unique_ptr<ControlServerData> data;
struct lws_context* context = nullptr;
};
static nlohmann::json dispatch_command(ControlServerData* data, const nlohmann::json& msg) {
if (!msg.contains("cmd")) {
spdlog::warn("Control: message missing 'cmd' field");
return {{"error", "missing 'cmd' field"}};
}
auto cmd = msg["cmd"].get<std::string>();
auto it = data->commands.find(cmd);
if (it != data->commands.end()) {
return it->second(msg);
} else {
spdlog::warn("Control: unknown command '{}'", cmd);
return {{"error", "unknown command: " + cmd}};
}
}
static int send_http_json(struct lws* wsi, const std::string& status, const std::string& body) {
auto hdr = "HTTP/1.1 " + status + "\r\n"
"Content-Type: application/json\r\n"
"Content-Length: " + std::to_string(body.size()) + "\r\n"
"Connection: close\r\n"
"\r\n";
std::vector<uint8_t> buf(LWS_PRE + hdr.size() + body.size());
std::memcpy(buf.data() + LWS_PRE, hdr.data(), hdr.size());
std::memcpy(buf.data() + LWS_PRE + hdr.size(), body.data(), body.size());
lws_write(wsi, buf.data() + LWS_PRE, hdr.size() + body.size(), LWS_WRITE_HTTP);
lws_http_transaction_completed(wsi);
return -1;
}
struct PerSession {
std::string http_body;
bool is_ws = false;
ControlServerData* data = nullptr;
};
static int callback_all(struct lws* wsi, enum lws_callback_reasons reason,
void* user, void* in, size_t len) {
auto* ps = static_cast<PerSession*>(user);
switch (reason) {
case LWS_CALLBACK_HTTP: {
new (ps) PerSession();
auto* vhost = lws_get_vhost(wsi);
ps->data = vhost ? static_cast<ControlServerData*>(lws_vhost_user(vhost)) : nullptr;
char* uri_ptr = nullptr;
int uri_len = 0;
int method = lws_http_get_uri_and_method(wsi, &uri_ptr, &uri_len);
std::string path(uri_ptr ? uri_ptr : "", uri_len > 0 ? uri_len : 0);
if (path == "/cmd" && method == LWSHUMETH_POST) {
int cl = lws_hdr_total_length(wsi, WSI_TOKEN_HTTP_CONTENT_LENGTH);
if (cl > 0) {
ps->http_body.reserve(cl);
}
return 0;
}
if (path == "/cmd" && (method == LWSHUMETH_GET || method == -1)) {
return send_http_json(wsi, "405 Method Not Allowed", R"({"error":"POST only"})");
}
return 0;
}
case LWS_CALLBACK_HTTP_BODY: {
ps->http_body.append(static_cast<char*>(in), len);
break;
}
case LWS_CALLBACK_HTTP_BODY_COMPLETION: {
if (!ps->data) {
return send_http_json(wsi, "500 Error", R"({"error":"no data"})");
}
try {
auto msg = nlohmann::json::parse(ps->http_body);
auto result = dispatch_command(ps->data, msg);
return send_http_json(wsi, "200 OK", result.dump());
} catch (const nlohmann::json::parse_error& e) {
return send_http_json(wsi, "400 Bad Request",
nlohmann::json({{"error", e.what()}}).dump());
}
}
case LWS_CALLBACK_ESTABLISHED: {
auto* vhost = lws_get_vhost(wsi);
ps->data = vhost ? static_cast<ControlServerData*>(lws_vhost_user(vhost)) : nullptr;
ps->is_ws = true;
if (ps->data) {
ps->data->client_wsi = wsi;
}
spdlog::info("Control WS: client connected");
break;
}
case LWS_CALLBACK_RECEIVE: {
if (!ps->data) {
break;
}
try {
auto msg = nlohmann::json::parse(static_cast<char*>(in), static_cast<char*>(in) + len);
dispatch_command(ps->data, msg);
} catch (const nlohmann::json::parse_error& e) {
spdlog::warn("Control WS: JSON parse error: {}", e.what());
}
break;
}
case LWS_CALLBACK_SERVER_WRITEABLE: {
if (!ps->data) {
break;
}
std::lock_guard<std::mutex> lock(ps->data->send_mutex);
while (!ps->data->send_queue.empty()) {
auto& msg = ps->data->send_queue.back();
std::vector<uint8_t> buf(LWS_PRE + msg.size());
std::memcpy(buf.data() + LWS_PRE, msg.data(), msg.size());
lws_write(wsi, buf.data() + LWS_PRE, msg.size(), LWS_WRITE_TEXT);
ps->data->send_queue.pop_back();
}
break;
}
case LWS_CALLBACK_CLOSED: {
if (ps->data) {
ps->data->client_wsi = nullptr;
}
spdlog::info("Control WS: client disconnected");
break;
}
default:
break;
}
return 0;
}
static const struct lws_protocols protocols[] = {
{
"http-only",
callback_all,
sizeof(PerSession),
0,
},
{
"dmf-control",
callback_all,
sizeof(PerSession),
65536,
},
{nullptr, nullptr, 0, 0},
};
static const struct lws_http_mount mounts[] = {
{
.mount_next = &mounts[1],
.mountpoint = "/cmd",
.origin = "",
.def = "",
.protocol = "http-only",
.cgienv = nullptr,
.extra_mimetypes = nullptr,
.interpret = nullptr,
.cgi_timeout = 0,
.cache_max_age = 0,
.auth_mask = 0,
.cache_reusable = 0,
.cache_revalidate = 0,
.cache_intermediaries = 0,
.origin_protocol = LWSMPRO_CALLBACK,
.mountpoint_len = 4,
},
{
.mount_next = nullptr,
.mountpoint = "/",
.origin = "",
.def = "",
.protocol = "dmf-control",
.cgienv = nullptr,
.extra_mimetypes = nullptr,
.interpret = nullptr,
.cgi_timeout = 0,
.cache_max_age = 0,
.auth_mask = 0,
.cache_reusable = 0,
.cache_revalidate = 0,
.cache_intermediaries = 0,
.origin_protocol = LWSMPRO_CALLBACK,
.mountpoint_len = 1,
},
};
ControlServer::ControlServer(uint16_t port, StatusCallback on_event)
: impl_(std::make_unique<Impl>()) {
impl_->port = port;
impl_->data = std::make_unique<ControlServerData>();
impl_->data->on_event = std::move(on_event);
struct lws_context_creation_info info;
std::memset(&info, 0, sizeof(info));
info.port = port;
info.protocols = protocols;
info.mounts = mounts;
info.user = impl_->data.get();
info.gid = -1;
info.uid = -1;
impl_->context = lws_create_context(&info);
if (!impl_->context) {
spdlog::error("Failed to create WS context on port {}", port);
throw std::runtime_error("Failed to create WS context");
}
spdlog::info("Control WS: listening on port {}", port);
}
ControlServer::~ControlServer() {
if (impl_->context) {
lws_context_destroy(impl_->context);
}
}
void ControlServer::register_command(const std::string& cmd, CommandHandler handler) {
impl_->data->commands[cmd] = std::move(handler);
}
void ControlServer::send_event(const nlohmann::json& event) {
auto data = event.dump();
{
std::lock_guard<std::mutex> lock(impl_->data->send_mutex);
impl_->data->send_queue.push_back(data);
}
if (impl_->data->client_wsi) {
lws_callback_on_writable(impl_->data->client_wsi);
}
}
void ControlServer::poll(int timeout_ms) {
lws_service(impl_->context, timeout_ms);
}
} // namespace dmf_node
+232
View File
@@ -0,0 +1,232 @@
#include <dmf-node/node_runner.hpp>
#include <dmf-node/control_server.hpp>
#include <mxl/flow.h>
#include <mxl/mxl.h>
#include <nlohmann/json.hpp>
#include <spdlog/spdlog.h>
#include <cstring>
#include <csignal>
#include <filesystem>
#include <thread>
namespace dmf_node {
static std::atomic<bool> g_running{true};
static void signal_handler(int /*signum*/) {
g_running = false;
}
bool NodeRunner::parse_args(int argc, char* argv[]) {
for (int i = 1; i < argc; ++i) {
std::string arg = argv[i];
if ((arg == "--node-id" || arg == "-n") && i + 1 < argc) {
node_id_ = argv[++i];
} else if ((arg == "--control-port" || arg == "-p") && i + 1 < argc) {
control_port_ = static_cast<uint16_t>(std::stoi(argv[++i]));
} else if ((arg == "--mxl-domain" || arg == "-d") && i + 1 < argc) {
mxl_domain_ = argv[++i];
} else if ((arg == "--config" || arg == "-c") && i + 1 < argc) {
config_str_ = argv[++i];
} else if (arg == "--help" || arg == "-h") {
spdlog::info("Usage: {} [options]", argv[0]);
spdlog::info(" --node-id, -n Node instance ID");
spdlog::info(" --control-port, -p WebSocket control port");
spdlog::info(" --mxl-domain, -d MXL domain path (default: /dev/shm/mxl)");
spdlog::info(" --config, -c Node configuration JSON");
return false;
}
}
if (node_id_.empty()) {
spdlog::error("--node-id is required");
return false;
}
if (control_port_ == 0) {
spdlog::error("--control-port is required");
return false;
}
return true;
}
int NodeRunner::exec(std::unique_ptr<Node> node) {
std::signal(SIGINT, signal_handler);
std::signal(SIGTERM, signal_handler);
spdlog::info("Starting node '{}' type='{}'", node_id_, node->type());
if (!std::filesystem::exists(mxl_domain_)) {
spdlog::error("MXL domain path does not exist: {}", mxl_domain_);
return 1;
}
mxl_instance_ = mxlCreateInstance(mxl_domain_.c_str(), nullptr);
if (!mxl_instance_) {
spdlog::error("Failed to create MXL instance on domain: {}", mxl_domain_);
return 1;
}
spdlog::info("MXL instance created on domain: {}", mxl_domain_);
mxlGarbageCollectFlows(mxl_instance_);
if (!config_str_.empty()) {
try {
auto cfg = nlohmann::json::parse(config_str_);
node->configure(cfg);
} catch (const nlohmann::json::parse_error& e) {
spdlog::error("Failed to parse config JSON: {}", e.what());
mxlDestroyInstance(mxl_instance_);
return 1;
}
}
auto control_server = std::make_unique<ControlServer>(control_port_, [](const nlohmann::json& /*event*/) {});
auto mk_ports = [](const std::vector<PortDef>& ports) {
auto arr = nlohmann::json::array();
for (const auto& p : ports) {
arr.push_back({{"id", p.id}, {"direction", p.direction == PortDirection::Input ? "input" : "output"}, {"media_type", static_cast<int>(p.media_type)}});
}
return arr;
};
struct FlowResource {
std::string port_id;
mxlFlowWriter writer = nullptr;
mxlFlowReader reader = nullptr;
};
std::vector<FlowResource> flow_resources;
control_server->register_command("add_writer", [&](const nlohmann::json& msg) -> nlohmann::json {
auto flow_id = msg["flow_id"].get<std::string>();
auto port_id = msg["port_id"].get<std::string>();
auto flow_def = msg["flow_def"].dump();
mxlFlowWriter writer = nullptr;
bool created = false;
mxlFlowConfigInfo config_info{};
auto status = mxlCreateFlowWriter(mxl_instance_, flow_def.c_str(), nullptr, &writer, &config_info, &created);
if (status != MXL_STATUS_OK || !writer) {
spdlog::error("Failed to create flow writer for flow {}: status={}", flow_id, static_cast<int>(status));
return {{"ok", false}, {"error", "failed to create flow writer"}};
}
spdlog::info("Created flow writer on port '{}' flow {} (created={})", port_id, flow_id, created);
flow_resources.push_back({port_id, writer, nullptr});
node->on_add_writer(port_id, writer);
return {{"ok", true}};
});
control_server->register_command("add_reader", [&](const nlohmann::json& msg) -> nlohmann::json {
auto flow_id = msg["flow_id"].get<std::string>();
auto port_id = msg["port_id"].get<std::string>();
mxlFlowReader reader = nullptr;
auto status = mxlCreateFlowReader(mxl_instance_, flow_id.c_str(), nullptr, &reader);
if (status != MXL_STATUS_OK || !reader) {
spdlog::error("Failed to create flow reader for flow {}: status={}", flow_id, static_cast<int>(status));
return {{"ok", false}, {"error", "failed to create flow reader"}};
}
spdlog::info("Created flow reader on port '{}' flow {}", port_id, flow_id);
flow_resources.push_back({port_id, nullptr, reader});
node->on_add_reader(port_id, reader);
return {{"ok", true}};
});
control_server->register_command("remove_writer", [&](const nlohmann::json& msg) -> nlohmann::json {
auto port_id = msg["port_id"].get<std::string>();
node->on_remove_writer(port_id);
for (auto it = flow_resources.begin(); it != flow_resources.end(); ++it) {
if (it->port_id == port_id && it->writer) {
mxlReleaseFlowWriter(mxl_instance_, it->writer);
flow_resources.erase(it);
break;
}
}
spdlog::info("Removed writer on port '{}'", port_id);
return {{"ok", true}};
});
control_server->register_command("remove_reader", [&](const nlohmann::json& msg) -> nlohmann::json {
auto port_id = msg["port_id"].get<std::string>();
node->on_remove_reader(port_id);
for (auto it = flow_resources.begin(); it != flow_resources.end(); ++it) {
if (it->port_id == port_id && it->reader) {
mxlReleaseFlowReader(mxl_instance_, it->reader);
flow_resources.erase(it);
break;
}
}
spdlog::info("Removed reader on port '{}'", port_id);
return {{"ok", true}};
});
control_server->register_command("configure", [&](const nlohmann::json& msg) -> nlohmann::json {
if (msg.contains("params")) {
node->configure(msg["params"]);
spdlog::info("Reconfigured node '{}'", node_id_);
}
return {{"ok", true}};
});
control_server->register_command("status", [&](const nlohmann::json& /*msg*/) -> nlohmann::json {
nlohmann::json resp;
resp["event"] = "status";
resp["node_id"] = node_id_;
resp["data"] = node->status();
control_server->send_event(resp);
return resp;
});
control_server->register_command("shutdown", [&](const nlohmann::json& /*msg*/) -> nlohmann::json {
spdlog::info("Shutdown command received");
g_running = false;
return {{"ok", true}};
});
nlohmann::json ready_event;
ready_event["event"] = "ready";
ready_event["type"] = node->type();
ready_event["node_id"] = node_id_;
ready_event["ports"] = mk_ports(node->input_ports());
for (const auto& p : node->output_ports()) {
ready_event["ports"].push_back({{"id", p.id}, {"direction", "output"}, {"media_type", static_cast<int>(p.media_type)}});
}
control_server->send_event(ready_event);
running_ = true;
spdlog::info("Node '{}' entering process loop", node_id_);
std::thread process_thread([&]() {
while (g_running && running_) {
node->process();
}
});
while (g_running && running_) {
control_server->poll(10);
}
running_ = false;
process_thread.join();
spdlog::info("Node '{}' shutting down", node_id_);
for (auto& res : flow_resources) {
if (res.writer) {
mxlReleaseFlowWriter(mxl_instance_, res.writer);
}
if (res.reader) {
mxlReleaseFlowReader(mxl_instance_, res.reader);
}
}
flow_resources.clear();
mxlDestroyInstance(mxl_instance_);
return 0;
}
} // namespace dmf_node
+21
View File
@@ -0,0 +1,21 @@
cmake_minimum_required(VERSION 3.24)
project(dmf-node-decklink-in LANGUAGES CXX)
set(DECKLINK_INCLUDE "${DECKLINK_SDK_DIR}/Linux/include")
add_executable(dmf-node-decklink-in
src/main.cpp
src/decklink_in_node.cpp
"${DECKLINK_SDK_DIR}/Linux/include/DeckLinkAPIDispatch.cpp"
)
target_include_directories(dmf-node-decklink-in PRIVATE
"${DECKLINK_INCLUDE}"
)
target_link_libraries(dmf-node-decklink-in PRIVATE
dmf-node
)
target_compile_options(dmf-node-decklink-in PRIVATE -Wno-unused-parameter)
+292
View File
@@ -0,0 +1,292 @@
#include "decklink_in_node.hpp"
#include <mxl/flow.h>
#include <mxl/mxl.h>
#include <mxl/time.h>
#include <DeckLinkAPIConfiguration.h>
#include <spdlog/spdlog.h>
#include <cstring>
namespace dmf_node {
DeckLinkInNode::~DeckLinkInNode() {
close_device();
}
void DeckLinkInNode::configure(const nlohmann::json& params) {
if (params.contains("device_index")) {
device_index_ = params["device_index"].get<int>();
}
if (params.contains("mode")) {
auto mode_str = params["mode"].get<std::string>();
if (mode_str == "1080i50") display_mode_ = bmdModeHD1080i50;
else if (mode_str == "1080p50") display_mode_ = bmdModeHD1080p50;
else if (mode_str == "1080p25") display_mode_ = bmdModeHD1080p25;
else if (mode_str == "1080i5994") display_mode_ = bmdModeHD1080i5994;
else if (mode_str == "1080p5994") display_mode_ = bmdModeHD1080p5994;
else if (mode_str == "1080p2997") display_mode_ = bmdModeHD1080p2997;
else if (mode_str == "720p50") display_mode_ = bmdModeHD720p50;
else if (mode_str == "720p5994") display_mode_ = bmdModeHD720p5994;
else {
spdlog::warn("DeckLink-in: unknown mode '{}', defaulting to 1080i50", mode_str);
}
}
if (params.contains("input_connection")) {
auto conn_str = params["input_connection"].get<std::string>();
if (conn_str == "sdi") input_connection_ = bmdVideoConnectionSDI;
else if (conn_str == "hdmi") input_connection_ = bmdVideoConnectionHDMI;
else if (conn_str == "optical_sdi") input_connection_ = bmdVideoConnectionOpticalSDI;
else if (conn_str == "component") input_connection_ = bmdVideoConnectionComponent;
else if (conn_str == "composite") input_connection_ = bmdVideoConnectionComposite;
else if (conn_str == "svideo") input_connection_ = bmdVideoConnectionSVideo;
else {
spdlog::warn("DeckLink-in: unknown input_connection '{}', defaulting to auto", conn_str);
}
}
if (!open_device()) {
spdlog::error("DeckLink-in: failed to open device during configure");
}
}
void DeckLinkInNode::on_add_writer(const std::string& port_id, mxlFlowWriter writer) {
if (port_id == "video_out") {
writer_ = writer;
auto now = mxlGetTime();
write_index_ = mxlTimestampToIndex(&grain_rate_, now);
spdlog::info("DeckLink-in: writer added, grain_rate={}/{}", grain_rate_.numerator, grain_rate_.denominator);
if (input_ && !capturing_) {
if (input_->StartStreams() != S_OK) {
spdlog::error("DeckLink-in: failed to start streams");
return;
}
capturing_ = true;
spdlog::info("DeckLink-in: capture started");
}
}
}
void DeckLinkInNode::on_remove_writer(const std::string& port_id) {
if (port_id == "video_out") {
close_device();
writer_.reset();
spdlog::info("DeckLink-in: writer removed");
}
}
bool DeckLinkInNode::open_device() {
auto* iter = CreateDeckLinkIteratorInstance();
if (!iter) {
spdlog::error("DeckLink-in: DeckLink drivers not found");
return false;
}
IDeckLink* device = nullptr;
for (int i = 0; i <= device_index_; ++i) {
if (iter->Next(&device) != S_OK) {
spdlog::error("DeckLink-in: device index {} not found", device_index_);
iter->Release();
return false;
}
if (i < device_index_) {
device->Release();
}
}
iter->Release();
const char* model_name = nullptr;
device->GetModelName(&model_name);
spdlog::info("DeckLink-in: opened device '{}'", model_name ? model_name : "unknown");
if (input_connection_ != bmdVideoConnectionUnspecified) {
IDeckLinkConfiguration* config = nullptr;
if (device->QueryInterface(IID_IDeckLinkConfiguration, (void**)&config) == S_OK && config) {
if (config->SetInt(bmdDeckLinkConfigVideoInputConnection, input_connection_) != S_OK) {
spdlog::warn("DeckLink-in: failed to set input connection");
} else {
spdlog::info("DeckLink-in: input connection configured");
}
config->Release();
}
}
if (device->QueryInterface(IID_IDeckLinkInput, (void**)&input_) != S_OK) {
spdlog::error("DeckLink-in: device has no input interface");
device->Release();
return false;
}
decklink_ = device;
callback_ = std::make_unique<CaptureCallback>(*this);
input_->SetCallback(callback_.get());
auto flags = bmdVideoInputEnableFormatDetection;
if (input_->EnableVideoInput(display_mode_, bmdFormat10BitYUV, flags) != S_OK) {
spdlog::error("DeckLink-in: failed to enable video input");
return false;
}
IDeckLinkDisplayMode* mode = nullptr;
if (input_->GetDisplayMode(display_mode_, &mode) == S_OK) {
frame_width_ = mode->GetWidth();
frame_height_ = mode->GetHeight();
BMDTimeValue duration = 0;
BMDTimeScale scale = 0;
mode->GetFrameRate(&duration, &scale);
mode->Release();
is_interlaced_ = (display_mode_ == bmdModeHD1080i50 ||
display_mode_ == bmdModeHD1080i5994);
grain_rate_.numerator = static_cast<int32_t>(scale);
grain_rate_.denominator = static_cast<int32_t>(duration);
auto g = std::__gcd(grain_rate_.numerator, grain_rate_.denominator);
grain_rate_.numerator /= g;
grain_rate_.denominator /= g;
spdlog::info("DeckLink-in: {}x{} @ {}/{} fps, interlaced={}",
frame_width_, frame_height_,
grain_rate_.numerator, grain_rate_.denominator,
is_interlaced_);
}
spdlog::info("DeckLink-in: device opened, waiting for writer to start capture");
return true;
}
void DeckLinkInNode::close_device() {
if (input_) {
if (capturing_) {
input_->StopStreams();
}
input_->DisableVideoInput();
input_->Release();
input_ = nullptr;
}
if (decklink_) {
decklink_->Release();
decklink_ = nullptr;
}
capturing_ = false;
callback_.reset();
spdlog::info("DeckLink-in: device closed");
}
void DeckLinkInNode::process() {
if (!writer_ || !capturing_) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
return;
}
auto deadline = mxlIndexToTimestamp(&grain_rate_, write_index_ + 1);
mxlSleepUntil(deadline);
void* src_data = nullptr;
long src_row_bytes = 0;
long src_width = 0;
long src_height = 0;
{
std::lock_guard<std::mutex> lock(frame_mutex_);
if (frame_data_) {
src_data = frame_data_;
src_row_bytes = frame_row_bytes_;
src_width = frame_width_;
src_height = frame_height_;
}
frame_ready_ = false;
}
if (!src_data) {
write_index_++;
return;
}
mxlGrainInfo out_grain{};
uint8_t* out_payload = nullptr;
auto status = mxlFlowWriterOpenGrain(*writer_, write_index_, &out_grain, &out_payload);
if (status != MXL_STATUS_OK) {
auto now = mxlGetTime();
auto current = mxlTimestampToIndex(&grain_rate_, now);
write_index_ = current + 1;
return;
}
auto dst_row_bytes = out_grain.grainSize / src_height;
auto copy_row_bytes = std::min(static_cast<long>(dst_row_bytes), src_row_bytes);
for (long y = 0; y < src_height && y < static_cast<long>(out_grain.grainSize / dst_row_bytes); ++y) {
std::memcpy(out_payload + y * dst_row_bytes,
static_cast<uint8_t*>(src_data) + y * src_row_bytes,
copy_row_bytes);
}
out_grain.validSlices = out_grain.totalSlices;
mxlFlowWriterCommitGrain(*writer_, &out_grain);
write_index_++;
grains_written_++;
if (grains_written_ == 1) {
spdlog::info("DeckLink-in: first grain written, index={}", write_index_ - 1);
}
}
nlohmann::json DeckLinkInNode::status() const {
return {
{"type", "decklink-in"},
{"grains_written", grains_written_},
{"write_index", write_index_},
{"capturing", capturing_.load()},
{"has_writer", writer_.has_value()},
{"grain_rate", {{"numerator", grain_rate_.numerator}, {"denominator", grain_rate_.denominator}}},
{"width", frame_width_},
{"height", frame_height_},
{"interlaced", is_interlaced_},
};
}
HRESULT DeckLinkInNode::CaptureCallback::VideoInputFormatChanged(
BMDVideoInputFormatChangedEvents, IDeckLinkDisplayMode* newMode, BMDDetectedVideoInputFormatFlags) {
spdlog::info("DeckLink-in: input format changed");
return S_OK;
}
HRESULT DeckLinkInNode::CaptureCallback::VideoInputFrameArrived(
IDeckLinkVideoInputFrame* videoFrame, IDeckLinkAudioInputPacket*) {
if (!videoFrame || (videoFrame->GetFlags() & bmdFrameHasNoInputSource)) {
return S_OK;
}
void* bytes = nullptr;
IDeckLinkVideoBuffer* buf = nullptr;
if (videoFrame->QueryInterface(IID_IDeckLinkVideoBuffer, (void**)&buf) == S_OK && buf) {
buf->StartAccess(bmdBufferAccessRead);
buf->GetBytes(&bytes);
buf->EndAccess(bmdBufferAccessRead);
buf->Release();
}
if (!bytes) {
return S_OK;
}
std::lock_guard<std::mutex> lock(owner_.frame_mutex_);
owner_.frame_data_ = bytes;
owner_.frame_row_bytes_ = videoFrame->GetRowBytes();
owner_.frame_width_ = videoFrame->GetWidth();
owner_.frame_height_ = videoFrame->GetHeight();
owner_.frame_ready_ = true;
return S_OK;
}
} // namespace dmf_node
@@ -0,0 +1,84 @@
#pragma once
#include <dmf-node/node.hpp>
#include <mxl/flow.h>
#include <mxl/time.h>
#include <DeckLinkAPI.h>
#include <atomic>
#include <mutex>
#include <optional>
#include <thread>
namespace dmf_node {
class DeckLinkInNode : public Node {
public:
DeckLinkInNode() = default;
~DeckLinkInNode();
std::string type() const override { return "decklink-in"; }
std::vector<PortDef> input_ports() const override { return {}; }
std::vector<PortDef> output_ports() const override {
return {{"video_out", PortDirection::Output, MediaType::VideoV210}};
}
void configure(const nlohmann::json& params) override;
void on_add_writer(const std::string& port_id, mxlFlowWriter writer) override;
void on_add_reader(const std::string& port_id, mxlFlowReader reader) override {}
void on_remove_writer(const std::string& port_id) override;
void on_remove_reader(const std::string& port_id) override {}
void process() override;
nlohmann::json status() const override;
private:
bool open_device();
void close_device();
std::optional<mxlFlowWriter> writer_;
mxlRational grain_rate_{25, 1};
uint64_t write_index_ = 0;
uint64_t grains_written_ = 0;
int device_index_ = 0;
BMDDisplayMode display_mode_ = bmdModeHD1080i50;
BMDVideoConnection input_connection_ = bmdVideoConnectionUnspecified;
IDeckLink* decklink_ = nullptr;
IDeckLinkInput* input_ = nullptr;
std::atomic<bool> capturing_{false};
bool is_interlaced_ = false;
std::mutex frame_mutex_;
void* frame_data_ = nullptr;
long frame_row_bytes_ = 0;
long frame_width_ = 1920;
long frame_height_ = 1080;
bool frame_ready_ = false;
class CaptureCallback : public IDeckLinkInputCallback {
public:
CaptureCallback(DeckLinkInNode& owner) : owner_(owner) {}
HRESULT STDMETHODCALLTYPE QueryInterface(REFIID, void**) override { return E_NOINTERFACE; }
ULONG STDMETHODCALLTYPE AddRef() override { return 1; }
ULONG STDMETHODCALLTYPE Release() override { return 1; }
HRESULT STDMETHODCALLTYPE VideoInputFormatChanged(
BMDVideoInputFormatChangedEvents, IDeckLinkDisplayMode*, BMDDetectedVideoInputFormatFlags) override;
HRESULT STDMETHODCALLTYPE VideoInputFrameArrived(
IDeckLinkVideoInputFrame* videoFrame, IDeckLinkAudioInputPacket*) override;
private:
DeckLinkInNode& owner_;
};
std::unique_ptr<CaptureCallback> callback_;
};
} // namespace dmf_node
+6
View File
@@ -0,0 +1,6 @@
#include <dmf-node/node_runner.hpp>
#include "decklink_in_node.hpp"
int main(int argc, char* argv[]) {
return dmf_node::NodeRunner::run<dmf_node::DeckLinkInNode>(argc, argv);
}
+21
View File
@@ -0,0 +1,21 @@
cmake_minimum_required(VERSION 3.24)
project(dmf-node-decklink-out LANGUAGES CXX)
set(DECKLINK_INCLUDE "${DECKLINK_SDK_DIR}/Linux/include")
add_executable(dmf-node-decklink-out
src/main.cpp
src/decklink_out_node.cpp
"${DECKLINK_SDK_DIR}/Linux/include/DeckLinkAPIDispatch.cpp"
)
target_include_directories(dmf-node-decklink-out PRIVATE
"${DECKLINK_INCLUDE}"
)
target_link_libraries(dmf-node-decklink-out PRIVATE
dmf-node
)
target_compile_options(dmf-node-decklink-out PRIVATE -Wno-unused-parameter)
@@ -0,0 +1,365 @@
#include "decklink_out_node.hpp"
#include <mxl/flow.h>
#include <mxl/mxl.h>
#include <mxl/time.h>
#include <DeckLinkAPIConfiguration.h>
#include <spdlog/spdlog.h>
#include <cstring>
namespace dmf_node {
DeckLinkOutNode::~DeckLinkOutNode() {
close_device();
}
void DeckLinkOutNode::configure(const nlohmann::json& params) {
if (params.contains("device_index")) {
device_index_ = params["device_index"].get<int>();
}
if (params.contains("mode")) {
auto mode_str = params["mode"].get<std::string>();
if (mode_str == "1080i50") display_mode_ = bmdModeHD1080i50;
else if (mode_str == "1080p50") display_mode_ = bmdModeHD1080p50;
else if (mode_str == "1080p25") display_mode_ = bmdModeHD1080p25;
else if (mode_str == "1080i5994") display_mode_ = bmdModeHD1080i5994;
else if (mode_str == "1080p5994") display_mode_ = bmdModeHD1080p5994;
else if (mode_str == "1080p2997") display_mode_ = bmdModeHD1080p2997;
else if (mode_str == "720p50") display_mode_ = bmdModeHD720p50;
else if (mode_str == "720p5994") display_mode_ = bmdModeHD720p5994;
else {
spdlog::warn("DeckLink-out: unknown mode '{}', defaulting to 1080i50", mode_str);
}
}
if (params.contains("output_connection")) {
auto conn_str = params["output_connection"].get<std::string>();
if (conn_str == "sdi") output_connection_ = bmdVideoConnectionSDI;
else if (conn_str == "hdmi") output_connection_ = bmdVideoConnectionHDMI;
else if (conn_str == "optical_sdi") output_connection_ = bmdVideoConnectionOpticalSDI;
else if (conn_str == "component") output_connection_ = bmdVideoConnectionComponent;
else if (conn_str == "composite") output_connection_ = bmdVideoConnectionComposite;
else if (conn_str == "svideo") output_connection_ = bmdVideoConnectionSVideo;
else {
spdlog::warn("DeckLink-out: unknown output_connection '{}', defaulting to auto", conn_str);
}
}
}
void DeckLinkOutNode::on_add_reader(const std::string& port_id, mxlFlowReader reader) {
if (port_id == "video_in") {
reader_ = reader;
mxlFlowConfigInfo config{};
mxlFlowReaderGetConfigInfo(*reader_, &config);
flow_rate_ = config.common.grainRate;
auto now = mxlGetTime();
auto current_index = mxlTimestampToIndex(&flow_rate_, now);
read_index_ = current_index - 2;
spdlog::info("DeckLink-out: reader added, flow_rate={}/{}", flow_rate_.numerator, flow_rate_.denominator);
if (!open_device()) {
spdlog::error("DeckLink-out: failed to open device");
return;
}
}
}
void DeckLinkOutNode::on_remove_reader(const std::string& port_id) {
if (port_id == "video_in") {
close_device();
reader_.reset();
spdlog::info("DeckLink-out: reader removed");
}
}
bool DeckLinkOutNode::open_device() {
auto* iter = CreateDeckLinkIteratorInstance();
if (!iter) {
spdlog::error("DeckLink-out: DeckLink drivers not found");
return false;
}
IDeckLink* device = nullptr;
for (int i = 0; i <= device_index_; ++i) {
if (iter->Next(&device) != S_OK) {
spdlog::error("DeckLink-out: device index {} not found", device_index_);
iter->Release();
return false;
}
if (i < device_index_) {
device->Release();
}
}
iter->Release();
const char* model_name = nullptr;
device->GetModelName(&model_name);
spdlog::info("DeckLink-out: opened device '{}'", model_name ? model_name : "unknown");
if (device->QueryInterface(IID_IDeckLinkOutput, (void**)&output_) != S_OK) {
spdlog::error("DeckLink-out: device has no output interface");
device->Release();
return false;
}
decklink_ = device;
if (output_connection_ != bmdVideoConnectionUnspecified) {
IDeckLinkConfiguration* config = nullptr;
if (decklink_->QueryInterface(IID_IDeckLinkConfiguration, (void**)&config) == S_OK && config) {
if (config->SetInt(bmdDeckLinkConfigVideoOutputConnection, output_connection_) != S_OK) {
spdlog::warn("DeckLink-out: failed to set output connection");
} else {
spdlog::info("DeckLink-out: set output connection configured");
}
config->Release();
} else {
spdlog::warn("DeckLink-out: IDeckLinkConfiguration not available");
}
}
callback_ = std::make_unique<OutputCallback>(*this);
output_->SetScheduledFrameCompletionCallback(callback_.get());
if (output_->EnableVideoOutput(display_mode_, bmdVideoOutputFlagDefault) != S_OK) {
spdlog::error("DeckLink-out: failed to enable video output");
return false;
}
IDeckLinkDisplayMode* mode = nullptr;
if (output_->GetDisplayMode(display_mode_, &mode) == S_OK) {
out_width_ = mode->GetWidth();
out_height_ = mode->GetHeight();
BMDTimeValue duration = 0;
BMDTimeScale scale = 0;
mode->GetFrameRate(&duration, &scale);
frame_duration_ = duration;
time_scale_ = scale;
mode->Release();
output_rate_.numerator = static_cast<int32_t>(scale);
output_rate_.denominator = static_cast<int32_t>(duration);
auto g = std::__gcd(output_rate_.numerator, output_rate_.denominator);
output_rate_.numerator /= g;
output_rate_.denominator /= g;
is_interlaced_ = (display_mode_ == bmdModeHD1080i50 ||
display_mode_ == bmdModeHD1080i5994);
spdlog::info("DeckLink-out: {}x{} @ {}/{} fps, interlaced={}",
out_width_, out_height_,
output_rate_.numerator, output_rate_.denominator,
is_interlaced_);
}
preroll_frames_ = 5;
for (int i = 0; i < preroll_frames_; ++i) {
IDeckLinkMutableVideoFrame* frame = nullptr;
int32_t out_row_bytes = out_width_ * 16 / 6;
if (output_->CreateVideoFrame(out_width_, out_height_, out_row_bytes,
bmdFormat10BitYUV, bmdFrameFlagDefault, &frame) != S_OK || !frame) {
out_row_bytes = out_width_ * 2;
if (output_->CreateVideoFrame(out_width_, out_height_, out_row_bytes,
bmdFormat8BitYUV, bmdFrameFlagDefault, &frame) != S_OK || !frame) {
spdlog::error("DeckLink-out: failed to create preroll frame");
return false;
}
spdlog::info("DeckLink-out: using 8BitYUV output format");
output_format_ = bmdFormat8BitYUV;
} else {
output_format_ = bmdFormat10BitYUV;
}
IDeckLinkVideoBuffer* buf = nullptr;
if (frame->QueryInterface(IID_IDeckLinkVideoBuffer, (void**)&buf) == S_OK && buf) {
buf->StartAccess(bmdBufferAccessWrite);
void* dst = nullptr;
buf->GetBytes(&dst);
if (dst) {
std::memset(dst, 0, out_row_bytes * out_height_);
}
buf->EndAccess(bmdBufferAccessWrite);
buf->Release();
}
auto stream_time = i * frame_duration_;
output_->ScheduleVideoFrame(frame, stream_time, frame_duration_, time_scale_);
frame->Release();
}
stream_frame_ = preroll_frames_;
if (output_->StartScheduledPlayback(0, time_scale_, 1.0) != S_OK) {
spdlog::error("DeckLink-out: failed to start scheduled playback");
return false;
}
playing_ = true;
spdlog::info("DeckLink-out: playback started");
return true;
}
void DeckLinkOutNode::close_device() {
if (output_) {
if (playing_) {
output_->StopScheduledPlayback(0, nullptr, time_scale_);
}
output_->DisableVideoOutput();
output_->Release();
output_ = nullptr;
}
if (decklink_) {
decklink_->Release();
decklink_ = nullptr;
}
playing_ = false;
callback_.reset();
{
std::lock_guard<std::mutex> lock(frame_pool_mutex_);
for (auto* f : frame_pool_) {
f->Release();
}
frame_pool_.clear();
}
spdlog::info("DeckLink-out: device closed");
}
void DeckLinkOutNode::schedule_frame(void* mxl_payload, long width, long height, long row_bytes) {
if (!output_) return;
int32_t out_row_bytes = 0;
if (output_format_ == bmdFormat10BitYUV) {
if (output_->RowBytesForPixelFormat(bmdFormat10BitYUV, width, &out_row_bytes) != S_OK || out_row_bytes <= 0) {
out_row_bytes = ((width + 5) / 6) * 16;
}
} else {
out_row_bytes = width * 2;
}
IDeckLinkMutableVideoFrame* frame = nullptr;
{
std::lock_guard<std::mutex> lock(frame_pool_mutex_);
if (!frame_pool_.empty()) {
frame = frame_pool_.back();
frame_pool_.pop_back();
}
}
if (!frame) {
return;
}
IDeckLinkVideoBuffer* buf = nullptr;
if (frame->QueryInterface(IID_IDeckLinkVideoBuffer, (void**)&buf) == S_OK && buf) {
buf->StartAccess(bmdBufferAccessWrite);
void* dst = nullptr;
buf->GetBytes(&dst);
if (dst) {
auto copy_row_bytes = std::min(static_cast<long>(out_row_bytes), row_bytes);
for (long y = 0; y < height; ++y) {
std::memcpy(static_cast<uint8_t*>(dst) + y * out_row_bytes,
static_cast<uint8_t*>(mxl_payload) + y * row_bytes,
copy_row_bytes);
}
}
buf->EndAccess(bmdBufferAccessWrite);
buf->Release();
}
auto stream_time = stream_frame_ * frame_duration_;
output_->ScheduleVideoFrame(frame, stream_time, frame_duration_, time_scale_);
frame->Release();
stream_frame_++;
}
void DeckLinkOutNode::process() {
if (!reader_ || !playing_) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
return;
}
if (first_frame_) {
auto now = mxlGetTime();
grains_read_ = mxlTimestampToIndex(&output_rate_, now);
first_frame_ = false;
spdlog::info("DeckLink-out: starting at output index {}", grains_read_);
}
auto deadline = mxlIndexToTimestamp(&output_rate_, grains_read_ + 1);
mxlSleepUntil(deadline);
auto out_timestamp = mxlIndexToTimestamp(&output_rate_, grains_read_);
auto source_index = mxlTimestampToIndex(&flow_rate_, out_timestamp);
mxlGrainInfo grain_info{};
uint8_t* payload = nullptr;
auto status = mxlFlowReaderGetGrain(*reader_, source_index, 5000000ULL, &grain_info, &payload);
if (status != MXL_STATUS_OK) {
if (status == MXL_ERR_OUT_OF_RANGE_TOO_LATE) {
auto now = mxlGetTime();
auto current_index = mxlTimestampToIndex(&flow_rate_, now);
spdlog::warn("DeckLink-out: grain too late, skipping to index {}", current_index - 2);
source_index = current_index - 2;
status = mxlFlowReaderGetGrain(*reader_, source_index, 5000000ULL, &grain_info, &payload);
if (status != MXL_STATUS_OK) {
grains_read_++;
return;
}
} else if (status == MXL_ERR_OUT_OF_RANGE_TOO_EARLY) {
return;
} else {
return;
}
}
mxlFlowConfigInfo config{};
mxlFlowReaderGetConfigInfo(*reader_, &config);
auto grain_size = grain_info.grainSize;
long src_row_bytes = (config.discrete.sliceSizes[0] > 0) ? static_cast<long>(config.discrete.sliceSizes[0]) : (grain_size / out_height_);
long src_height = (src_row_bytes > 0) ? (grain_size / src_row_bytes) : out_height_;
long src_width = (src_row_bytes * 3) / 8;
schedule_frame(payload, src_width, src_height, src_row_bytes);
grains_read_++;
if (grains_read_ == 1) {
spdlog::info("DeckLink-out: first grain output, source_index={}", source_index);
}
}
nlohmann::json DeckLinkOutNode::status() const {
return {
{"type", "decklink-out"},
{"grains_read", grains_read_},
{"read_index", read_index_},
{"playing", playing_.load()},
{"has_reader", reader_.has_value()},
{"flow_rate", {{"numerator", flow_rate_.numerator}, {"denominator", flow_rate_.denominator}}},
{"output_rate", {{"numerator", output_rate_.numerator}, {"denominator", output_rate_.denominator}}},
{"width", out_width_},
{"height", out_height_},
};
}
HRESULT DeckLinkOutNode::OutputCallback::ScheduledFrameCompleted(
IDeckLinkVideoFrame* completedFrame, BMDOutputFrameCompletionResult result) {
if (completedFrame) {
completedFrame->AddRef();
std::lock_guard<std::mutex> lock(owner_.frame_pool_mutex_);
owner_.frame_pool_.push_back(static_cast<IDeckLinkMutableVideoFrame*>(completedFrame));
owner_.frame_pool_cv_.notify_one();
}
return S_OK;
}
HRESULT DeckLinkOutNode::OutputCallback::ScheduledPlaybackHasStopped() {
owner_.playing_ = false;
spdlog::info("DeckLink-out: scheduled playback stopped");
return S_OK;
}
} // namespace dmf_node
@@ -0,0 +1,93 @@
#pragma once
#include <dmf-node/node.hpp>
#include <mxl/flow.h>
#include <mxl/time.h>
#include <DeckLinkAPI.h>
#include <atomic>
#include <condition_variable>
#include <mutex>
#include <optional>
#include <thread>
#include <vector>
namespace dmf_node {
class DeckLinkOutNode : public Node {
public:
DeckLinkOutNode() = default;
~DeckLinkOutNode();
std::string type() const override { return "decklink-out"; }
std::vector<PortDef> input_ports() const override {
return {{"video_in", PortDirection::Input, MediaType::VideoV210}};
}
std::vector<PortDef> output_ports() const override { return {}; }
void configure(const nlohmann::json& params) override;
void on_add_writer(const std::string& port_id, mxlFlowWriter writer) override {}
void on_add_reader(const std::string& port_id, mxlFlowReader reader) override;
void on_remove_writer(const std::string& port_id) override {}
void on_remove_reader(const std::string& port_id) override;
void process() override;
nlohmann::json status() const override;
private:
bool open_device();
void close_device();
void schedule_frame(void* mxl_payload, long width, long height, long row_bytes);
std::optional<mxlFlowReader> reader_;
mxlRational flow_rate_{50, 1};
mxlRational output_rate_{50, 1};
uint64_t read_index_ = 0;
uint64_t grains_read_ = 0;
bool first_frame_ = true;
int device_index_ = 0;
BMDDisplayMode display_mode_ = bmdModeHD1080i50;
BMDVideoConnection output_connection_ = bmdVideoConnectionUnspecified;
IDeckLink* decklink_ = nullptr;
IDeckLinkOutput* output_ = nullptr;
BMDTimeValue frame_duration_ = 1000;
BMDTimeScale time_scale_ = 50000;
BMDPixelFormat output_format_ = bmdFormat10BitYUV;
int preroll_frames_ = 3;
std::atomic<bool> playing_{false};
int create_fail_count_ = 0;
uint64_t stream_frame_ = 0;
bool is_interlaced_ = false;
long out_width_ = 1920;
long out_height_ = 1080;
class OutputCallback : public IDeckLinkVideoOutputCallback {
public:
OutputCallback(DeckLinkOutNode& owner) : owner_(owner) {}
HRESULT STDMETHODCALLTYPE QueryInterface(REFIID, void**) override { return E_NOINTERFACE; }
ULONG STDMETHODCALLTYPE AddRef() override { return 1; }
ULONG STDMETHODCALLTYPE Release() override { return 1; }
HRESULT STDMETHODCALLTYPE ScheduledFrameCompleted(
IDeckLinkVideoFrame* completedFrame, BMDOutputFrameCompletionResult result) override;
HRESULT STDMETHODCALLTYPE ScheduledPlaybackHasStopped() override;
private:
DeckLinkOutNode& owner_;
};
std::unique_ptr<OutputCallback> callback_;
std::vector<IDeckLinkMutableVideoFrame*> frame_pool_;
std::mutex frame_pool_mutex_;
std::condition_variable frame_pool_cv_;
IDeckLinkMutableVideoFrame* scheduled_frame_ = nullptr;
};
} // namespace dmf_node
+6
View File
@@ -0,0 +1,6 @@
#include <dmf-node/node_runner.hpp>
#include "decklink_out_node.hpp"
int main(int argc, char* argv[]) {
return dmf_node::NodeRunner::run<dmf_node::DeckLinkOutNode>(argc, argv);
}
+6
View File
@@ -0,0 +1,6 @@
add_executable(dmf-node-passthrough
src/main.cpp
src/passthrough_node.cpp
)
target_link_libraries(dmf-node-passthrough PRIVATE dmf-node)
+6
View File
@@ -0,0 +1,6 @@
#include <dmf-node/node_runner.hpp>
#include "passthrough_node.hpp"
int main(int argc, char* argv[]) {
return dmf_node::NodeRunner::run<dmf_node::PassthroughNode>(argc, argv);
}
+112
View File
@@ -0,0 +1,112 @@
#include "passthrough_node.hpp"
#include <mxl/flow.h>
#include <mxl/mxl.h>
#include <mxl/time.h>
#include <spdlog/spdlog.h>
#include <cstring>
namespace dmf_node {
void PassthroughNode::on_add_writer(const std::string& port_id, mxlFlowWriter writer) {
if (port_id == "video_out") {
writer_ = writer;
spdlog::info("Passthrough: writer added on video_out");
}
}
void PassthroughNode::on_add_reader(const std::string& port_id, mxlFlowReader reader) {
if (port_id == "video_in") {
reader_ = reader;
mxlFlowConfigInfo config{};
mxlFlowReaderGetConfigInfo(*reader_, &config);
grain_rate_ = config.common.grainRate;
auto now = mxlGetTime();
auto current_index = mxlTimestampToIndex(&grain_rate_, now);
read_index_ = current_index - READ_DELAY_GRAINS;
aligned_ = true;
spdlog::info("Passthrough: reader added, grain_rate={}/{}, read_index={}, delay={} grains",
grain_rate_.numerator, grain_rate_.denominator, read_index_, READ_DELAY_GRAINS);
}
}
void PassthroughNode::on_remove_writer(const std::string& port_id) {
if (port_id == "video_out") {
writer_.reset();
spdlog::info("Passthrough: writer removed from video_out");
}
}
void PassthroughNode::on_remove_reader(const std::string& port_id) {
if (port_id == "video_in") {
reader_.reset();
spdlog::info("Passthrough: reader removed from video_in");
}
}
void PassthroughNode::process() {
if (!reader_ || !writer_ || !aligned_) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
return;
}
auto deadline = mxlIndexToTimestamp(&grain_rate_, read_index_ + 1);
mxlSleepUntil(deadline);
mxlGrainInfo grain_info{};
uint8_t* payload = nullptr;
auto status = mxlFlowReaderGetGrain(*reader_, read_index_, 5000000ULL, &grain_info, &payload);
if (status != MXL_STATUS_OK) {
if (status == MXL_ERR_OUT_OF_RANGE_TOO_LATE || status == MXL_ERR_OUT_OF_RANGE_TOO_EARLY) {
auto now = mxlGetTime();
auto current_index = mxlTimestampToIndex(&grain_rate_, now);
read_index_ = current_index - READ_DELAY_GRAINS;
if (grains_processed_ == 0) {
spdlog::warn("Passthrough: realigned to index {}", read_index_);
}
}
return;
}
mxlGrainInfo out_grain{};
uint8_t* out_payload = nullptr;
status = mxlFlowWriterOpenGrain(*writer_, grain_info.index, &out_grain, &out_payload);
if (status != MXL_STATUS_OK) {
spdlog::warn("Passthrough: failed to open output grain at index {}: {}",
grain_info.index, static_cast<int>(status));
read_index_++;
return;
}
auto copy_size = std::min(grain_info.grainSize, out_grain.grainSize);
std::memcpy(out_payload, payload, copy_size);
out_grain.validSlices = grain_info.validSlices;
out_grain.flags = grain_info.flags;
mxlFlowWriterCommitGrain(*writer_, &out_grain);
read_index_++;
grains_processed_++;
if (grains_processed_ == 1) {
spdlog::info("Passthrough: first grain processed, index={}", grain_info.index);
}
}
nlohmann::json PassthroughNode::status() const {
return {
{"type", "passthrough"},
{"grains_processed", grains_processed_},
{"read_index", read_index_},
{"has_reader", reader_.has_value()},
{"has_writer", writer_.has_value()},
};
}
} // namespace dmf_node
@@ -0,0 +1,47 @@
#pragma once
#include <dmf-node/node.hpp>
#include <mxl/flow.h>
#include <optional>
#include <thread>
namespace dmf_node {
class PassthroughNode : public Node {
public:
PassthroughNode() = default;
std::string type() const override { return "passthrough"; }
std::vector<PortDef> input_ports() const override {
return {{"video_in", PortDirection::Input, MediaType::VideoV210}};
}
std::vector<PortDef> output_ports() const override {
return {{"video_out", PortDirection::Output, MediaType::VideoV210}};
}
void configure(const nlohmann::json& /*params*/) override {}
void on_add_writer(const std::string& port_id, mxlFlowWriter writer) override;
void on_add_reader(const std::string& port_id, mxlFlowReader reader) override;
void on_remove_writer(const std::string& port_id) override;
void on_remove_reader(const std::string& port_id) override;
void process() override;
nlohmann::json status() const override;
private:
static constexpr int64_t READ_DELAY_GRAINS = 2;
std::optional<mxlFlowReader> reader_;
std::optional<mxlFlowWriter> writer_;
mxlRational grain_rate_{50, 1};
uint64_t read_index_ = 0;
uint64_t grains_processed_ = 0;
bool aligned_ = false;
};
} // namespace dmf_node
Executable
+51
View File
@@ -0,0 +1,51 @@
#!/bin/bash
set -e
ENGINE_PORT=${ENGINE_PORT:-9000}
MXL_DOMAIN=${MXL_DOMAIN:-/tmp/dmf-mxl}
FLOW_ID=${FLOW_ID:-a0000001-0000-0000-0000-000000000001}
BASE_URL="http://127.0.0.1:${ENGINE_PORT}"
log() { echo "=== $1 ==="; }
log "Adding passthrough node"
curl -s -X POST "${BASE_URL}/api/graph/nodes" \
-H "Content-Type: application/json" \
-d '{"type":"passthrough","id":"pass1"}' | python3 -m json.tool 2>/dev/null || echo ""
log "Starting graph"
curl -s -X POST "${BASE_URL}/api/graph/start" | python3 -m json.tool 2>/dev/null || echo ""
sleep 1
log "Connecting input to flow ${FLOW_ID}"
curl -s -X POST "${BASE_URL}/api/graph/nodes/pass1/connect-input" \
-H "Content-Type: application/json" \
-d "{\"port_id\":\"video_in\",\"flow_id\":\"${FLOW_ID}\"}" | python3 -m json.tool 2>/dev/null || echo ""
log "Connecting output (creates new MXL flow)"
OUTPUT=$(curl -s -X POST "${BASE_URL}/api/graph/nodes/pass1/connect-output" \
-H "Content-Type: application/json" \
-d '{"port_id":"video_out"}')
echo "$OUTPUT" | python3 -m json.tool 2>/dev/null || echo "$OUTPUT"
OUTPUT_FLOW_ID=$(echo "$OUTPUT" | python3 -c "import sys,json; print(json.load(sys.stdin)['flow_id'])" 2>/dev/null)
if [ -n "$OUTPUT_FLOW_ID" ]; then
log "Output flow ID: ${OUTPUT_FLOW_ID}"
log "To verify with mxl-gst-sink:"
echo " mxl-gst-sink -d ${MXL_DOMAIN} -v ${OUTPUT_FLOW_ID}"
fi
sleep 3
log "Checking MXL domain"
mxl-info --domain "${MXL_DOMAIN}" 2>/dev/null || true
log "Sending status command"
curl -s -X POST "${BASE_URL}/api/graph/nodes/pass1/command" \
-H "Content-Type: application/json" \
-d '{"cmd":"status"}' | python3 -m json.tool 2>/dev/null || echo ""
log "Graph state"
curl -s "${BASE_URL}/api/graph" | python3 -m json.tool 2>/dev/null || echo ""
+8
View File
@@ -0,0 +1,8 @@
add_executable(dmf-test-graph
test_graph.cpp
)
target_link_libraries(dmf-test-graph PRIVATE
dmf-engine
Catch2::Catch2WithMain
)
+89
View File
@@ -0,0 +1,89 @@
#include <dmf-engine/graph.hpp>
#include <dmf-engine/flow_manager.hpp>
#include <catch2/catch_test_macros.hpp>
TEST_CASE("Graph add and remove nodes", "[graph]") {
dmf_engine::Graph graph;
auto id1 = graph.add_node("passthrough");
auto id2 = graph.add_node("test-source");
REQUIRE(graph.get_node(id1) != nullptr);
REQUIRE(graph.get_node(id2) != nullptr);
REQUIRE(graph.get_node(id1)->type == "passthrough");
REQUIRE(graph.get_node(id2)->type == "test-source");
auto nodes = graph.get_nodes();
REQUIRE(nodes.size() == 2);
graph.remove_node(id1);
REQUIRE(graph.get_node(id1) == nullptr);
REQUIRE(graph.get_nodes().size() == 1);
}
TEST_CASE("Graph add and remove edges", "[graph]") {
dmf_engine::Graph graph;
auto n1 = graph.add_node("test-source");
auto n2 = graph.add_node("passthrough");
auto edge_id = graph.add_edge(n1, "video_out", n2, "video_in", "flow1", {});
REQUIRE(graph.get_edges().size() == 1);
graph.remove_edge(edge_id);
REQUIRE(graph.get_edges().size() == 0);
}
TEST_CASE("Graph remove node also removes edges", "[graph]") {
dmf_engine::Graph graph;
auto n1 = graph.add_node("test-source");
auto n2 = graph.add_node("passthrough");
graph.add_edge(n1, "video_out", n2, "video_in", "flow1", {});
graph.add_edge(n2, "video_out", n1, "video_in", "flow2", {});
REQUIRE(graph.get_edges().size() == 2);
graph.remove_node(n1);
REQUIRE(graph.get_edges().size() == 0);
}
TEST_CASE("FlowManager creates unique IDs", "[flow_manager]") {
dmf_engine::FlowManager fm1;
dmf_engine::FlowManager fm2;
auto id1 = fm1.create_flow_id();
auto id2 = fm1.create_flow_id();
auto id3 = fm2.create_flow_id();
REQUIRE(id1 != id2);
REQUIRE(id2 != id3);
}
TEST_CASE("FlowManager creates V210 flow definition", "[flow_manager]") {
dmf_engine::FlowManager fm;
auto flow_id = fm.create_flow_id();
auto def = fm.create_v210_flow_def(flow_id, 1920, 1080, 50, 1);
REQUIRE(def["id"] == flow_id);
REQUIRE(def["format"] == "urn:x-nmos:format:video");
REQUIRE(def["media_type"] == "video/v210");
REQUIRE(def["frame_width"] == 1920);
REQUIRE(def["frame_height"] == 1080);
}
TEST_CASE("Graph serialize", "[graph]") {
dmf_engine::Graph graph;
auto n1 = graph.add_node("test-source");
auto n2 = graph.add_node("passthrough");
graph.add_edge(n1, "video_out", n2, "video_in", "flow1", {});
auto j = graph.serialize();
REQUIRE(j.contains("nodes"));
REQUIRE(j.contains("edges"));
REQUIRE(j["nodes"].size() == 2);
REQUIRE(j["edges"].size() == 1);
}
+20
View File
@@ -0,0 +1,20 @@
{
"name": "dmf-studio",
"version": "0.1.0",
"dependencies": [
"libwebsockets",
"nlohmann-json",
"spdlog",
"fmt",
"catch2",
{
"name": "stduuid",
"version>=": "1.2.3",
"features": ["system-gen", "gsl-span"]
},
"picojson",
"cli11",
"ada-url"
],
"builtin-baseline": "432412eecac55981cf608cc12e5b4bd91768ec57"
}