Merge develop: Phase 1 complete — Framework + Passthrough
This commit is contained in:
+12
@@ -1 +1,13 @@
|
|||||||
ref_arch.pdf
|
ref_arch.pdf
|
||||||
|
build/
|
||||||
|
.cache/
|
||||||
|
CMakeUserPresets.json
|
||||||
|
compile_commands.json
|
||||||
|
.vcpkg/
|
||||||
|
vcpkg_installed/
|
||||||
|
node_modules/
|
||||||
|
web/dist/
|
||||||
|
*.o
|
||||||
|
*.a
|
||||||
|
*.so
|
||||||
|
*.d
|
||||||
|
|||||||
@@ -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()
|
||||||
@@ -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
|
||||||
|
```
|
||||||
@@ -0,0 +1,5 @@
|
|||||||
|
add_executable(dmf-studio-engine
|
||||||
|
src/main.cpp
|
||||||
|
)
|
||||||
|
|
||||||
|
target_link_libraries(dmf-studio-engine PRIVATE dmf-engine)
|
||||||
@@ -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;
|
||||||
|
}
|
||||||
+1
@@ -0,0 +1 @@
|
|||||||
|
/home/itten/DMF/mxl
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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)
|
||||||
@@ -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);
|
||||||
|
}
|
||||||
@@ -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
|
||||||
@@ -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 ""
|
||||||
@@ -0,0 +1,8 @@
|
|||||||
|
add_executable(dmf-test-graph
|
||||||
|
test_graph.cpp
|
||||||
|
)
|
||||||
|
|
||||||
|
target_link_libraries(dmf-test-graph PRIVATE
|
||||||
|
dmf-engine
|
||||||
|
Catch2::Catch2WithMain
|
||||||
|
)
|
||||||
@@ -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
@@ -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"
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user