Files
DMF-Studio/libs/dmf-engine/src/process_manager.cpp
T
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

114 lines
3.3 KiB
C++

#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