|
@@ -0,0 +1,156 @@
|
|
|
|
|
+#include "stream.h"
|
|
|
|
|
+#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;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+/* ---- 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;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+/* 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;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+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;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+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);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ 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;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void stream_cleanup(void){ stream_stop_listen(); }
|