Loading src/main.cpp +12 −4 Original line number Diff line number Diff line Loading @@ -38,11 +38,17 @@ using namespace std; * @return None */ void clean_up(pthread_t *threads, size_t sz) { for (size_t i = 0; i < sz; i++) pthread_cancel(threads[i]); int ret_val; for (size_t i = 0; i < sz; i++) { ret_val = pthread_cancel(threads[i]); if (!ret_val) cout << "pthread_cancel returned " << ret_val << ", sensor port " << i << endl; } } int main(int argc, char *argv[]) { int ret_val; string opt; map<string, string> args; pthread_t threads[SENSOR_PORTS]; Loading Loading @@ -72,8 +78,10 @@ int main(int argc, char *argv[]) { streamers[i] = new Streamer(args, i); pthread_attr_init(&attr); if (!pthread_create(&threads[i], &attr, Streamer::pthread_f, (void *) streamers[i])) { cerr << "Can not spawn streamer thread for port " << i << endl; ret_val = pthread_create(&threads[i], &attr, Streamer::pthread_f, (void *) streamers[i]); if (ret_val != 0) { cerr << "Can not spawn streamer thread for port " << i; cerr << ", pthread_create returned " << ret_val << endl; clean_up(threads, SENSOR_PORTS); exit(EXIT_FAILURE); } Loading src/parameters.h +1 −0 Original line number Diff line number Diff line Loading @@ -47,6 +47,7 @@ public: off_t lseek(off_t offset, int whence) { return ::lseek(fd_fparmsall, offset, whence); } bool daemon_enabled(void); void setPValue(unsigned long *val_array, int count); inline int get_port_num() const {return sensor_port;} protected: // static Parameters *_parameters; Loading src/streamer.cpp +116 −108 Original line number Diff line number Diff line Loading @@ -37,15 +37,23 @@ using namespace std; #define RTSP_DEBUG_2 #ifdef RTSP_DEBUG #define D(a) a #define D(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D(a) #define D(s_port, a) #endif #ifdef RTSP_DEBUG_2 #define D2(a) a #define D2(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D2(a) #define D2(s_port, a) #endif //Streamer *Streamer::_streamer = NULL; Loading @@ -70,7 +78,7 @@ Streamer::Streamer(const map<string, string> &_args, int port_num) { fps = atof(args["f"].c_str()); if (fps < 0.1) fps = 0; D( cout << "use fps: " << fps << endl;) D(sensor_port, cout << "use fps: " << fps << endl); video->fps(fps); } rtsp_server = NULL; Loading @@ -79,10 +87,10 @@ D( cout << "use fps: " << fps << endl;) void Streamer::audio_init(void) { if (audio != NULL) { D(cerr << "delete audio" << endl;) D(sensor_port, cerr << "delete audio" << endl); delete audio; } D(cout << "audio_enabled == " << session->process_audio << endl;) D(sensor_port, cout << "audio_enabled == " << session->process_audio << endl); audio = new Audio(session->process_audio, params, session->audio.sample_rate, session->audio.channels); if (audio->present() && session->process_audio) { session->process_audio = true; Loading @@ -109,7 +117,7 @@ int Streamer::f_handler(void *ptr, RTSP_Server *rtsp_server, RTSP_Server::event } int Streamer::update_settings(bool apply) { D( cerr << "update_settings" << endl;) D(sensor_port, cerr << "update_settings" << endl); // check settings, normalize its, return 1 if was changed // update settings at application if apply = 1 and parameters change isn't on-fly safe, update parameters always Loading Loading @@ -185,7 +193,7 @@ D( cerr << "update_settings" << endl;) if (ip > a_max) ip = a_max; a.s_addr = htonl(ip); D( cerr << "multicast ip asked: " << inet_ntoa(a) << endl;) D(sensor_port, cerr << "multicast ip asked: " << inet_ntoa(a) << endl); if (apply) { session->rtp_out.ip_cached = ip; session->rtp_out.ip_custom = true; Loading @@ -196,7 +204,7 @@ D( cerr << "multicast ip asked: " << inet_ntoa(a) << endl;) } } else { struct in_addr a = Socket::mcast_from_local(); D( cerr << "multicast ip generated: " << inet_ntoa(a) << endl;) D(sensor_port, cerr << "multicast ip generated: " << inet_ntoa(a) << endl); if (apply) { session->rtp_out.ip_custom = false; session->rtp_out.ip = inet_ntoa(a); Loading @@ -204,8 +212,8 @@ D( cerr << "multicast ip generated: " << inet_ntoa(a) << endl;) } transport_was_changed = true; } D( if(apply)) D( cerr << "actual multicast IP: " << session->rtp_out.ip << endl;) //D( if(apply)) D(sensor_port, if (apply) cerr << "actual multicast IP: " << session->rtp_out.ip << endl); // port int port = params->getGPValue(P_STROP_MCAST_PORT); if (port != session->rtp_out.port_video) { Loading Loading @@ -269,7 +277,7 @@ D( cerr << "actual multicast IP: " << session->rtp_out.ip << endl;) session->process_audio = audio_proc; session->audio.sample_rate = audio_rate; session->audio.channels = audio_channels; D2( cerr << "Audio was changed. Should restart it" << endl;) D2(sensor_port, cerr << "Audio was changed. Should restart it" << endl); audio_init(); audio_restarted = true; // if audio enable was asked, check what soundcard really is connected Loading Loading @@ -344,7 +352,7 @@ D2( cerr << "Audio was changed. Should restart it" << endl;) int Streamer::handler(RTSP_Server *rtsp_server, RTSP_Server::event event) { static bool _play = false; D( cerr << "event: running= " << running << " ";) D(sensor_port, cerr << "event: running= " << running << " "); switch (event) { case RTSP_Server::DESCRIBE: /// Update frame size, fps before starting new stream (generating SDP file) update_settings(true); Loading @@ -352,12 +360,13 @@ D( cerr << "event: running= " << running << " ";) case RTSP_Server::PARAMS_WAS_CHANGED: /// Update frame size, fps before starting new stream (generating SDP file) return (update_settings(false) || !(params->daemon_enabled())); case RTSP_Server::PLAY: D( cerr << "==PLAY==";) D(sensor_port, cerr << "==PLAY=="); if (connected_count == 0) { int ttl = -1; if (session->rtp_out.multicast) ttl = atoi(session->rtp_out.ttl.c_str()); video->Start(session->rtp_out.ip, session->rtp_out.port_video, session->video.fps_scale, ttl); video->Start(session->rtp_out.ip, session->rtp_out.port_video, session->video.fps_scale, ttl); if (audio != NULL) audio->Start(session->rtp_out.ip, session->rtp_out.port_audio, ttl); } Loading @@ -366,7 +375,7 @@ D( cerr << "==PLAY==";) running = true; break; case RTSP_Server::PAUSE: D( cerr << "PAUSE";) D(sensor_port, cerr << "PAUSE"); connected_count--; if (connected_count <= 0) { video->Stop(); Loading @@ -378,9 +387,9 @@ D( cerr << "PAUSE";) } break; case RTSP_Server::TEARDOWN: D( cerr << "TEARDOWN";) D(sensor_port, cerr << "TEARDOWN"); if (!running) { D( cerr << " was not running";) D(sensor_port, cerr << " was not running"); break; } connected_count--; Loading @@ -394,9 +403,9 @@ D( cerr << " was not running";) } break; case RTSP_Server::RESET: D( cerr << "RESET";) D(sensor_port, cerr << "RESET"); if (!running) { D( cerr << " was not running";) D(sensor_port, cerr << " was not running"); break; } video->Stop(); Loading @@ -413,16 +422,15 @@ D( cerr << "IS_DAEMON_ENABLED video->isDaemonEnabled(-1)=" << video->isDaemonEn break; */ default: D( cerr << "unknown == " << event;) D(sensor_port, cerr << "unknown == " << event); break; } D( cerr << endl;) D(sensor_port, cerr << endl); return 0; } void Streamer::Main(void) { D( cerr << "start Main for sensor port " << sensor_port << endl;) string def_mcast = "232.1.1.1"; D(sensor_port, cerr << "start Main for sensor port " << sensor_port << endl); int def_port = 20020; string def_ttl = "2"; Loading @@ -442,18 +450,18 @@ void Streamer::Main(void) { update_settings(true); /// Got here if is and was enabled (may use more actions instead of just "continue" // start RTSP server D2( cerr << "start server" << endl;) D2(sensor_port, cerr << "start server" << endl); if (rtsp_server == NULL) rtsp_server = new RTSP_Server(Streamer::f_handler, (void *) this, params, session); rtsp_server->main(); D2( cerr << "server was stopped" << endl;) D2( cerr << "stop video" << endl;) D2(sensor_port, cerr << "server was stopped" << endl); D2(sensor_port, cerr << "stop video" << endl); video->Stop(); D2( cerr << "stop audio" << endl;) D2(sensor_port, cerr << "stop audio" << endl); if (audio != NULL) { audio->Stop(); // free audio resource - other app can use soundcard D2( cerr << "delete audio" << endl;) D2(sensor_port, cerr << "delete audio" << endl); delete audio; audio = NULL; } Loading src/video.cpp +122 −117 Original line number Diff line number Diff line Loading @@ -45,19 +45,31 @@ using namespace std; #define VIDEO_DEBUG_3 // for FPS monitoring #ifdef VIDEO_DEBUG #define D(a) a #define D(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D(a) #endif #ifdef VIDEO_DEBUG_2 #define D2(a) a #define D2(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D2(a) #endif #ifdef VIDEO_DEBUG_3 #define D3(a) a #define D3(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D3(a) #endif Loading Loading @@ -87,8 +99,7 @@ static const char *jhead_file_names[] = { Video::Video(int port, Parameters *pars) { string err_msg; D( cerr << "Video::Video() on port " << port << endl;) D( cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << endl;) D(sensor_port, cerr << "Video::Video() on sensor port " << port << endl); params = pars; sensor_port = port; stream_name = "video"; Loading @@ -107,21 +118,16 @@ Video::Video(int port, Parameters *pars) { err_msg = "can't mmap " + *circbuf_file_names[sensor_port]; throw runtime_error(err_msg); } cout << "<-- 1" << endl; // buffer_ptr_s = (unsigned long *) mmap(buffer_ptr + (buffer_length >> 2), buffer_length, // PROT_READ, MAP_FIXED | MAP_SHARED, fd_circbuf, 0); /// preventing buffer rollovers buffer_ptr_s = (unsigned long *) mmap(buffer_ptr + (buffer_length >> 2), 100 * 4096, PROT_READ, MAP_FIXED | MAP_SHARED, fd_circbuf, 0); /// preventing buffer rollovers cout << "<-- 2" << endl; if ((int) buffer_ptr_s == -1) { err_msg = "can't create second mmap for " + *circbuf_file_names[sensor_port]; throw runtime_error(err_msg); } cout << "<-- 3" << endl; /// Skip several frames if it is just booted /// May get stuck here if compressor is off, it should be enabled externally D( cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " frame=" << params->getGPValue(G_THIS_FRAME) << " buffer_length=" << buffer_length << endl;) D(sensor_port, cerr << " frame=" << params->getGPValue(G_THIS_FRAME) << " buffer_length=" << buffer_length << endl); while (params->getGPValue(G_THIS_FRAME) < 10) { lseek(fd_circbuf, LSEEK_CIRC_TOWP, SEEK_END); /// get to the end of buffer lseek(fd_circbuf, LSEEK_CIRC_WAIT, SEEK_END); /// wait frame got ready there Loading @@ -129,7 +135,7 @@ Video::Video(int port, Parameters *pars) { /// One more wait always to make sure compressor is actually running lseek(fd_circbuf, LSEEK_CIRC_WAIT, SEEK_END); lseek(fd_circbuf, LSEEK_CIRC_WAIT, SEEK_END); D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " frame=" << params->getGPValue(G_THIS_FRAME) << " buffer_length=" << buffer_length <<endl;) D(sensor_port, cerr << " frame=" << params->getGPValue(G_THIS_FRAME) << " buffer_length=" << buffer_length <<endl); fd_jpeghead = open(jhead_file_names[sensor_port], O_RDWR); if (fd_jpeghead < 0) { err_msg = "can't open " + *jhead_file_names[sensor_port]; Loading @@ -147,7 +153,7 @@ Video::Video(int port, Parameters *pars) { // create thread... init_pthread((void *) this); D( cerr << __FILE__<< ":" << __FUNCTION__ << ":" << __LINE__ << endl;) D(sensor_port, cerr << "finish constructor" << endl); } Video::~Video(void) { Loading @@ -169,7 +175,7 @@ Video::~Video(void) { /// Compressor should be turned on outside of the streamer #define TURN_COMPRESSOR_ON 0 void Video::Start(string ip, long port, int _fps_scale, int ttl) { D( cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << "_play=" << _play << endl;) D(sensor_port, cerr << "_play=" << _play << endl); if (_play) { cerr << "ERROR-->> wrong usage: Video()->Start() when already play!!!" << endl; return; Loading Loading @@ -226,9 +232,9 @@ bool Video::waitDaemonEnabled(int daemonBit) { // <0 - use default unsigned long this_frame = params->getGPValue(G_THIS_FRAME); /// No semaphors, so it is possible to miss event and wait until the streamer will be re-enabled before sending message, /// but it seems not so terrible D(cerr << " lseek(fd_circbuf" << fd_circbuf << ", LSEEK_DAEMON_CIRCBUF+lastDaemonBit, SEEK_END)... " << endl;) D(sensor_port, cerr << " lseek(fd_circbuf" << sensor_port << ", LSEEK_DAEMON_CIRCBUF+lastDaemonBit, SEEK_END)... " << endl); lseek(fd_circbuf, LSEEK_DAEMON_CIRCBUF + lastDaemonBit, SEEK_END); /// D(cerr << "...done" << endl;) D(sensor_port, cerr << "...done" << endl); if (this_frame == params->getGPValue(G_THIS_FRAME)) return true; Loading Loading @@ -293,11 +299,13 @@ long Video::getFramePars(struct interframe_params_t *frame_pars, long before, lo if (metadata_start < 0) metadata_start += buffer_length; /// copy the interframe data (timestamps are not yet there) D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " before=" << before << " metadata_start=" << metadata_start << endl;) D(sensor_port, cerr << " before=" << before << " metadata_start=" << metadata_start << endl); memcpy(frame_pars, &char_buffer_ptr[metadata_start], 32); long jpeg_len = frame_pars->frame_length; //! frame_pars->frame_length is now the length of bitstream if (frame_pars->signffff != 0xffff) { cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " Wrong signature in getFramePars() (broken frame), frame_pars->signffff="<< frame_pars->signffff << endl; cerr << __FILE__ << ":" << __FUNCTION__ << ":" << __LINE__ << " Wrong signature in getFramePars() (broken frame), frame_pars->signffff=" << frame_pars->signffff << endl; int i; long * dd = (long *) frame_pars; cerr << hex << (metadata_start / 4) << ": "; Loading Loading @@ -342,13 +350,9 @@ struct video_desc_t Video::get_current_desc(bool with_fps) { fps *= 1000000.0; fps += frame_pars.timestamp_usec; fps -= prev_pars.timestamp_usec; D3( float _f = fps;) fps = 1000000.0 / fps; video_desc.fps = (used_fps = fps); D( cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " fps=" << fps << endl;) // cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " fps=" << video_desc.fps << endl; D3(if(_f == 0)) D3(cerr << "delta == " << _f << endl << endl;) D(sensor_port, cerr << " fps=" << fps << endl); } } video_desc.valid = true; Loading Loading @@ -423,7 +427,8 @@ long Video::capture(void) { /// Each time the latest acquired frame is considered, so we do not need to save frmae poointer additionally if ((fp->width != used_width) || (fp->height != used_height)) { for (before = 1; before <= (int) params->getGPValue(G_SKIP_DIFF_FRAME); before++) { if(((frameStartByteIndex = getFramePars(&frame_pars, before))) && (frame_pars.width == used_width) && (frame_pars.height == used_height)) { if (((frameStartByteIndex = getFramePars(&frame_pars, before))) && (frame_pars.width == used_width) && (frame_pars.height == used_height)) { /// substitute older frame instead of the latest one. Leave wrong timestamp? /// copying code above (may need some cleanup). Maybe - just move earlier so there will be no code duplication? latestAvailableFrame_ptr = frameStartByteIndex; Loading Loading @@ -452,16 +457,16 @@ long Video::capture(void) { } } if (before > (int) params->getGPValue(G_SKIP_DIFF_FRAME)) { D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__<< " Killing stream because of frame size change " << endl;) D(sensor_port, cerr << " Killing stream because of frame size change " << endl); return -SIZE_CHANGE; /// It seems that frame size is changed for good, need to restart the stream } D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__<< " Waiting for the original frame size to be restored , using " << before << " frames ago" << endl;) D(sensor_port, cerr << " Waiting for the original frame size to be restored , using " << before << " frames ago" << endl); } ///long Video::getFramePars(struct interframe_params_t * frame_pars, long before) { ///getGPValue(unsigned long GPNumber) quality = fp->quality2; if (qtables_include && quality != f_quality) { D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__<< " Updating quality tables, new quality is " << quality << endl;) D(sensor_port, cerr << " Updating quality tables, new quality is " << quality << endl); lseek(fd_jpeghead, frameStartByteIndex | 2, SEEK_END); /// '||2' indicates that we need just quantization tables, not full JPEG header read(fd_jpeghead, (void *) &qtable[0], 128); } Loading Loading @@ -491,7 +496,7 @@ long Video::process(void) { return 0; /// now never here } else { if (frame_len < 0) { D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__<< "capture returned negative" << frame_len << endl;) D(sensor_port, cerr << "capture returned negative" << frame_len << endl); // return false; return frame_len; /// attention (restart) is needed } Loading Loading
src/main.cpp +12 −4 Original line number Diff line number Diff line Loading @@ -38,11 +38,17 @@ using namespace std; * @return None */ void clean_up(pthread_t *threads, size_t sz) { for (size_t i = 0; i < sz; i++) pthread_cancel(threads[i]); int ret_val; for (size_t i = 0; i < sz; i++) { ret_val = pthread_cancel(threads[i]); if (!ret_val) cout << "pthread_cancel returned " << ret_val << ", sensor port " << i << endl; } } int main(int argc, char *argv[]) { int ret_val; string opt; map<string, string> args; pthread_t threads[SENSOR_PORTS]; Loading Loading @@ -72,8 +78,10 @@ int main(int argc, char *argv[]) { streamers[i] = new Streamer(args, i); pthread_attr_init(&attr); if (!pthread_create(&threads[i], &attr, Streamer::pthread_f, (void *) streamers[i])) { cerr << "Can not spawn streamer thread for port " << i << endl; ret_val = pthread_create(&threads[i], &attr, Streamer::pthread_f, (void *) streamers[i]); if (ret_val != 0) { cerr << "Can not spawn streamer thread for port " << i; cerr << ", pthread_create returned " << ret_val << endl; clean_up(threads, SENSOR_PORTS); exit(EXIT_FAILURE); } Loading
src/parameters.h +1 −0 Original line number Diff line number Diff line Loading @@ -47,6 +47,7 @@ public: off_t lseek(off_t offset, int whence) { return ::lseek(fd_fparmsall, offset, whence); } bool daemon_enabled(void); void setPValue(unsigned long *val_array, int count); inline int get_port_num() const {return sensor_port;} protected: // static Parameters *_parameters; Loading
src/streamer.cpp +116 −108 Original line number Diff line number Diff line Loading @@ -37,15 +37,23 @@ using namespace std; #define RTSP_DEBUG_2 #ifdef RTSP_DEBUG #define D(a) a #define D(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D(a) #define D(s_port, a) #endif #ifdef RTSP_DEBUG_2 #define D2(a) a #define D2(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D2(a) #define D2(s_port, a) #endif //Streamer *Streamer::_streamer = NULL; Loading @@ -70,7 +78,7 @@ Streamer::Streamer(const map<string, string> &_args, int port_num) { fps = atof(args["f"].c_str()); if (fps < 0.1) fps = 0; D( cout << "use fps: " << fps << endl;) D(sensor_port, cout << "use fps: " << fps << endl); video->fps(fps); } rtsp_server = NULL; Loading @@ -79,10 +87,10 @@ D( cout << "use fps: " << fps << endl;) void Streamer::audio_init(void) { if (audio != NULL) { D(cerr << "delete audio" << endl;) D(sensor_port, cerr << "delete audio" << endl); delete audio; } D(cout << "audio_enabled == " << session->process_audio << endl;) D(sensor_port, cout << "audio_enabled == " << session->process_audio << endl); audio = new Audio(session->process_audio, params, session->audio.sample_rate, session->audio.channels); if (audio->present() && session->process_audio) { session->process_audio = true; Loading @@ -109,7 +117,7 @@ int Streamer::f_handler(void *ptr, RTSP_Server *rtsp_server, RTSP_Server::event } int Streamer::update_settings(bool apply) { D( cerr << "update_settings" << endl;) D(sensor_port, cerr << "update_settings" << endl); // check settings, normalize its, return 1 if was changed // update settings at application if apply = 1 and parameters change isn't on-fly safe, update parameters always Loading Loading @@ -185,7 +193,7 @@ D( cerr << "update_settings" << endl;) if (ip > a_max) ip = a_max; a.s_addr = htonl(ip); D( cerr << "multicast ip asked: " << inet_ntoa(a) << endl;) D(sensor_port, cerr << "multicast ip asked: " << inet_ntoa(a) << endl); if (apply) { session->rtp_out.ip_cached = ip; session->rtp_out.ip_custom = true; Loading @@ -196,7 +204,7 @@ D( cerr << "multicast ip asked: " << inet_ntoa(a) << endl;) } } else { struct in_addr a = Socket::mcast_from_local(); D( cerr << "multicast ip generated: " << inet_ntoa(a) << endl;) D(sensor_port, cerr << "multicast ip generated: " << inet_ntoa(a) << endl); if (apply) { session->rtp_out.ip_custom = false; session->rtp_out.ip = inet_ntoa(a); Loading @@ -204,8 +212,8 @@ D( cerr << "multicast ip generated: " << inet_ntoa(a) << endl;) } transport_was_changed = true; } D( if(apply)) D( cerr << "actual multicast IP: " << session->rtp_out.ip << endl;) //D( if(apply)) D(sensor_port, if (apply) cerr << "actual multicast IP: " << session->rtp_out.ip << endl); // port int port = params->getGPValue(P_STROP_MCAST_PORT); if (port != session->rtp_out.port_video) { Loading Loading @@ -269,7 +277,7 @@ D( cerr << "actual multicast IP: " << session->rtp_out.ip << endl;) session->process_audio = audio_proc; session->audio.sample_rate = audio_rate; session->audio.channels = audio_channels; D2( cerr << "Audio was changed. Should restart it" << endl;) D2(sensor_port, cerr << "Audio was changed. Should restart it" << endl); audio_init(); audio_restarted = true; // if audio enable was asked, check what soundcard really is connected Loading Loading @@ -344,7 +352,7 @@ D2( cerr << "Audio was changed. Should restart it" << endl;) int Streamer::handler(RTSP_Server *rtsp_server, RTSP_Server::event event) { static bool _play = false; D( cerr << "event: running= " << running << " ";) D(sensor_port, cerr << "event: running= " << running << " "); switch (event) { case RTSP_Server::DESCRIBE: /// Update frame size, fps before starting new stream (generating SDP file) update_settings(true); Loading @@ -352,12 +360,13 @@ D( cerr << "event: running= " << running << " ";) case RTSP_Server::PARAMS_WAS_CHANGED: /// Update frame size, fps before starting new stream (generating SDP file) return (update_settings(false) || !(params->daemon_enabled())); case RTSP_Server::PLAY: D( cerr << "==PLAY==";) D(sensor_port, cerr << "==PLAY=="); if (connected_count == 0) { int ttl = -1; if (session->rtp_out.multicast) ttl = atoi(session->rtp_out.ttl.c_str()); video->Start(session->rtp_out.ip, session->rtp_out.port_video, session->video.fps_scale, ttl); video->Start(session->rtp_out.ip, session->rtp_out.port_video, session->video.fps_scale, ttl); if (audio != NULL) audio->Start(session->rtp_out.ip, session->rtp_out.port_audio, ttl); } Loading @@ -366,7 +375,7 @@ D( cerr << "==PLAY==";) running = true; break; case RTSP_Server::PAUSE: D( cerr << "PAUSE";) D(sensor_port, cerr << "PAUSE"); connected_count--; if (connected_count <= 0) { video->Stop(); Loading @@ -378,9 +387,9 @@ D( cerr << "PAUSE";) } break; case RTSP_Server::TEARDOWN: D( cerr << "TEARDOWN";) D(sensor_port, cerr << "TEARDOWN"); if (!running) { D( cerr << " was not running";) D(sensor_port, cerr << " was not running"); break; } connected_count--; Loading @@ -394,9 +403,9 @@ D( cerr << " was not running";) } break; case RTSP_Server::RESET: D( cerr << "RESET";) D(sensor_port, cerr << "RESET"); if (!running) { D( cerr << " was not running";) D(sensor_port, cerr << " was not running"); break; } video->Stop(); Loading @@ -413,16 +422,15 @@ D( cerr << "IS_DAEMON_ENABLED video->isDaemonEnabled(-1)=" << video->isDaemonEn break; */ default: D( cerr << "unknown == " << event;) D(sensor_port, cerr << "unknown == " << event); break; } D( cerr << endl;) D(sensor_port, cerr << endl); return 0; } void Streamer::Main(void) { D( cerr << "start Main for sensor port " << sensor_port << endl;) string def_mcast = "232.1.1.1"; D(sensor_port, cerr << "start Main for sensor port " << sensor_port << endl); int def_port = 20020; string def_ttl = "2"; Loading @@ -442,18 +450,18 @@ void Streamer::Main(void) { update_settings(true); /// Got here if is and was enabled (may use more actions instead of just "continue" // start RTSP server D2( cerr << "start server" << endl;) D2(sensor_port, cerr << "start server" << endl); if (rtsp_server == NULL) rtsp_server = new RTSP_Server(Streamer::f_handler, (void *) this, params, session); rtsp_server->main(); D2( cerr << "server was stopped" << endl;) D2( cerr << "stop video" << endl;) D2(sensor_port, cerr << "server was stopped" << endl); D2(sensor_port, cerr << "stop video" << endl); video->Stop(); D2( cerr << "stop audio" << endl;) D2(sensor_port, cerr << "stop audio" << endl); if (audio != NULL) { audio->Stop(); // free audio resource - other app can use soundcard D2( cerr << "delete audio" << endl;) D2(sensor_port, cerr << "delete audio" << endl); delete audio; audio = NULL; } Loading
src/video.cpp +122 −117 Original line number Diff line number Diff line Loading @@ -45,19 +45,31 @@ using namespace std; #define VIDEO_DEBUG_3 // for FPS monitoring #ifdef VIDEO_DEBUG #define D(a) a #define D(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D(a) #endif #ifdef VIDEO_DEBUG_2 #define D2(a) a #define D2(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D2(a) #endif #ifdef VIDEO_DEBUG_3 #define D3(a) a #define D3(s_port, a) \ do { \ cerr << __FILE__ << ": " << __FUNCTION__ << ": " << __LINE__ << ": sensor port: " << s_port << " "; \ a; \ } while (0) #else #define D3(a) #endif Loading Loading @@ -87,8 +99,7 @@ static const char *jhead_file_names[] = { Video::Video(int port, Parameters *pars) { string err_msg; D( cerr << "Video::Video() on port " << port << endl;) D( cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << endl;) D(sensor_port, cerr << "Video::Video() on sensor port " << port << endl); params = pars; sensor_port = port; stream_name = "video"; Loading @@ -107,21 +118,16 @@ Video::Video(int port, Parameters *pars) { err_msg = "can't mmap " + *circbuf_file_names[sensor_port]; throw runtime_error(err_msg); } cout << "<-- 1" << endl; // buffer_ptr_s = (unsigned long *) mmap(buffer_ptr + (buffer_length >> 2), buffer_length, // PROT_READ, MAP_FIXED | MAP_SHARED, fd_circbuf, 0); /// preventing buffer rollovers buffer_ptr_s = (unsigned long *) mmap(buffer_ptr + (buffer_length >> 2), 100 * 4096, PROT_READ, MAP_FIXED | MAP_SHARED, fd_circbuf, 0); /// preventing buffer rollovers cout << "<-- 2" << endl; if ((int) buffer_ptr_s == -1) { err_msg = "can't create second mmap for " + *circbuf_file_names[sensor_port]; throw runtime_error(err_msg); } cout << "<-- 3" << endl; /// Skip several frames if it is just booted /// May get stuck here if compressor is off, it should be enabled externally D( cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " frame=" << params->getGPValue(G_THIS_FRAME) << " buffer_length=" << buffer_length << endl;) D(sensor_port, cerr << " frame=" << params->getGPValue(G_THIS_FRAME) << " buffer_length=" << buffer_length << endl); while (params->getGPValue(G_THIS_FRAME) < 10) { lseek(fd_circbuf, LSEEK_CIRC_TOWP, SEEK_END); /// get to the end of buffer lseek(fd_circbuf, LSEEK_CIRC_WAIT, SEEK_END); /// wait frame got ready there Loading @@ -129,7 +135,7 @@ Video::Video(int port, Parameters *pars) { /// One more wait always to make sure compressor is actually running lseek(fd_circbuf, LSEEK_CIRC_WAIT, SEEK_END); lseek(fd_circbuf, LSEEK_CIRC_WAIT, SEEK_END); D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " frame=" << params->getGPValue(G_THIS_FRAME) << " buffer_length=" << buffer_length <<endl;) D(sensor_port, cerr << " frame=" << params->getGPValue(G_THIS_FRAME) << " buffer_length=" << buffer_length <<endl); fd_jpeghead = open(jhead_file_names[sensor_port], O_RDWR); if (fd_jpeghead < 0) { err_msg = "can't open " + *jhead_file_names[sensor_port]; Loading @@ -147,7 +153,7 @@ Video::Video(int port, Parameters *pars) { // create thread... init_pthread((void *) this); D( cerr << __FILE__<< ":" << __FUNCTION__ << ":" << __LINE__ << endl;) D(sensor_port, cerr << "finish constructor" << endl); } Video::~Video(void) { Loading @@ -169,7 +175,7 @@ Video::~Video(void) { /// Compressor should be turned on outside of the streamer #define TURN_COMPRESSOR_ON 0 void Video::Start(string ip, long port, int _fps_scale, int ttl) { D( cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << "_play=" << _play << endl;) D(sensor_port, cerr << "_play=" << _play << endl); if (_play) { cerr << "ERROR-->> wrong usage: Video()->Start() when already play!!!" << endl; return; Loading Loading @@ -226,9 +232,9 @@ bool Video::waitDaemonEnabled(int daemonBit) { // <0 - use default unsigned long this_frame = params->getGPValue(G_THIS_FRAME); /// No semaphors, so it is possible to miss event and wait until the streamer will be re-enabled before sending message, /// but it seems not so terrible D(cerr << " lseek(fd_circbuf" << fd_circbuf << ", LSEEK_DAEMON_CIRCBUF+lastDaemonBit, SEEK_END)... " << endl;) D(sensor_port, cerr << " lseek(fd_circbuf" << sensor_port << ", LSEEK_DAEMON_CIRCBUF+lastDaemonBit, SEEK_END)... " << endl); lseek(fd_circbuf, LSEEK_DAEMON_CIRCBUF + lastDaemonBit, SEEK_END); /// D(cerr << "...done" << endl;) D(sensor_port, cerr << "...done" << endl); if (this_frame == params->getGPValue(G_THIS_FRAME)) return true; Loading Loading @@ -293,11 +299,13 @@ long Video::getFramePars(struct interframe_params_t *frame_pars, long before, lo if (metadata_start < 0) metadata_start += buffer_length; /// copy the interframe data (timestamps are not yet there) D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " before=" << before << " metadata_start=" << metadata_start << endl;) D(sensor_port, cerr << " before=" << before << " metadata_start=" << metadata_start << endl); memcpy(frame_pars, &char_buffer_ptr[metadata_start], 32); long jpeg_len = frame_pars->frame_length; //! frame_pars->frame_length is now the length of bitstream if (frame_pars->signffff != 0xffff) { cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " Wrong signature in getFramePars() (broken frame), frame_pars->signffff="<< frame_pars->signffff << endl; cerr << __FILE__ << ":" << __FUNCTION__ << ":" << __LINE__ << " Wrong signature in getFramePars() (broken frame), frame_pars->signffff=" << frame_pars->signffff << endl; int i; long * dd = (long *) frame_pars; cerr << hex << (metadata_start / 4) << ": "; Loading Loading @@ -342,13 +350,9 @@ struct video_desc_t Video::get_current_desc(bool with_fps) { fps *= 1000000.0; fps += frame_pars.timestamp_usec; fps -= prev_pars.timestamp_usec; D3( float _f = fps;) fps = 1000000.0 / fps; video_desc.fps = (used_fps = fps); D( cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " fps=" << fps << endl;) // cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__ << " fps=" << video_desc.fps << endl; D3(if(_f == 0)) D3(cerr << "delta == " << _f << endl << endl;) D(sensor_port, cerr << " fps=" << fps << endl); } } video_desc.valid = true; Loading Loading @@ -423,7 +427,8 @@ long Video::capture(void) { /// Each time the latest acquired frame is considered, so we do not need to save frmae poointer additionally if ((fp->width != used_width) || (fp->height != used_height)) { for (before = 1; before <= (int) params->getGPValue(G_SKIP_DIFF_FRAME); before++) { if(((frameStartByteIndex = getFramePars(&frame_pars, before))) && (frame_pars.width == used_width) && (frame_pars.height == used_height)) { if (((frameStartByteIndex = getFramePars(&frame_pars, before))) && (frame_pars.width == used_width) && (frame_pars.height == used_height)) { /// substitute older frame instead of the latest one. Leave wrong timestamp? /// copying code above (may need some cleanup). Maybe - just move earlier so there will be no code duplication? latestAvailableFrame_ptr = frameStartByteIndex; Loading Loading @@ -452,16 +457,16 @@ long Video::capture(void) { } } if (before > (int) params->getGPValue(G_SKIP_DIFF_FRAME)) { D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__<< " Killing stream because of frame size change " << endl;) D(sensor_port, cerr << " Killing stream because of frame size change " << endl); return -SIZE_CHANGE; /// It seems that frame size is changed for good, need to restart the stream } D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__<< " Waiting for the original frame size to be restored , using " << before << " frames ago" << endl;) D(sensor_port, cerr << " Waiting for the original frame size to be restored , using " << before << " frames ago" << endl); } ///long Video::getFramePars(struct interframe_params_t * frame_pars, long before) { ///getGPValue(unsigned long GPNumber) quality = fp->quality2; if (qtables_include && quality != f_quality) { D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__<< " Updating quality tables, new quality is " << quality << endl;) D(sensor_port, cerr << " Updating quality tables, new quality is " << quality << endl); lseek(fd_jpeghead, frameStartByteIndex | 2, SEEK_END); /// '||2' indicates that we need just quantization tables, not full JPEG header read(fd_jpeghead, (void *) &qtable[0], 128); } Loading Loading @@ -491,7 +496,7 @@ long Video::process(void) { return 0; /// now never here } else { if (frame_len < 0) { D(cerr << __FILE__<< ":"<< __FUNCTION__ << ":" <<__LINE__<< "capture returned negative" << frame_len << endl;) D(sensor_port, cerr << "capture returned negative" << frame_len << endl); // return false; return frame_len; /// attention (restart) is needed } Loading