Files
zoneminder/src/zm_rtsp_server.cpp
T
Isaac ConnorandClaude Opus 5 1f70965a98 fix: retry deadlocked queries with bounded backoff instead of failing them
zmDbDo, zmDbDoInsert and zmDbDoUpdate all decided what to do with a failed
query by testing for ER_LOCK_WAIT_TIMEOUT alone, which got both halves of
lock contention wrong:

- ER_LOCK_DEADLOCK was not retried at all. InnoDB resolves a deadlock by
  rolling one side back and expects that side to re-run; instead the query
  was logged and abandoned. Event creation goes through zmDbDoInsert, which
  is where this actually bites.

- ER_LOCK_WAIT_TIMEOUT re-ran immediately and forever, with no delay and no
  attempt limit, so two writers deadlocking against each other kept
  colliding on the same schedule.

Both now go through one retry decision: five attempts with a jittered
50ms-doubling backoff, then give up and report. The jitter is what stops
two contending sessions waking together and repeating the deadlock.

The wait happens under db_mutex, which every other database user in the
process is blocked on, so the budget is deliberately small -- about 3.1s
across all five attempts. That is still far less than the unbounded
ER_LOCK_WAIT_TIMEOUT loop it replaces, where each round costs a full
innodb_lock_wait_timeout. The change in behaviour is that a query which
would eventually have won after many minutes is now abandoned; it is
reported at Error with the attempt count.

The error is now logged once, when giving up, rather than on every round.
mysql_errno is read next to mysql_error rather than after the logging
call, so it cannot be clobbered in between.

Two other things, both small and both in the same area:

zmDbEscapeString called mysql_real_escape_string unconditionally, and that
reads the character set off the connection, so a closed handle sends it
into freed state. It now falls back to escaping the injection-relevant
characters itself. That fallback is only correct because the connection is
utf8mb4, where no byte of a multi-byte sequence is ASCII and so no sequence
can absorb a trailing backslash; the comment says so, since it would be
wrong for a character set like GBK. It deliberately does not take
db_mutex to read the flag: the logger calls this from Error(), and
zmDbFetch reaches Error() while holding db_mutex, so locking here would
self-deadlock the process on any failed query.

zm_rtsp_server was the one daemon closing the database without stopping
the queue that writes through it. Fixed at the call site rather than
inside zmDbClose, which holds db_mutex while the queue thread needs that
same mutex to drain -- joining it from in there would deadlock.

Tests: tests/zm_db_contention.cpp covers the retry budget only, 2011
assertions. Verified it fails when the budget is removed. The retry loop
itself needs a database and two contending sessions and is NOT covered;
it wants verifying against a real server under contention before this is
relied on. Full suite 12171 assertions in 133 test cases. Builds clean.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_015Y6FieTwEXuLhhR4e2yiax
2026-09-04 07:15:30 -04:00

399 lines
15 KiB
C++

//
// ZoneMinder RTSP Daemon
// Copyright (C) 2021 Isaac Connor
//
// This program is free software; you can redistribute it and/or
// modify it under the terms of the GNU General Public License
// as published by the Free Software Foundation; either version 2
// of the License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU General Public License for more details.
//
// You should have received a copy of the GNU General Public License
// along with this program; if not, write to the Free Software
// Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA.
//
/*
=head1 NAME
zm_rtsp_server - The ZoneMinder Server
=head1 SYNOPSIS
zmc -m <monitor_id>
zmc --monitor <monitor_id>
zmc -h
zmc --help
zmc -v
zmc --version
=head1 DESCRIPTION
This binary's job is to connect to fifo's provided by local zmc processes
and provide that stream over rtsp
=head1 OPTIONS
-m, --monitor_id - ID of a monitor to stream
-h, --help - Display usage information
-v, --version - Print the installed version of ZoneMinder
=cut
*/
#include "zm.h"
#include "zm_db.h"
#include "zm_define.h"
#include "zm_monitor.h"
#include "zm_rtsp_server_authenticator.h"
#include "zm_rtsp_server_fifo_h264_source.h"
#include "zm_rtsp_server_fifo_av1_source.h"
#include "zm_rtsp_server_fifo_adts_source.h"
#include "xop/G711USource.h"
#include "zm_signal.h"
#include "zm_time.h"
#include "zm_utils.h"
#include <getopt.h>
#include <iostream>
#include <vector>
#include "xop/RtspServer.h"
void Usage() {
fprintf(stderr, "zm_rtsp_server -m <monitor_id>\n");
fprintf(stderr, "Options:\n");
fprintf(stderr, " -m, --monitor <monitor_id> : We default to all monitors use this to specify just one\n");
fprintf(stderr, " -h, --help : This screen\n");
fprintf(stderr, " -v, --version : Report the installed version of ZoneMinder\n");
exit(0);
}
int main(int argc, char *argv[]) {
self = argv[0];
srand(getpid() * time(nullptr));
int monitor_id = -1;
static struct option long_options[] = {
{"monitor", 1, nullptr, 'm'},
{"help", 0, nullptr, 'h'},
{"version", 0, nullptr, 'v'},
{nullptr, 0, nullptr, 0}
};
while (1) {
int option_index = 0;
int c = getopt_long(argc, argv, "m:h:v", long_options, &option_index);
if (c == -1)
break;
switch (c) {
case 'm':
monitor_id = atoi(optarg);
break;
case 'h':
case '?':
Usage();
break;
case 'v':
std::cout << ZM_VERSION << "\n";
exit(0);
default:
// fprintf(stderr, "?? getopt returned character code 0%o ??\n", c);
break;
}
}
if (optind < argc) {
fprintf(stderr, "Extraneous options, ");
while (optind < argc)
printf("%s ", argv[optind++]);
printf("\n");
Usage();
}
const char *log_id_string = "zm_rtsp_server";
///std::string log_id_string = std::string("zm_rtsp_server");
///if ( monitor_id > 0 ) log_id_string += stringtf("_m%d", monitor_id);
logInit(log_id_string);
zmLoadStaticConfig();
zmDbConnect();
zmLoadDBConfig();
logInit(log_id_string);
if (!config.min_rtsp_port) {
Debug(1, "Not starting RTSP server because min_rtsp_port not set");
exit(-1);
}
HwCapsDetect();
std::string where = "`Deleted` = 0 AND `Capturing` != 'None' AND `RTSPServer` != false";
if (staticConfig.SERVER_ID)
where += stringtf(" AND `ServerId`=%d", staticConfig.SERVER_ID);
if (monitor_id > 0)
where += stringtf(" AND `Id`=%d", monitor_id);
Info("Starting RTSP Server version %s", ZM_VERSION);
zmSetDefaultHupHandler();
zmSetDefaultTermHandler();
zmSetDefaultDieHandler();
std::shared_ptr<xop::EventLoop> eventLoop(new xop::EventLoop());
std::shared_ptr<xop::RtspServer> rtspServer = xop::RtspServer::Create(eventLoop.get());
rtspServer->SetVersion("ZoneMinder RTSP Server");
if (config.opt_use_auth) {
std::shared_ptr<ZM_RtspServer_Authenticator> authenticator(new ZM_RtspServer_Authenticator());
rtspServer->SetAuthenticator(authenticator);
}
if (!rtspServer->Start("0.0.0.0", config.min_rtsp_port)) {
Debug(1, "Failed starting RTSP server on port %d", config.min_rtsp_port);
exit(-1);
}
std::unordered_map<unsigned int, xop::MediaSession *> sessions;
std::unordered_map<unsigned int, ZoneMinderFifoSource *> video_sources;
std::unordered_map<unsigned int, ZoneMinderFifoSource *> audio_sources;
std::unordered_map<unsigned int, std::shared_ptr<Monitor>> monitors;
while (!zm_terminate) {
std::unordered_map<unsigned int, std::shared_ptr<Monitor>> old_monitors = monitors;
std::vector<std::shared_ptr<Monitor>> new_monitors = Monitor::LoadMonitors(where, Monitor::QUERY);
for (const auto &monitor : new_monitors) {
auto old_monitor_it = old_monitors.find(monitor->Id());
if (old_monitor_it != old_monitors.end()
and
(old_monitor_it->second->GetRTSPStreamName() == monitor->GetRTSPStreamName())
) {
Debug(1, "Found monitor in oldmonitors, clearing it");
old_monitors.erase(old_monitor_it);
} else {
Debug(1, "Adding monitor %d to monitors", monitor->Id());
monitors[monitor->Id()] = monitor;
}
}
// Remove monitors that are no longer doing rtsp
for (auto it = old_monitors.begin(); it != old_monitors.end(); ++it) {
auto mid = it->first;
auto &monitor = it->second;
Debug(1, "Removing %d %s from monitors", monitor->Id(), monitor->Name());
monitors.erase(mid);
if (sessions.find(mid) != sessions.end()) {
if (video_sources.find(monitor->Id()) != video_sources.end()) {
delete video_sources[monitor->Id()];
video_sources.erase(monitor->Id());
}
if (audio_sources.find(monitor->Id()) != audio_sources.end()) {
delete audio_sources[monitor->Id()];
audio_sources.erase(monitor->Id());
}
rtspServer->RemoveSession(sessions[mid]->GetMediaSessionId());
sessions.erase(mid);
}
}
for (auto it = monitors.begin(); it != monitors.end(); ++it) {
auto &monitor = it->second;
if (!monitor->ShmValid()) {
Debug(1, "!ShmValid");
monitor->disconnect();
if (!monitor->connect()) {
Warning("Couldn't connect to monitor %d", monitor->Id());
if (sessions.find(monitor->Id()) != sessions.end()) {
// Delete (not just erase) the FifoSources first: their dtors stop
// and join the read/write threads. RemoveSession then destroys
// the MediaSession which owns the xop H264/H265/AV1Source; any
// still-running ReadRun thread holds a raw m_h264Source pointer
// and will crash in SetSPS on the next SPS NAL.
if (video_sources.find(monitor->Id()) != video_sources.end()) {
delete video_sources[monitor->Id()];
video_sources.erase(monitor->Id());
}
if (audio_sources.find(monitor->Id()) != audio_sources.end()) {
delete audio_sources[monitor->Id()];
audio_sources.erase(monitor->Id());
}
rtspServer->RemoveSession(sessions[monitor->Id()]->GetMediaSessionId());
sessions.erase(monitor->Id());
}
monitor->Reload(); // This is to pickup change of colours, width, height, etc
continue;
} // end if failed to connect
} // end if !ShmValid
if (sessions.end() == sessions.find(monitor->Id())) {
Debug(1, "Monitor not found in sessions, opening it");
std::string videoFifoPath = monitor->GetVideoFifoPath();
if (videoFifoPath.empty()) {
Debug(1, "video fifo is empty. Skipping.");
continue;
}
std::string streamname = monitor->GetRTSPStreamName();
xop::MediaSession *session = sessions[monitor->Id()] = xop::MediaSession::CreateNew(streamname);
if (!session) {
Error("Unable to create session for %s", streamname.c_str());
continue;
}
session->AddNotifyConnectedCallback([] (xop::MediaSessionId sessionId, const std::string &peer_ip, uint16_t peer_port) {
Debug(1, "RTSP client connect, ip=%s, port=%hu", peer_ip.c_str(), peer_port);
});
session->AddNotifyDisconnectedCallback([](xop::MediaSessionId sessionId, const std::string &peer_ip, uint16_t peer_port) {
Debug(1, "RTSP client disconnect, ip=%s, port=%hu", peer_ip.c_str(), peer_port);
});
rtspServer->AddSession(session);
//char *url = rtspServer->rtspURL(session);
//Debug(1, "url is %s for stream %s", url, streamname.c_str());
//delete[] url;
monitors[monitor->Id()] = monitor;
Debug(1, "Adding video fifo %s", videoFifoPath.c_str());
ZoneMinderFifoVideoSource *videoSource = nullptr;
if (std::string::npos != videoFifoPath.find("h264")) {
xop::H264Source *h264Source = xop::H264Source::CreateNew();
h264Source->SetResolution(monitor->Width(), monitor->Height());
session->AddSource(xop::channel_0, h264Source);
H264_ZoneMinderFifoSource *h264FifoSource = new H264_ZoneMinderFifoSource(rtspServer, session->GetMediaSessionId(), xop::channel_0, videoFifoPath);
h264FifoSource->setH264Source(h264Source); // Allow FIFO source to set SPS/PPS
videoSource = h264FifoSource;
} else if (
std::string::npos != videoFifoPath.find("hevc")
or
std::string::npos != videoFifoPath.find("h265")) {
xop::H265Source *h265Source = xop::H265Source::CreateNew();
h265Source->SetResolution(monitor->Width(), monitor->Height());
session->AddSource(xop::channel_0, h265Source);
H265_ZoneMinderFifoSource *h265FifoSource = new H265_ZoneMinderFifoSource(rtspServer, session->GetMediaSessionId(), xop::channel_0, videoFifoPath);
h265FifoSource->setH265Source(h265Source); // Allow FIFO source to set VPS/SPS/PPS
videoSource = h265FifoSource;
} else if (std::string::npos != videoFifoPath.find("av1")) {
xop::AV1Source *av1Source = xop::AV1Source::CreateNew();
av1Source->SetResolution(monitor->Width(), monitor->Height());
session->AddSource(xop::channel_0, av1Source);
AV1_ZoneMinderFifoSource *av1FifoSource = new AV1_ZoneMinderFifoSource(rtspServer, session->GetMediaSessionId(), xop::channel_0, videoFifoPath);
av1FifoSource->setAV1Source(av1Source); // Allow FIFO source to set sequence header
videoSource = av1FifoSource;
} else {
Warning("Unknown format in %s", videoFifoPath.c_str());
}
if (videoSource == nullptr) {
Error("Unable to create source for %s", videoFifoPath.c_str());
rtspServer->RemoveSession(sessions[monitor->Id()]->GetMediaSessionId());
sessions.erase(monitor->Id());
continue;
}
video_sources[monitor->Id()] = videoSource;
videoSource->setWidth(monitor->Width());
videoSource->setHeight(monitor->Height());
std::string audioFifoPath = monitor->GetAudioFifoPath();
if (audioFifoPath.empty()) {
Debug(1, "audio fifo is empty. Skipping.");
continue;
}
Debug(1, "Adding audio fifo %s", audioFifoPath.c_str());
ZoneMinderFifoAudioSource *audioSource = nullptr;
if (std::string::npos != audioFifoPath.find("aac")) {
Debug(1, "Adding aac source at %dHz %d channels",
monitor->GetAudioFrequency(), monitor->GetAudioChannels());
session->AddSource(xop::channel_1, xop::AACSource::CreateNew(
monitor->GetAudioFrequency(),
monitor->GetAudioChannels(),
false /* has_adts */));
audioSource = new ADTS_ZoneMinderFifoSource(rtspServer,
session->GetMediaSessionId(), xop::channel_1, audioFifoPath);
audioSource->setFrequency(monitor->GetAudioFrequency());
audioSource->setChannels(monitor->GetAudioChannels());
} else if (std::string::npos != audioFifoPath.find("pcm_alaw")) {
Debug(1, "Adding G711A source at %dHz %d channels",
monitor->GetAudioFrequency(), monitor->GetAudioChannels());
session->AddSource(xop::channel_1, xop::G711ASource::CreateNew());
audioSource = new ADTS_ZoneMinderFifoSource(rtspServer,
session->GetMediaSessionId(), xop::channel_1, audioFifoPath);
audioSource->setFrequency(monitor->GetAudioFrequency());
audioSource->setChannels(monitor->GetAudioChannels());
} else if (std::string::npos != audioFifoPath.find("pcm_mulaw")) {
Debug(1, "Adding G711U source at %dHz %d channels",
monitor->GetAudioFrequency(), monitor->GetAudioChannels());
session->AddSource(xop::channel_1, xop::G711USource::CreateNew());
audioSource = new ADTS_ZoneMinderFifoSource(rtspServer,
session->GetMediaSessionId(), xop::channel_1, audioFifoPath);
audioSource->setFrequency(monitor->GetAudioFrequency());
audioSource->setChannels(monitor->GetAudioChannels());
} else {
Warning("Unknown format in %s", audioFifoPath.c_str());
}
if (audioSource == nullptr) {
Error("Unable to create source");
}
audio_sources[monitor->Id()] = audioSource;
} // end if ! sessions[monitor->Id()]
if (sessions[monitor->Id()]->GetNumClient() > 0) {
SystemTimePoint now = std::chrono::system_clock::now();
monitor->setLastViewed(now);
}
} // end foreach monitor
sleep(10);
if (zm_reload) {
Info("Reloading configuration");
logTerm();
zmLoadDBConfig();
logInit(log_id_string);
zm_reload = false;
} // end if zm_reload
} // end while !zm_terminate
Info("RTSP Server shutting down");
for (const std::pair<const unsigned int, std::shared_ptr<Monitor>> &mon_pair : monitors) {
unsigned int i = mon_pair.first;
if (video_sources.find(i) != video_sources.end()) {
delete video_sources[i];
}
if (audio_sources.find(i) != audio_sources.end()) {
delete audio_sources[i];
}
if (sessions.find(i) != sessions.end()) {
Debug(1, "Removing session for %s", mon_pair.second->Name());
rtspServer->RemoveSession(sessions[i]->GetMediaSessionId());
sessions.erase(i);
}
} // end foreach monitor
rtspServer->Stop();
sessions.clear();
Image::Deinitialise();
logTerm();
// Stop the queue before closing the connection it writes through; every other
// daemon already does. Outside zmDbClose because that holds db_mutex, and the
// queue thread wants that same mutex to drain, so joining it from in there
// would deadlock.
dbQueue.stop();
zmDbClose();
return 0;
}