13 Commits

Author SHA1 Message Date
Johanness e4f3498617 Merge develop: Phase 1 complete — Framework + Passthrough 2026-05-26 23:04:51 +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
34 changed files with 2311 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
+31
View File
@@ -0,0 +1,31 @@
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)
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)
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
```
+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 @@
/home/itten/DMF/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,21 @@
#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);
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
+427
View File
@@ -0,0 +1,427 @@
#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;
}
auto flow_id = fm.create_flow_id();
auto flow_def = fm.create_v210_flow_def(flow_id, 1920, 1080, 50, 1);
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
@@ -0,0 +1,88 @@
#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;
}
} // 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<void(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
+264
View File
@@ -0,0 +1,264 @@
#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 void dispatch_command(ControlServerData* data, const nlohmann::json& msg) {
if (!msg.contains("cmd")) {
spdlog::warn("Control: message missing 'cmd' field");
return;
}
auto cmd = msg["cmd"].get<std::string>();
auto it = data->commands.find(cmd);
if (it != data->commands.end()) {
it->second(msg);
} else {
spdlog::warn("Control: 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);
dispatch_command(ps->data, msg);
return send_http_json(wsi, "200 OK", R"({"ok":true})");
} 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
+225
View File
@@ -0,0 +1,225 @@
#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) {
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;
}
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);
});
control_server->register_command("add_reader", [&](const nlohmann::json& msg) {
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;
}
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);
});
control_server->register_command("remove_writer", [&](const nlohmann::json& msg) {
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);
});
control_server->register_command("remove_reader", [&](const nlohmann::json& msg) {
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);
});
control_server->register_command("configure", [&](const nlohmann::json& msg) {
if (msg.contains("params")) {
node->configure(msg["params"]);
spdlog::info("Reconfigured node '{}'", node_id_);
}
});
control_server->register_command("status", [&](const nlohmann::json& /*msg*/) {
nlohmann::json resp;
resp["event"] = "status";
resp["node_id"] = node_id_;
resp["data"] = node->status();
control_server->send_event(resp);
});
control_server->register_command("shutdown", [&](const nlohmann::json& /*msg*/) {
spdlog::info("Shutdown command received");
g_running = false;
});
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
+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"
}