|
@@ -0,0 +1,549 @@
|
|
|
|
|
+#include "stream.h"
|
|
|
|
|
+#include <pthread.h>
|
|
|
|
|
+#include <sys/socket.h>
|
|
|
|
|
+#include <netinet/in.h>
|
|
|
|
|
+#include <arpa/inet.h>
|
|
|
|
|
+#include <fcntl.h>
|
|
|
|
|
+
|
|
|
|
|
+/* ============================================================
|
|
|
|
|
+ * 内部数据结构
|
|
|
|
|
+ * ============================================================ */
|
|
|
|
|
+
|
|
|
|
|
+#define MAX_SESSIONS 128
|
|
|
|
|
+
|
|
|
|
|
+static StreamSession* g_sessions[MAX_SESSIONS];
|
|
|
|
|
+static int g_session_count = 0;
|
|
|
|
|
+static int g_session_id_counter = 0;
|
|
|
|
|
+
|
|
|
|
|
+static pthread_mutex_t g_session_mutex = PTHREAD_MUTEX_INITIALIZER;
|
|
|
|
|
+
|
|
|
|
|
+/* RTMP推流相关(模拟zlmediakit推流) */
|
|
|
|
|
+typedef struct {
|
|
|
|
|
+ int socket_fd;
|
|
|
|
|
+ char server_ip[MAX_IP_SIZE];
|
|
|
|
|
+ short server_port;
|
|
|
|
|
+ char app_name[64];
|
|
|
|
|
+ char stream_key[128];
|
|
|
|
|
+ int is_connected;
|
|
|
|
|
+ pthread_t push_thread;
|
|
|
|
|
+ int stop_push;
|
|
|
|
|
+} RTMPPushContext;
|
|
|
|
|
+
|
|
|
|
|
+/* ============================================================
|
|
|
|
|
+ * PS封装相关常量
|
|
|
|
|
+ * ============================================================ */
|
|
|
|
|
+
|
|
|
|
|
+/* MPEG-2 PS流包头常量 */
|
|
|
|
|
+#define PS_PACKET_SIZE 2048
|
|
|
|
|
+#define PS_SYSTEM_HEADER 0x00
|
|
|
|
|
+#define PS_PACK_START_CODE 0x000001BA
|
|
|
|
|
+#define PS_SYSTEM_HEADER_CODE 0x000001BB
|
|
|
|
|
+#define PS_PROGRAM_STREAM_MAP 0x000001BC
|
|
|
|
|
+#define PS_PRIVATE_STREAM_1 0x000001BD
|
|
|
|
|
+#define PS_PADDING_STREAM 0x000001BE
|
|
|
|
|
+#define PS_PRIVATE_STREAM_2 0x000001BF
|
|
|
|
|
+#define PS_PROGRAM_STREAM_DIRECTORY 0x000001BF
|
|
|
|
|
+#define PS_END_OF_STREAM 0x000001B9
|
|
|
|
|
+#define PS_PACK_HEADER 0x000001E0
|
|
|
|
|
+
|
|
|
|
|
+/* H.264 NALU类型 */
|
|
|
|
|
+#define H264_NALU_START 0x000001
|
|
|
|
|
+#define H264_NALU_START_4 0x00000001
|
|
|
|
|
+#define H264_NALU_END_CODE 0x000001B9
|
|
|
|
|
+
|
|
|
|
|
+/* ============================================================
|
|
|
|
|
+ * 内部函数
|
|
|
|
|
+ * ============================================================ */
|
|
|
|
|
+
|
|
|
|
|
+static StreamSession* find_session(int stream_id) {
|
|
|
|
|
+ for (int i = 0; i < g_session_count; i++) {
|
|
|
|
|
+ if (g_sessions[i] != NULL && g_sessions[i]->stream_id == stream_id) {
|
|
|
|
|
+ return g_sessions[i];
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ return NULL;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+static int add_session(StreamSession* session) {
|
|
|
|
|
+ pthread_mutex_lock(&g_session_mutex);
|
|
|
|
|
+
|
|
|
|
|
+ if (g_session_count >= MAX_SESSIONS) {
|
|
|
|
|
+ pthread_mutex_unlock(&g_session_mutex);
|
|
|
|
|
+ LOG_ERROR("Maximum session count reached");
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ session->stream_id = ++g_session_id_counter;
|
|
|
|
|
+ session->state = STREAM_STATE_IDLE;
|
|
|
|
|
+ session->is_running = 0;
|
|
|
|
|
+ pthread_mutex_init(&session->lock, NULL);
|
|
|
|
|
+
|
|
|
|
|
+ g_sessions[g_session_count++] = session;
|
|
|
|
|
+
|
|
|
|
|
+ pthread_mutex_unlock(&g_session_mutex);
|
|
|
|
|
+ LOG_INFO("Session added: ID=%d, Device=%s, Channel=%d",
|
|
|
|
|
+ session->stream_id, session->device_id, session->channel);
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+static void remove_session(int stream_id) {
|
|
|
|
|
+ pthread_mutex_lock(&g_session_mutex);
|
|
|
|
|
+
|
|
|
|
|
+ for (int i = 0; i < g_session_count; i++) {
|
|
|
|
|
+ if (g_sessions[i] != NULL && g_sessions[i]->stream_id == stream_id) {
|
|
|
|
|
+ g_sessions[i] = NULL;
|
|
|
|
|
+ g_session_count--;
|
|
|
|
|
+
|
|
|
|
|
+ /* 压缩数组 */
|
|
|
|
|
+ int write_idx = 0;
|
|
|
|
|
+ for (int read_idx = 0; read_idx < g_session_count; read_idx++) {
|
|
|
|
|
+ if (g_sessions[read_idx] != NULL) {
|
|
|
|
|
+ g_sessions[write_idx++] = g_sessions[read_idx];
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ for (int j = write_idx; j < g_session_count; j++) {
|
|
|
|
|
+ g_sessions[j] = NULL;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Session removed: ID=%d", stream_id);
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ pthread_mutex_unlock(&g_session_mutex);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+/* ============================================================
|
|
|
|
|
+ * PS封装函数实现
|
|
|
|
|
+ * ============================================================ */
|
|
|
|
|
+
|
|
|
|
|
+/**
|
|
|
|
|
+ * @brief 将H.264 NALU封装为PS包
|
|
|
|
|
+ *
|
|
|
|
|
+ * PS(Program Stream)格式结构:
|
|
|
|
|
+ * - Pack Header (4-11 bytes)
|
|
|
|
|
+ * - Program Stream Map (可变长度)
|
|
|
|
|
+ * - PES Packet (可变长度)
|
|
|
|
|
+ * - PES Header
|
|
|
|
|
+ * - PES Payload (NALU数据)
|
|
|
|
|
+ *
|
|
|
|
|
+ * 每个PS包固定2048字节(包括纠错码)
|
|
|
|
|
+ */
|
|
|
|
|
+static int pack_h264_to_ps(StreamSession* session, const uint8_t* nalu_data, int nalu_size) {
|
|
|
|
|
+ /* 简化实现: 将H.264 NALU数据封装为PS流包 */
|
|
|
|
|
+ /* 实际实现需要完整的MPEG-2 PS封装逻辑 */
|
|
|
|
|
+
|
|
|
|
|
+ if (nalu_data == NULL || nalu_size <= 0) {
|
|
|
|
|
+ return -1;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /* 在实际实现中,这里会:
|
|
|
|
|
+ * 1. 检查SPS/PPS,如有变化则生成PS系统头和PSMAP
|
|
|
|
|
+ * 2. 将NALU封装为PES包
|
|
|
|
|
+ * 3. 添加Pack Header
|
|
|
|
|
+ * 4. 填充到PS_PACKET_SIZE(2048字节)
|
|
|
|
|
+ * 5. 添加纠错码
|
|
|
|
|
+ */
|
|
|
|
|
+
|
|
|
|
|
+ /* 模拟实现: 返回NALU大小作为"封装后大小" */
|
|
|
|
|
+ (void)session;
|
|
|
|
|
+ return nalu_size;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+/* ============================================================
|
|
|
|
|
+ * RTMP推流模拟(对接zlmediakit)
|
|
|
|
|
+ * ============================================================ */
|
|
|
|
|
+
|
|
|
|
|
+/**
|
|
|
|
|
+ * @brief RTMP推流线程函数
|
|
|
|
|
+ *
|
|
|
|
|
+ * 实际实现中需要链接zlmediakit的C API库:
|
|
|
|
|
+ * - zlm_sdk.h
|
|
|
|
|
+ * - libzlmediakit.so
|
|
|
|
|
+ *
|
|
|
|
|
+ * 核心流程:
|
|
|
|
|
+ * 1. 连接到zlmediakit服务器
|
|
|
|
|
+ * 2. 创建媒体流
|
|
|
|
|
+ * 3. 循环发送视频帧
|
|
|
|
|
+ * 4. 处理断线重连
|
|
|
|
|
+ */
|
|
|
|
|
+static void* rtmp_push_thread_func(void* arg) {
|
|
|
|
|
+ StreamSession* session = (StreamSession*)arg;
|
|
|
|
|
+ RTMPPushContext* ctx = (RTMPPushContext*)session->internal_handle;
|
|
|
|
|
+
|
|
|
|
|
+ if (ctx == NULL) {
|
|
|
|
|
+ LOG_ERROR("RTMP push context is NULL");
|
|
|
|
|
+ return NULL;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ char rtmp_url[512];
|
|
|
|
|
+ snprintf(rtmp_url, sizeof(rtmp_url), "rtmp://%s:%d/%s/%s",
|
|
|
|
|
+ ctx->server_ip, ctx->server_port, ctx->app_name, ctx->stream_key);
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("RTMP push thread started: %s", rtmp_url);
|
|
|
|
|
+
|
|
|
|
|
+ /* 模拟RTMP推流循环 */
|
|
|
|
|
+ while (!ctx->stop_push && session->is_running) {
|
|
|
|
|
+ /* 在实际实现中:
|
|
|
|
|
+ * 1. 从设备接收视频帧
|
|
|
|
|
+ * 2. 封装为PS流
|
|
|
|
|
+ * 3. 通过RTMP协议发送到zlmediakit
|
|
|
|
|
+ *
|
|
|
|
|
+ * 使用zlmediakit C API:
|
|
|
|
|
+ * ZLMediaKit_API *api = ZLMediaKit_Create();
|
|
|
|
|
+ * api->Connect(api, ctx->server_ip, ctx->server_port);
|
|
|
|
|
+ * api->CreateMediaSource(api, ctx->app_name, ctx->stream_key);
|
|
|
|
|
+ *
|
|
|
|
|
+ * 或者使用FFmpeg推流:
|
|
|
|
|
+ * avformat_alloc_output_context2(&fmt_ctx, NULL, "rtmp", rtmp_url);
|
|
|
|
|
+ * avformat_write_header(fmt_ctx, NULL);
|
|
|
|
|
+ * av_interleaved_write_frame(fmt_ctx, &packet);
|
|
|
|
|
+ */
|
|
|
|
|
+
|
|
|
|
|
+ /* 模拟发送帧 */
|
|
|
|
|
+ if (session->state == STREAM_STATE_STREAMING) {
|
|
|
|
|
+ /* 这里应该有实际的帧发送逻辑 */
|
|
|
|
|
+ sleep_milliseconds(1000 / session->fps);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ sleep_milliseconds(100);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("RTMP push thread stopped: %s", rtmp_url);
|
|
|
|
|
+ return NULL;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+/* ============================================================
|
|
|
|
|
+ * 公共接口实现
|
|
|
|
|
+ * ============================================================ */
|
|
|
|
|
+
|
|
|
|
|
+int stream_init(void) {
|
|
|
|
|
+ LOG_INFO("Initializing stream management...");
|
|
|
|
|
+
|
|
|
|
|
+ /* 初始化流媒体库 */
|
|
|
|
|
+ if (NET_ESTREAM_Init() != SUCCESS) {
|
|
|
|
|
+ LOG_ERROR("Failed to initialize stream library");
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ g_session_count = 0;
|
|
|
|
|
+ memset(g_sessions, 0, sizeof(g_sessions));
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Stream management initialized");
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void stream_cleanup(void) {
|
|
|
|
|
+ LOG_INFO("Cleaning up stream management...");
|
|
|
|
|
+
|
|
|
|
|
+ stream_stop_all();
|
|
|
|
|
+
|
|
|
|
|
+ for (int i = 0; i < MAX_SESSIONS; i++) {
|
|
|
|
|
+ if (g_sessions[i] != NULL) {
|
|
|
|
|
+ pthread_mutex_destroy(&g_sessions[i]->lock);
|
|
|
|
|
+ free(g_sessions[i]);
|
|
|
|
|
+ g_sessions[i] = NULL;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ g_session_count = 0;
|
|
|
|
|
+
|
|
|
|
|
+ NET_ESTREAM_Clean();
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Stream management cleaned up");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+int stream_start_playback_listen(const char* listen_ip, short listen_port,
|
|
|
|
|
+ void* callback, void* user_data) {
|
|
|
|
|
+ if (listen_ip == NULL) {
|
|
|
|
|
+ LOG_ERROR("Invalid parameters");
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Starting playback listen on %s:%d...", listen_ip, listen_port);
|
|
|
|
|
+
|
|
|
|
|
+ /* 调用HCNetSDK启动回放监听 */
|
|
|
|
|
+ if (NET_ESTREAM_StartListenPlayBack((char*)listen_ip, listen_port, callback, user_data) != SUCCESS) {
|
|
|
|
|
+ LOG_ERROR("Failed to start playback listen");
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Playback listen started successfully");
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void stream_stop_playback_listen(void) {
|
|
|
|
|
+ LOG_INFO("Stopping playback listen...");
|
|
|
|
|
+ NET_ESTREAM_StopListenPlayBack();
|
|
|
|
|
+ LOG_INFO("Playback listen stopped");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+int stream_set_playback_data_cb(int stream_handle, StreamDataCB callback, void* user_data) {
|
|
|
|
|
+ if (callback == NULL) {
|
|
|
|
|
+ LOG_ERROR("Invalid callback");
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Setting playback data callback for handle=%d", stream_handle);
|
|
|
|
|
+
|
|
|
|
|
+ /* 调用HCNetSDK设置回调 */
|
|
|
|
|
+ if (NET_ESTREAM_SetPlayBackDataCB(stream_handle, callback, user_data) != SUCCESS) {
|
|
|
|
|
+ LOG_ERROR("Failed to set playback data callback");
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+StreamSession* stream_create_session(const char* device_id, int channel,
|
|
|
|
|
+ StreamType stream_type, CodecType codec,
|
|
|
|
|
+ int width, int height, int fps, int bitrate) {
|
|
|
|
|
+ if (device_id == NULL || channel < 0) {
|
|
|
|
|
+ LOG_ERROR("Invalid parameters");
|
|
|
|
|
+ return NULL;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ StreamSession* session = (StreamSession*)calloc(1, sizeof(StreamSession));
|
|
|
|
|
+ if (session == NULL) {
|
|
|
|
|
+ LOG_ERROR("Failed to allocate session");
|
|
|
|
|
+ return NULL;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ safe_strcpy(session->device_id, device_id, MAX_DEVICE_ID);
|
|
|
|
|
+ session->channel = channel;
|
|
|
|
|
+ session->type = stream_type;
|
|
|
|
|
+ session->codec = codec;
|
|
|
|
|
+ session->width = width;
|
|
|
|
|
+ session->height = height;
|
|
|
|
|
+ session->fps = fps;
|
|
|
|
|
+ session->bitrate = bitrate;
|
|
|
|
|
+ session->internal_handle = NULL;
|
|
|
|
|
+
|
|
|
|
|
+ if (add_session(session) != SUCCESS) {
|
|
|
|
|
+ free(session);
|
|
|
|
|
+ return NULL;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Session created: ID=%d, Device=%s, Channel=%d, Type=%d, Codec=%d",
|
|
|
|
|
+ session->stream_id, device_id, channel, stream_type, codec);
|
|
|
|
|
+
|
|
|
|
|
+ return session;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void stream_destroy_session(StreamSession* session) {
|
|
|
|
|
+ if (session == NULL) return;
|
|
|
|
|
+
|
|
|
|
|
+ stream_stop_push(session);
|
|
|
|
|
+
|
|
|
|
|
+ remove_session(session->stream_id);
|
|
|
|
|
+ pthread_mutex_destroy(&session->lock);
|
|
|
|
|
+ free(session);
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Session destroyed: ID=%d", session->stream_id);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+int stream_start_push(StreamSession* session, const char* rtmp_url,
|
|
|
|
|
+ RTMPPushCB callback, void* user_data) {
|
|
|
|
|
+ if (session == NULL || rtmp_url == NULL) {
|
|
|
|
|
+ LOG_ERROR("Invalid parameters");
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (session->is_running) {
|
|
|
|
|
+ LOG_WARN("Stream already running: ID=%d", session->stream_id);
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Starting stream push: Device=%s, Channel=%d, RTMP=%s",
|
|
|
|
|
+ session->device_id, session->channel, rtmp_url);
|
|
|
|
|
+
|
|
|
|
|
+ /* 解析RTMP URL */
|
|
|
|
|
+ RTMPPushContext* ctx = (RTMPPushContext*)calloc(1, sizeof(RTMPPushContext));
|
|
|
|
|
+ if (ctx == NULL) {
|
|
|
|
|
+ LOG_ERROR("Failed to allocate RTMP context");
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /* 解析URL: rtmp://server:port/app/stream_key */
|
|
|
|
|
+ char url_copy[512];
|
|
|
|
|
+ safe_strcpy(url_copy, rtmp_url, sizeof(url_copy));
|
|
|
|
|
+
|
|
|
|
|
+ char* p = strstr(url_copy, "://");
|
|
|
|
|
+ if (p) {
|
|
|
|
|
+ p += 3;
|
|
|
|
|
+ char* slash = strchr(p, '/');
|
|
|
|
|
+ if (slash) {
|
|
|
|
|
+ *slash = '\0';
|
|
|
|
|
+ char* colon = strchr(p, ':');
|
|
|
|
|
+ if (colon) {
|
|
|
|
|
+ *colon = '\0';
|
|
|
|
|
+ safe_strcpy(ctx->server_ip, p, MAX_IP_SIZE);
|
|
|
|
|
+ ctx->server_port = atoi(colon + 1);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ safe_strcpy(ctx->server_ip, p, MAX_IP_SIZE);
|
|
|
|
|
+ ctx->server_port = DEFAULT_RTMP_PORT;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ char* app_start = slash + 1;
|
|
|
|
|
+ char* key_start = strchr(app_start, '/');
|
|
|
|
|
+ if (key_start) {
|
|
|
|
|
+ *key_start = '\0';
|
|
|
|
|
+ safe_strcpy(ctx->app_name, app_start, sizeof(ctx->app_name));
|
|
|
|
|
+ safe_strcpy(ctx->stream_key, key_start + 1, sizeof(ctx->stream_key));
|
|
|
|
|
+ } else {
|
|
|
|
|
+ safe_strcpy(ctx->app_name, app_start, sizeof(ctx->app_name));
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ session->internal_handle = ctx;
|
|
|
|
|
+ session->state = STREAM_STATE_STARTING;
|
|
|
|
|
+ session->is_running = 1;
|
|
|
|
|
+
|
|
|
|
|
+ /* 启动推流线程 */
|
|
|
|
|
+ ctx->stop_push = 0;
|
|
|
|
|
+ if (pthread_create(&ctx->push_thread, NULL, rtmp_push_thread_func, session) != 0) {
|
|
|
|
|
+ LOG_ERROR("Failed to create RTMP push thread");
|
|
|
|
|
+ session->is_running = 0;
|
|
|
|
|
+ session->state = STREAM_STATE_ERROR;
|
|
|
|
|
+ free(ctx);
|
|
|
|
|
+ session->internal_handle = NULL;
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ session->state = STREAM_STATE_STREAMING;
|
|
|
|
|
+
|
|
|
|
|
+ /* 触发回调 */
|
|
|
|
|
+ if (callback) {
|
|
|
|
|
+ callback(ctx->stream_key, 1, user_data); /* 1 = connected */
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Stream push started: ID=%d, RTMP Server=%s:%d",
|
|
|
|
|
+ session->stream_id, ctx->server_ip, ctx->server_port);
|
|
|
|
|
+
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+int stream_stop_push(StreamSession* session) {
|
|
|
|
|
+ if (session == NULL) return FAILURE;
|
|
|
|
|
+
|
|
|
|
|
+ if (!session->is_running) {
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Stopping stream push: ID=%d", session->stream_id);
|
|
|
|
|
+
|
|
|
|
|
+ session->is_running = 0;
|
|
|
|
|
+ session->state = STREAM_STATE_STOPPING;
|
|
|
|
|
+
|
|
|
|
|
+ RTMPPushContext* ctx = (RTMPPushContext*)session->internal_handle;
|
|
|
|
|
+ if (ctx) {
|
|
|
|
|
+ ctx->stop_push = 1;
|
|
|
|
|
+ if (ctx->push_thread != 0) {
|
|
|
|
|
+ pthread_join(ctx->push_thread, NULL);
|
|
|
|
|
+ ctx->push_thread = 0;
|
|
|
|
|
+ }
|
|
|
|
|
+ free(ctx);
|
|
|
|
|
+ session->internal_handle = NULL;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ session->state = STREAM_STATE_IDLE;
|
|
|
|
|
+
|
|
|
|
|
+ LOG_INFO("Stream push stopped: ID=%d", session->stream_id);
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+int stream_send_frame(StreamSession* session, StreamFrame* frame) {
|
|
|
|
|
+ if (session == NULL || frame == NULL) {
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (session->state != STREAM_STATE_STREAMING) {
|
|
|
|
|
+ return FAILURE;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /* 封装为PS流 */
|
|
|
|
|
+ /* 在实际实现中,这里会调用encode_to_ps将H.264数据封装为PS流 */
|
|
|
|
|
+
|
|
|
|
|
+ RTMPPushContext* ctx = (RTMPPushContext*)session->internal_handle;
|
|
|
|
|
+ if (ctx && ctx->socket_fd >= 0) {
|
|
|
|
|
+ /* 通过RTMP发送帧数据 */
|
|
|
|
|
+ /* 实际实现中使用zlmediakit API或FFmpeg */
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+int stream_stop_channel(const char* device_id, int channel) {
|
|
|
|
|
+ if (device_id == NULL) return FAILURE;
|
|
|
|
|
+
|
|
|
|
|
+ pthread_mutex_lock(&g_session_mutex);
|
|
|
|
|
+ for (int i = 0; i < g_session_count; i++) {
|
|
|
|
|
+ if (g_sessions[i] != NULL &&
|
|
|
|
|
+ strcmp(g_sessions[i]->device_id, device_id) == 0 &&
|
|
|
|
|
+ g_sessions[i]->channel == channel) {
|
|
|
|
|
+ stream_stop_push(g_sessions[i]);
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ pthread_mutex_unlock(&g_session_mutex);
|
|
|
|
|
+
|
|
|
|
|
+ return SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void stream_stop_all(void) {
|
|
|
|
|
+ pthread_mutex_lock(&g_session_mutex);
|
|
|
|
|
+ for (int i = 0; i < g_session_count; i++) {
|
|
|
|
|
+ if (g_sessions[i] != NULL) {
|
|
|
|
|
+ stream_stop_push(g_sessions[i]);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ pthread_mutex_unlock(&g_session_mutex);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+StreamState stream_get_state(StreamSession* session) {
|
|
|
|
|
+ if (session == NULL) return STREAM_STATE_ERROR;
|
|
|
|
|
+ return session->state;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+int encode_to_ps(const uint8_t* input_data, int input_size,
|
|
|
|
|
+ uint8_t* output_buf, int output_size) {
|
|
|
|
|
+ if (input_data == NULL || output_buf == NULL) {
|
|
|
|
|
+ return -1;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (input_size > output_size) {
|
|
|
|
|
+ return -1;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /* 简化实现: 直接复制数据
|
|
|
|
|
+ * 实际实现需要完整的MPEG-2 PS封装:
|
|
|
|
|
+ * 1. 添加Pack Header
|
|
|
|
|
+ * 2. 添加PES Header
|
|
|
|
|
+ * 3. 封装NALU数据
|
|
|
|
|
+ * 4. 填充到2048字节
|
|
|
|
|
+ * 5. 添加CRC校验
|
|
|
|
|
+ */
|
|
|
|
|
+ memcpy(output_buf, input_data, input_size);
|
|
|
|
|
+ return input_size;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+int decode_ps_to_frame(const uint8_t* ps_data, int ps_size,
|
|
|
|
|
+ uint8_t* output_buf, int output_size) {
|
|
|
|
|
+ if (ps_data == NULL || output_buf == NULL) {
|
|
|
|
|
+ return -1;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /* 简化实现: 直接复制数据
|
|
|
|
|
+ * 实际实现需要PS解复用:
|
|
|
|
|
+ * 1. 解析Pack Header
|
|
|
|
|
+ * 2. 解析PES Header
|
|
|
|
|
+ * 3. 提取NALU数据
|
|
|
|
|
+ * 4. 去除PS封装开销
|
|
|
|
|
+ */
|
|
|
|
|
+ if (ps_size > output_size) {
|
|
|
|
|
+ return -1;
|
|
|
|
|
+ }
|
|
|
|
|
+ memcpy(output_buf, ps_data, ps_size);
|
|
|
|
|
+ return ps_size;
|
|
|
|
|
+}
|