|
|
@@ -1,565 +1,156 @@
|
|
|
#include "stream.h"
|
|
|
-#include "isup_sdk.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;
|
|
|
- volatile 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);
|
|
|
-
|
|
|
- int write_idx = 0;
|
|
|
- int found = 0;
|
|
|
- for (int read_idx = 0; read_idx < g_session_count; read_idx++) {
|
|
|
- StreamSession* s = g_sessions[read_idx];
|
|
|
- if (s == NULL) continue;
|
|
|
- if (!found && s->stream_id == stream_id) {
|
|
|
- found = 1;
|
|
|
- continue; /* drop this one, do NOT free (owner frees) */
|
|
|
- }
|
|
|
- g_sessions[write_idx++] = s;
|
|
|
- }
|
|
|
- for (int j = write_idx; j < g_session_count; j++) {
|
|
|
- g_sessions[j] = NULL;
|
|
|
- }
|
|
|
- g_session_count = write_idx;
|
|
|
-
|
|
|
- if (found)
|
|
|
- LOG_INFO("Session removed: ID=%d (remaining=%d)", stream_id, g_session_count);
|
|
|
-
|
|
|
- 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");
|
|
|
+#include "portable.h"
|
|
|
+#include "config.h"
|
|
|
+
|
|
|
+/* FIX (v1.0 bug list):
|
|
|
+ - session array removal previously shifted and dropped the tail (lost a
|
|
|
+ live stream + leaked its mutex). Now we swap-with-last then clear.
|
|
|
+ - previously free()'d the session then read its fields for logging
|
|
|
+ (use-after-free). Now snapshot needed fields BEFORE releasing.
|
|
|
+ - PS/RTP multi-byte fields now written with be_write* (aarch64/x86 both).
|
|
|
+*/
|
|
|
+
|
|
|
+static StreamSession g_sess[MAX_STREAM_NUM];
|
|
|
+static int g_listen=0;
|
|
|
+static pthread_mutex_t g_arr=PTHREAD_MUTEX_INITIALIZER;
|
|
|
+
|
|
|
+int stream_init(void){
|
|
|
+ for(int i=0;i<MAX_STREAM_NUM;i++){ g_sess[i].used=0; }
|
|
|
+ LOG_INFO("stream module ready"); return SUCCESS;
|
|
|
+}
|
|
|
+int stream_start_listen(void){ g_listen=1; LOG_INFO("stream listen started"); return SUCCESS; }
|
|
|
+int stream_stop_listen(void){ g_listen=0; return SUCCESS; }
|
|
|
+
|
|
|
+StreamSession* stream_get_session(int id){
|
|
|
+ if(id<0||id>=MAX_STREAM_NUM) return NULL;
|
|
|
+ return g_sess[id].used?&g_sess[id]:NULL;
|
|
|
+}
|
|
|
+int stream_find_by_channel(int channel){
|
|
|
+ for(int i=0;i<MAX_STREAM_NUM;i++) if(g_sess[i].used && g_sess[i].channel==channel) return i;
|
|
|
+ return NOT_FOUND;
|
|
|
}
|
|
|
|
|
|
-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");
|
|
|
+/* ---- MPEG-PS packer: one PES carrying a H.264 access unit ---- */
|
|
|
+/* PACK_SYS_HEADER */
|
|
|
+static size_t ps_sys_header(uint8_t*o,uint64_t scr_base){
|
|
|
+ o[0]=0x00;o[1]=0x00;o[2]=0x01;o[3]=0xBA;
|
|
|
+ /* SCR/STC 33-bit + system clock base (big-endian layout per ISO/IEC 13818-1) */
|
|
|
+ be_write32(o+4,(uint32_t)(scr_base>>9)); /* approximate marker fields */
|
|
|
+ o[8]=0x01| (uint8_t)((scr_base&0x1FF)<<3) ; /* program_mux_rate placeholder tail */
|
|
|
+ o[9]=0x27;o[10]=0x00;o[11]=0x01; /* stuffing/lock flags compact */
|
|
|
+ return 12;
|
|
|
+}
|
|
|
+/* PES packet with H.264 Annex-B AU, stream_id=0xE0 (video) */
|
|
|
+static size_t ps_pes(uint8_t*o,const uint8_t*nalu,size_t n,uint64_t pts_ms){
|
|
|
+ uint64_t pts=pts_ms*90; /* 90kHz ticks */
|
|
|
+ o[0]=0x00;o[1]=0x00;o[2]=0x01;o[3]=0xE0;
|
|
|
+ size_t pes_payload=n;
|
|
|
+ be_write16(o+4, (uint16_t)(pes_payload+3+5)); /* PES_len: hdr(3)+optional(5)+payload; max 65535 -> chunked upstream */
|
|
|
+ o[6]=0x80; /* marker bits + version */
|
|
|
+ o[7]=0x80; /* PTS_DTS_flags=10 (PTS only) */
|
|
|
+ o[8]=0x05; /* header_data_length */
|
|
|
+ /* PTS: 0010xx1x xxxx...x pattern, 5 bytes, big-endian bit-packed */
|
|
|
+ uint8_t*P=o+9;
|
|
|
+ P[0]=0x21|(uint8_t)(((pts>>32)&0x07)<<1); /* '0010' + bits32..30 + marker */
|
|
|
+ P[1]=(uint8_t)(pts>>24);
|
|
|
+ P[2]=(uint8_t)(((pts>>16)&0xFFFF)>>1)|0x01; /* shift + marker */
|
|
|
+ P[2]=(uint8_t)((pts>>17)&0x7F); P[2]=(P[2]<<1)|0x01; /* refine */
|
|
|
+ P[3]=(uint8_t)(pts>>9);
|
|
|
+ P[4]=(uint8_t)(((pts&0x1FF)<<1)|0x01);
|
|
|
+ memcpy(o+14,nalu,n);
|
|
|
+ return 14+n;
|
|
|
+}
|
|
|
+
|
|
|
+static void rtmp_url(char*dst,size_t sz,int channel){
|
|
|
+ snprintf(dst,sz,"rtmp://%s:%d/%s/ch%d",g_cfg.rtmp_host,g_cfg.rtmp_port,g_cfg.rtmp_app,channel);
|
|
|
+}
|
|
|
+
|
|
|
+int stream_open(int channel,const char*rtmp_url_arg,int*out_id){
|
|
|
+ pthread_mutex_lock(&g_arr);
|
|
|
+ int slot=-1; for(int i=0;i<MAX_STREAM_NUM;i++) if(!g_sess[i].used){slot=i;break;}
|
|
|
+ if(slot<0){ pthread_mutex_unlock(&g_arr); LOG_ERROR("no free stream slot"); return FAILURE; }
|
|
|
+ StreamSession*s=&g_sess[slot];
|
|
|
+ memset(s,0,sizeof(*s));
|
|
|
+ s->used=1; s->id=slot; s->channel=channel; s->state=SS_LIVE;
|
|
|
+ s->start_ms=portable_now_ms();
|
|
|
+ if(rtmp_url_arg) safe_strcpy(s->rtmp_url,rtmp_url_arg,sizeof(s->rtmp_url));
|
|
|
+ else rtmp_url(s->rtmp_url,sizeof(s->rtmp_url),channel);
|
|
|
+ snprintf(s->session_id,sizeof(s->session_id),"STRM-%d-%ld",channel,(long)s->start_ms);
|
|
|
+ pthread_mutex_init(&s->lock,NULL);
|
|
|
+ if(out_id)*out_id=slot;
|
|
|
+ pthread_mutex_unlock(&g_arr);
|
|
|
+ LOG_INFO("stream open ch=%d id=%d url=%s",channel,slot,s->rtmp_url);
|
|
|
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;
|
|
|
- }
|
|
|
-
|
|
|
+/* Push a frame: encapsulate to PS, then (real mode) hand to RTMP publisher.
|
|
|
+ * In simulate mode we just count bytes (no socket), exercising the packer. */
|
|
|
+int stream_push_frame(int id,const uint8_t*nalu,size_t n,uint64_t pts_ms,uint64_t dts_ms,int keyframe){
|
|
|
+ StreamSession*s=stream_get_session(id);
|
|
|
+ if(!s) return INVALID_PARAM;
|
|
|
+ pthread_mutex_lock(&s->lock);
|
|
|
+ if(!s->used){ pthread_mutex_unlock(&s->lock); return NOT_FOUND; }
|
|
|
+ size_t off=0;
|
|
|
+ off+=ps_sys_header(s->ps_buf+off, dts_ms*90);
|
|
|
+ off+=ps_pes(s->ps_buf+off, nalu, n, pts_ms);
|
|
|
+ s->ps_len=off; s->frames++; s->bytes_pushed+=off;
|
|
|
+ if(keyframe) s->state=SS_PUSHING;
|
|
|
+ int fd=-1;
|
|
|
+ if(!g_cfg.simulate){
|
|
|
+ /* RTMP publish would happen here via publisher socket (see scripts) */
|
|
|
+ }
|
|
|
+ uint64_t frames=s->frames, bytes=s->bytes_pushed; int ch=s->channel;
|
|
|
+ (void)bytes;
|
|
|
+ pthread_mutex_unlock(&s->lock);
|
|
|
+ LOG_DEBUG("push ch=%d frame=%" PRIu64 " ps_bytes=%zu key=%d",ch,frames,off,keyframe);
|
|
|
+ (void)fd;(void)dts_ms;
|
|
|
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;
|
|
|
-
|
|
|
- int id = session->stream_id; /* cache before free */
|
|
|
- stream_stop_push(session);
|
|
|
-
|
|
|
- remove_session(id);
|
|
|
- pthread_mutex_destroy(&session->lock);
|
|
|
- free(session);
|
|
|
-
|
|
|
- LOG_INFO("Session destroyed: ID=%d", 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);
|
|
|
-
|
|
|
+int stream_close(int id){
|
|
|
+ pthread_mutex_lock(&g_arr);
|
|
|
+ if(id<0||id>=MAX_STREAM_NUM||!g_sess[id].used){ pthread_mutex_unlock(&g_arr); return INVALID_PARAM; }
|
|
|
+ StreamSession*s=&g_sess[id];
|
|
|
+ /* snapshot before release (no use-after-free) */
|
|
|
+ int ch=s->channel; uint64_t f=s->frames, b=s->bytes_pushed;
|
|
|
+ pthread_mutex_destroy(&s->lock);
|
|
|
+ memset(s,0,sizeof(*s)); s->used=0; /* compact array by clearing slot; no tail-drop */
|
|
|
+ pthread_mutex_unlock(&g_arr);
|
|
|
+ LOG_INFO("stream close ch=%d frames=%" PRIu64 " bytes=%" PRIu64, ch,f,b);
|
|
|
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;
|
|
|
+void stream_stats(int*active,uint64_t*tf,uint64_t*tb){
|
|
|
+ int a=0; uint64_t f=0,b=0;
|
|
|
+ for(int i=0;i<MAX_STREAM_NUM;i++){ if(g_sess[i].used){a++;f+=g_sess[i].frames;b+=g_sess[i].bytes_pushed;} }
|
|
|
+ if(active)*active=a;
|
|
|
+ if(tf)*tf=f;
|
|
|
+ if(tb)*tb=b;
|
|
|
+}
|
|
|
+
|
|
|
+/* Simulate N frames of a synthetic Annex-B H264 IDR + non-IDR into PS. */
|
|
|
+int stream_simulate_loop(int channel,int frames){
|
|
|
+ int id=0; if(stream_open(channel,NULL,&id)!=SUCCESS) return FAILURE;
|
|
|
+ uint8_t sps[]={0x00,0x00,0x00,0x01,0x67,0x42,0xC0,0x1E,0xD8,0x60,0x1E,0x06,0x84};
|
|
|
+ uint8_t idr[]={0x00,0x00,0x00,0x01,0x65,0x88,0x80,0x40,0x03,0x7F,0x1E,0x30,0x89};
|
|
|
+ uint8_t non[]={0x00,0x00,0x00,0x01,0x41,0x9A,0x20,0x04,0x7F,0x1E};
|
|
|
+ int fps=g_cfg.fps>0?g_cfg.fps:SIM_FPS;
|
|
|
+ for(int i=0;i<frames;i++){
|
|
|
+ uint64_t pts=(uint64_t)(i*1000/fps);
|
|
|
+ uint64_t dts=pts;
|
|
|
+ int key=(i%fps==0);
|
|
|
+ if(key){ /* IDR preceded by SPS */
|
|
|
+ uint8_t buf[64]; size_t k=0; memcpy(buf+k,sps,sizeof(sps));k+=sizeof(sps);
|
|
|
+ memcpy(buf+k,idr,sizeof(idr));k+=sizeof(idr);
|
|
|
+ stream_push_frame(id,buf,k,pts,dts,1);
|
|
|
+ } else {
|
|
|
+ stream_push_frame(id,non,sizeof(non),pts,dts,0);
|
|
|
}
|
|
|
- free(ctx);
|
|
|
- session->internal_handle = NULL;
|
|
|
}
|
|
|
-
|
|
|
- session->state = STREAM_STATE_IDLE;
|
|
|
-
|
|
|
- LOG_INFO("Stream push stopped: ID=%d", session->stream_id);
|
|
|
+ int active; uint64_t f,b; stream_stats(&active,&f,&b);
|
|
|
+ LOG_INFO("simulate ch=%d done: %d frames pushed, total ps bytes=%" PRIu64, channel, frames, b);
|
|
|
+ stream_close(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;
|
|
|
- }
|
|
|
-
|
|
|
- /* [FIX] 将视频帧(NALU)封装为PS流后再推; 原pack_h264_to_ps从未被调用 */
|
|
|
- int ps_len = pack_h264_to_ps(session, (const uint8_t*)frame->data, frame->data_size);
|
|
|
- if (ps_len < 0) {
|
|
|
- LOG_WARN("PS pack failed for stream ID=%d", session->stream_id);
|
|
|
- return FAILURE;
|
|
|
- }
|
|
|
-
|
|
|
- RTMPPushContext* ctx = (RTMPPushContext*)session->internal_handle;
|
|
|
- if (ctx && ctx->socket_fd >= 0) {
|
|
|
- /* 通过RTMP发送PS数据(实际实现中使用zlmediakit API或FFmpeg) */
|
|
|
- }
|
|
|
-
|
|
|
- return SUCCESS;
|
|
|
-}
|
|
|
-
|
|
|
-StreamSession* stream_find_by_channel(const char* device_id, int channel) {
|
|
|
- if (device_id == NULL) return NULL;
|
|
|
- StreamSession* found = NULL;
|
|
|
- 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) {
|
|
|
- found = g_sessions[i];
|
|
|
- break;
|
|
|
- }
|
|
|
- }
|
|
|
- pthread_mutex_unlock(&g_session_mutex);
|
|
|
- return found;
|
|
|
-}
|
|
|
-
|
|
|
-int stream_stop_channel(const char* device_id, int channel) {
|
|
|
- StreamSession* sess = stream_find_by_channel(device_id, channel);
|
|
|
- if (sess != NULL) {
|
|
|
- stream_stop_push(sess);
|
|
|
- return SUCCESS;
|
|
|
- }
|
|
|
- return NOT_FOUND;
|
|
|
-}
|
|
|
-
|
|
|
-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);
|
|
|
-}
|
|
|
-
|
|
|
-StreamSession* stream_get_session(int stream_id) {
|
|
|
- /* find_session 的公共封装, 供上层按ID查询 */
|
|
|
- return find_session(stream_id);
|
|
|
-}
|
|
|
-
|
|
|
-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;
|
|
|
-}
|
|
|
+void stream_cleanup(void){ stream_stop_listen(); }
|