MythTV master
ExternalStreamHandler.cpp
Go to the documentation of this file.
1// -*- Mode: c++ -*-
2
3#include <QtGlobal>
4#if QT_VERSION >= QT_VERSION_CHECK(6,5,0)
5#include <QtSystemDetection>
6#endif
7
8// POSIX headers
9#include <thread>
10#include <iostream>
11#include <fcntl.h>
12#include <unistd.h>
13#include <algorithm>
14#ifndef Q_OS_WINDOWS
15#include <poll.h>
16#include <sys/ioctl.h>
17#endif
18
19#ifdef Q_OS_ANDROID
20#include <sys/wait.h>
21#endif
22
23// Qt headers
24#include <QString>
25#include <QFile>
26#include <QJsonDocument>
27#include <QJsonObject>
28
29// MythTV headers
30#include "libmythbase/mythconfig.h"
34
35#include "ExternalChannel.h"
37//#include "ThreadedFileWriter.h"
38#include "cardutil.h"
39#include "dtvsignalmonitor.h"
40#include "mpeg/mpegstreamdata.h"
42
43#define LOC QString("ExternSH[%1](%2): ").arg(m_inputId).arg(m_loc)
44
45ExternIO::ExternIO(const QString & app,
46 const QStringList & args)
47 : m_app(QFileInfo(app)),
48#if QT_VERSION < QT_VERSION_CHECK(6,0,0)
49 m_status(&m_statusBuf, QIODevice::ReadWrite)
50#else
51 m_status(&m_statusBuf, QIODeviceBase::ReadWrite)
52#endif
53{
54 if (!m_app.exists())
55 {
56 m_error = QString("ExternIO: '%1' does not exist.").arg(app);
57 return;
58 }
59 if (!m_app.isReadable() || !m_app.isFile())
60 {
61 m_error = QString("ExternIO: '%1' is not readable.")
62 .arg(m_app.canonicalFilePath());
63 return;
64 }
65 if (!m_app.isExecutable())
66 {
67 m_error = QString("ExternIO: '%1' is not executable.")
68 .arg(m_app.canonicalFilePath());
69 return;
70 }
71
72 m_args = args;
73 m_args.prepend(m_app.baseName());
74
75 m_status.setString(&m_statusBuf);
76}
77
79{
83
84 // waitpid(m_pid, &status, 0);
85 delete[] m_buffer;
86}
87
88bool ExternIO::Ready([[maybe_unused]] int fd,
89 [[maybe_unused]] std::chrono::milliseconds timeout,
90 [[maybe_unused]] const QString & what)
91{
92#ifndef Q_OS_WINDOWS
93 std::array<struct pollfd,2> m_poll {};
94
95 m_poll[0].fd = fd;
96 m_poll[0].events = POLLIN | POLLPRI;
97 int ret = poll(m_poll.data(), 1, timeout.count());
98
99 if (m_poll[0].revents & POLLHUP)
100 {
101 m_error = what + " poll eof (POLLHUP)";
102 return false;
103 }
104 if (m_poll[0].revents & POLLNVAL)
105 {
106 LOG(VB_GENERAL, LOG_ERR, "poll error");
107 return false;
108 }
109 if (m_poll[0].revents & POLLIN)
110 {
111 if (ret > 0)
112 return true;
113
114 if ((EOVERFLOW == errno))
115 m_error = "poll overflow";
116 return false;
117 }
118#endif // !defined( Q_OS_WINDOWS )
119 return false;
120}
121
122int ExternIO::Read(QByteArray & buffer, int maxlen, std::chrono::milliseconds timeout)
123{
124 if (Error())
125 {
126 LOG(VB_RECORD, LOG_ERR,
127 QString("ExternIO::Read: already in error state: '%1'")
128 .arg(m_error));
129 return 0;
130 }
131
132 if (!Ready(m_appOut, timeout, "data"))
133 return 0;
134
135 if (m_bufSize < maxlen)
136 {
137 m_bufSize = maxlen;
138 delete [] m_buffer;
139 m_buffer = new char[m_bufSize];
140 }
141
142 int len = read(m_appOut, m_buffer, maxlen);
143
144 if (len < 0)
145 {
146 if (errno == EAGAIN)
147 {
148 if (++m_errCnt > kMaxErrorCnt)
149 {
150 m_error = "Failed to read from External Recorder: " + ENO;
151 LOG(VB_RECORD, LOG_WARNING,
152 "External Recorder not ready. Giving up.");
153 }
154 else
155 {
156 LOG(VB_RECORD, LOG_WARNING,
157 QString("External Recorder not ready. Will retry (%1/%2).")
158 .arg(m_errCnt).arg(kMaxErrorCnt));
159 std::this_thread::sleep_for(100ms);
160 }
161 }
162 else
163 {
164 m_error = "Failed to read from External Recorder: " + ENO;
165 LOG(VB_RECORD, LOG_ERR, m_error);
166 }
167 }
168 else
169 {
170 m_errCnt = 0;
171 }
172
173 if (len == 0)
174 return 0;
175
176 buffer.append(m_buffer, len);
177
178 LOG(VB_RECORD, LOG_DEBUG,
179 QString("ExternIO::Read '%1' bytes, buffer size %2")
180 .arg(len).arg(buffer.size()));
181
182 return len;
183}
184
185QByteArray ExternIO::GetStatus(std::chrono::milliseconds timeout)
186{
187 if (Error())
188 {
189 LOG(VB_RECORD, LOG_ERR,
190 QString("ExternIO::GetStatus: already in error state: '%1'")
191 .arg(m_error));
192 return {};
193 }
194
195 std::chrono::milliseconds waitfor = m_status.atEnd() ? timeout : 0ms;
196 if (Ready(m_appErr, waitfor, "status"))
197 {
198 std::array<char,2048> buffer {};
199 int len = read(m_appErr, buffer.data(), buffer.size());
200 m_status << QString::fromLatin1(buffer.data(), len);
201 }
202
203 if (m_status.atEnd())
204 return {};
205
206 QString msg = m_status.readLine();
207
208 LOG(VB_RECORD, LOG_DEBUG, QString("ExternIO::GetStatus '%1'")
209 .arg(msg));
210
211 return msg.toUtf8();
212}
213
214int ExternIO::Write(const QByteArray & buffer)
215{
216 if (Error())
217 {
218 LOG(VB_RECORD, LOG_ERR,
219 QString("ExternIO::Write: already in error state: '%1'")
220 .arg(m_error));
221 return -1;
222 }
223
224 LOG(VB_RECORD, LOG_DEBUG, QString("ExternIO::Write('%1')")
225 .arg(QString(buffer).simplified()));
226
227 int len = write(m_appIn, buffer.constData(), buffer.size());
228 if (len != buffer.size())
229 {
230 if (len > 0)
231 {
232 LOG(VB_RECORD, LOG_WARNING,
233 QString("ExternIO::Write: only wrote %1 of %2 bytes '%3'")
234 .arg(len).arg(buffer.size()).arg(QString(buffer)));
235 }
236 else
237 {
238 m_error = QString("ExternIO: Failed to write '%1' to app's stdin: ")
239 .arg(QString(buffer)) + ENO;
240 return -1;
241 }
242 }
243
244 return len;
245}
246
248{
249 LOG(VB_RECORD, LOG_INFO, QString("ExternIO::Run()"));
250
251 Fork();
252 GetStatus(10ms);
253
254 return true;
255}
256
257/* Return true if the process is not, or is no longer running */
258bool ExternIO::KillIfRunning([[maybe_unused]] const QString & cmd)
259{
260#ifdef Q_OS_BSD4
261 return false;
262#elif defined( Q_OS_WINDOWS )
263 return false;
264#else
265 QString grp = QString("pgrep -x -f -- \"%1\" 2>&1 > /dev/null").arg(cmd);
266 QString kil = QString("pkill --signal 15 -x -f -- \"%1\" 2>&1 > /dev/null")
267 .arg(cmd);
268
269 int res_grp = system(grp.toUtf8().constData());
270 if (WEXITSTATUS(res_grp) == 1)
271 {
272 LOG(VB_RECORD, LOG_DEBUG, QString("'%1' not running.").arg(cmd));
273 return true;
274 }
275
276 LOG(VB_RECORD, LOG_WARNING, QString("'%1' already running, killing...")
277 .arg(cmd));
278 int res_kil = system(kil.toUtf8().constData());
279 if (WEXITSTATUS(res_kil) == 1)
280 LOG(VB_GENERAL, LOG_WARNING, QString("'%1' failed: %2")
281 .arg(kil, ENO));
282
283 res_grp = system(grp.toUtf8().constData());
284 if (WEXITSTATUS(res_grp) == 1)
285 {
286 LOG(WEXITSTATUS(res_kil) == 0 ? VB_RECORD : VB_GENERAL, LOG_WARNING,
287 QString("'%1' terminated.").arg(cmd));
288 return true;
289 }
290
291 std::this_thread::sleep_for(50ms);
292
293 kil = QString("pkill --signal 9 -x -f \"%1\" 2>&1 > /dev/null").arg(cmd);
294 res_kil = system(kil.toUtf8().constData());
295 if (WEXITSTATUS(res_kil) > 0)
296 LOG(VB_GENERAL, LOG_WARNING, QString("'%1' failed: %2")
297 .arg(kil, ENO));
298
299 res_grp = system(grp.toUtf8().constData());
300 LOG(WEXITSTATUS(res_kil) == 0 ? VB_RECORD : VB_GENERAL, LOG_WARNING,
301 QString("'%1' %2.")
302 .arg(cmd, WEXITSTATUS(res_grp) == 0 ? "sill running" : "terminated"));
303
304 return (WEXITSTATUS(res_grp) != 0);
305#endif
306}
307
309{
310#ifndef Q_OS_WINDOWS
311 if (Error())
312 {
313 LOG(VB_RECORD, LOG_INFO, QString("ExternIO in bad state: '%1'")
314 .arg(m_error));
315 return;
316 }
317
318 QString full_command = QString("%1").arg(m_args.join(" "));
319
320 if (!KillIfRunning(full_command))
321 {
322 // Give it one more chance.
323 std::this_thread::sleep_for(50ms);
324 if (!KillIfRunning(full_command))
325 {
326 m_error = QString("Unable to kill existing '%1'.")
327 .arg(full_command);
328 LOG(VB_GENERAL, LOG_ERR, m_error);
329 return;
330 }
331 }
332
333
334 LOG(VB_RECORD, LOG_INFO, QString("ExternIO::Fork '%1'").arg(full_command));
335
336 std::array<int,2> in = {-1, -1};
337 std::array<int,2> out = {-1, -1};
338 std::array<int,2> err = {-1, -1};
339
340 if (pipe(in.data()) < 0)
341 {
342 m_error = "pipe(in) failed: " + ENO;
343 return;
344 }
345 if (pipe(out.data()) < 0)
346 {
347 m_error = "pipe(out) failed: " + ENO;
348 close(in[0]);
349 close(in[1]);
350 return;
351 }
352 if (pipe(err.data()) < 0)
353 {
354 m_error = "pipe(err) failed: " + ENO;
355 close(in[0]);
356 close(in[1]);
357 close(out[0]);
358 close(out[1]);
359 return;
360 }
361
362 m_pid = fork();
363 if (m_pid < 0)
364 {
365 // Failed
366 m_error = "fork() failed: " + ENO;
367 return;
368 }
369 if (m_pid > 0)
370 {
371 // Parent
372 close(in[0]);
373 close(out[1]);
374 close(err[1]);
375 m_appIn = in[1];
376 m_appOut = out[0];
377 m_appErr = err[0];
378
379 bool error = false;
380 error = (fcntl(m_appIn, F_SETFL, O_NONBLOCK) == -1);
381 error |= (fcntl(m_appOut, F_SETFL, O_NONBLOCK) == -1);
382 error |= (fcntl(m_appErr, F_SETFL, O_NONBLOCK) == -1);
383
384 if (error)
385 {
386 LOG(VB_GENERAL, LOG_WARNING,
387 "ExternIO::Fork(): Failed to set O_NONBLOCK for FD: " + ENO);
388 std::this_thread::sleep_for(2s);
390 }
391
392 LOG(VB_RECORD, LOG_INFO, "Spawned");
393 return;
394 }
395
396 // Child
397 close(in[1]);
398 close(out[0]);
399 close(err[0]);
400 if (dup2( in[0], 0) < 0)
401 {
402 std::cerr << "dup2(stdin) failed: " << strerror(errno);
404 }
405 else if (dup2(out[1], 1) < 0)
406 {
407 std::cerr << "dup2(stdout) failed: " << strerror(errno);
409 }
410 else if (dup2(err[1], 2) < 0)
411 {
412 std::cerr << "dup2(stderr) failed: " << strerror(errno);
414 }
415
416 /* Close all open file descriptors except stdin/stdout/stderr */
417#if HAVE_CLOSE_RANGE
418 close_range(3, sysconf(_SC_OPEN_MAX) - 1, 0);
419#else
420 for (int i = sysconf(_SC_OPEN_MAX) - 1; i > 2; --i)
421 close(i);
422#endif
423
424 /* Set the process group id to be the same as the pid of this
425 * child process. This ensures that any subprocesses launched by this
426 * process can be killed along with the process itself. */
427 if (setpgid(0,0) < 0)
428 {
429 std::cerr << "ExternIO: "
430 << "setpgid() failed: "
431 << strerror(errno) << '\n';
432 }
433
434 /* run command */
435 char *command = strdup(m_app.canonicalFilePath()
436 .toUtf8().constData());
437 if (command == nullptr)
438 {
439 std::cerr << "ExternIO: strdup() failed: " << strerror(errno) << '\n';
441 }
442
443 // Copy QStringList to char**
444 char **arguments = new char*[m_args.size() + 1];
445 for (int i = 0; i < m_args.size(); ++i)
446 {
447 int len = m_args[i].size() + 1;
448 arguments[i] = new char[len];
449 memcpy(arguments[i], m_args[i].toStdString().c_str(), len);
450 }
451 arguments[m_args.size()] = nullptr;
452
453 if (execv(command, arguments) < 0)
454 {
455 // Can't use LOG due to locking fun.
456 std::cerr << "ExternIO: "
457 << "execv() failed: "
458 << strerror(errno) << '\n';
459 }
460 else
461 {
462 std::cerr << "ExternIO: "
463 << "execv() should not be here?: "
464 << strerror(errno) << '\n';
465 }
466
467#endif // !defined( Q_OS_WINDOWS )
468
469 /* Failed to exec */
470 _exit(GENERIC_EXIT_DAEMONIZING_ERROR); // this exit is ok
471}
472
473
474QMap<int, ExternalStreamHandler*> ExternalStreamHandler::s_handlers;
477
479 int inputid, int majorid)
480{
481 QMutexLocker locker(&s_handlersLock);
482
483 QMap<int, ExternalStreamHandler*>::iterator it = s_handlers.find(majorid);
484
485 if (it == s_handlers.end())
486 {
487 auto *newhandler = new ExternalStreamHandler(devname, inputid, majorid);
488 s_handlers[majorid] = newhandler;
489 s_handlersRefCnt[majorid] = 1;
490
491 LOG(VB_RECORD, LOG_INFO,
492 QString("ExternSH[%1:%2]: Creating new stream handler for %3 "
493 "(1 in use)")
494 .arg(inputid).arg(majorid).arg(devname));
495 }
496 else
497 {
498 ++s_handlersRefCnt[majorid];
499 uint rcount = s_handlersRefCnt[majorid];
500 LOG(VB_RECORD, LOG_INFO,
501 QString("ExternSH[%1:%2]: Using existing stream handler for %3")
502 .arg(inputid).arg(majorid).arg(devname) +
503 QString(" (%1 in use)").arg(rcount));
504 }
505
506 return s_handlers[majorid];
507}
508
510 int inputid)
511{
512 QMutexLocker locker(&s_handlersLock);
513
514 int majorid = ref->m_majorId;
515
516 QMap<int, uint>::iterator rit = s_handlersRefCnt.find(majorid);
517 if (rit == s_handlersRefCnt.end())
518 return;
519
520 QMap<int, ExternalStreamHandler*>::iterator it =
521 s_handlers.find(majorid);
522
523 if (*rit > 1)
524 {
525 ref = nullptr;
526 --(*rit);
527
528 LOG(VB_RECORD, LOG_INFO,
529 QString("ExternSH[%1:%2]: Return handler (%3 still in use)")
530 .arg(inputid).arg(majorid).arg(*rit));
531
532 return;
533 }
534
535 if ((it != s_handlers.end()) && (*it == ref))
536 {
537 LOG(VB_RECORD, LOG_INFO,
538 QString("ExternSH[%1:%2]: Closing handler (0 in use)")
539 .arg(inputid).arg(majorid));
540 delete *it;
541 s_handlers.erase(it);
542 }
543 else
544 {
545 LOG(VB_GENERAL, LOG_ERR,
546 QString("ExternSH[%1:%2]: Error: No handler to return!")
547 .arg(inputid).arg(majorid));
548 }
549
550 s_handlersRefCnt.erase(rit);
551 ref = nullptr;
552}
553
554/*
555 ExternalStreamHandler
556 */
557
559 int inputid,
560 int majorid)
561 : StreamHandler(path, inputid)
562 , m_loc(m_device)
563 , m_majorId(majorid)
564{
565 setObjectName("ExternSH");
566
567 m_args = path.split(' ',Qt::SkipEmptyParts) +
568 logPropagateArgs.split(' ', Qt::SkipEmptyParts);
569 //NOLINTNEXTLINE(cppcoreguidelines-prefer-member-initializer)
570 m_app = m_args.first();
571 m_args.removeFirst();
572
573 // Pass one (and only one) 'quiet'
574 if (!m_args.contains("--quiet") && !m_args.contains("-q"))
575 m_args << "--quiet";
576
577 m_args << "--inputid" << QString::number(majorid);
578 LOG(VB_RECORD, LOG_INFO, LOC + QString("args \"%1\"")
579 .arg(m_args.join(" ")));
580
581 if (!OpenApp())
582 {
583 LOG(VB_GENERAL, LOG_ERR, LOC +
584 QString("Failed to start %1").arg(m_device));
585 }
586}
587
589{
590 return m_streamingCnt.loadAcquire();
591}
592
594{
595 QString result;
596 QString ready_cmd;
597 QByteArray buffer;
598 int sz = 0;
599 uint len = 0;
600 uint read_len = 0;
601 uint restart_cnt = 0;
602 MythTimer status_timer;
603 MythTimer nodata_timer;
604
605 bool good_data = false;
606 uint data_proc_err = 0;
607 uint data_short_err = 0;
608
609 if (!m_io)
610 {
611 LOG(VB_GENERAL, LOG_ERR, LOC +
612 QString("%1 is not running.").arg(m_device));
613 }
614
615 status_timer.start();
616
617 RunProlog();
618
619 LOG(VB_RECORD, LOG_INFO, LOC + "run(): begin");
620
621 SetRunning(true, true, false);
622
623 if (m_pollMode)
624 ready_cmd = "SendBytes";
625 else
626 ready_cmd = "XON";
627
628 uint remainder = 0;
629 while (m_runningDesired && !m_bError)
630 {
631 if (!IsTSOpen())
632 {
633 LOG(VB_RECORD, LOG_WARNING, LOC + "TS not open yet.");
634 std::this_thread::sleep_for(10ms);
635 continue;
636 }
637
638 if (StreamingCount() == 0)
639 {
640 std::this_thread::sleep_for(10ms);
641 continue;
642 }
643
645
646 if (!m_xon || m_pollMode)
647 {
648 if (buffer.size() > TOO_FAST_SIZE)
649 {
650 LOG(VB_RECORD, LOG_WARNING, LOC +
651 "Internal buffer too full to accept more data from "
652 "external application.");
653 }
654 else
655 {
656 if (!ProcessCommand(ready_cmd, result))
657 {
658 if (result.startsWith("ERR"))
659 {
660 LOG(VB_GENERAL, LOG_ERR, LOC +
661 QString("Aborting: %1 -> %2")
662 .arg(ready_cmd, result));
663 m_bError = true;
664 continue;
665 }
666
667 if (restart_cnt++)
668 std::this_thread::sleep_for(20s);
669 if (!RestartStream())
670 {
671 LOG(VB_RECORD, LOG_ERR, LOC +
672 "Failed to restart stream.");
673 m_bError = true;
674 }
675 continue;
676 }
677 m_xon = true;
678 }
679 }
680
681 if (m_xon)
682 {
683 if (status_timer.elapsed() >= 2s)
684 {
685 // Since we may never need to send the XOFF
686 // command, occationally check to see if the
687 // External recorder needs to report an issue.
688 if (Monitor())
689 {
690 if (restart_cnt++)
691 std::this_thread::sleep_for(20s);
692 if (!RestartStream())
693 {
694 LOG(VB_RECORD, LOG_ERR, LOC +
695 "Failed to restart stream.");
696 m_bError = true;
697 }
698 continue;
699 }
700
701 status_timer.restart();
702 }
703
704 if (buffer.size() > TOO_FAST_SIZE)
705 {
706 if (!m_pollMode)
707 {
708 // Data is comming a little too fast, so XOFF
709 // to give us time to process it.
710 if (!ProcessCommand(QString("XOFF"), result))
711 {
712 if (result.startsWith("ERR"))
713 {
714 LOG(VB_GENERAL, LOG_ERR, LOC +
715 QString("Aborting: XOFF -> %2")
716 .arg(result));
717 m_bError = true;
718 }
719 }
720 m_xon = false;
721 }
722 }
723
724 read_len = 0;
725 if (m_io != nullptr)
726 {
727 sz = PACKET_SIZE - remainder;
728 if (sz > 0)
729 read_len = m_io->Read(buffer, sz, 100ms);
730 }
731 }
732 else
733 {
734 read_len = 0;
735 }
736
737 if (read_len == 0)
738 {
739 if (!nodata_timer.isRunning())
740 {
741 nodata_timer.start();
742 }
743 else
744 {
745 if (nodata_timer.elapsed() >= 50s)
746 {
747 LOG(VB_GENERAL, LOG_WARNING, LOC +
748 "No data for 50 seconds, Restarting stream.");
749 if (!RestartStream())
750 {
751 LOG(VB_RECORD, LOG_ERR, LOC +
752 "Failed to restart stream.");
753 m_bError = true;
754 }
755 nodata_timer.stop();
756 continue;
757 }
758 }
759
760 std::this_thread::sleep_for(50ms);
761
762 // HLS type streams may only produce data every ~10 seconds
763 if (nodata_timer.elapsed() < 12s && buffer.size() < TS_PACKET_SIZE)
764 continue;
765 }
766 else
767 {
768 nodata_timer.stop();
769 restart_cnt = 0;
770 }
771
772 if (m_io == nullptr)
773 {
774 LOG(VB_GENERAL, LOG_ERR, LOC + "I/O thread has disappeared!");
775 m_bError = true;
776 break;
777 }
778 if (m_io->Error())
779 {
780 LOG(VB_GENERAL, LOG_ERR, LOC +
781 QString("Fatal Error from External Recorder: %1")
782 .arg(m_io->ErrorString()));
783 CloseApp();
784 m_bError = true;
785 break;
786 }
787
788 len = remainder = buffer.size();
789
790 if (len == 0)
791 continue;
792
793 if (len < TS_PACKET_SIZE)
794 {
795 if (m_xon && data_short_err++ == 0)
796 LOG(VB_RECORD, LOG_INFO, LOC + "Waiting for a full TS packet.");
797 std::this_thread::sleep_for(50us);
798 continue;
799 }
800 if (data_short_err)
801 {
802 if (data_short_err > 1)
803 {
804 LOG(VB_RECORD, LOG_INFO, LOC +
805 QString("Waited for a full TS packet %1 times.")
806 .arg(data_short_err));
807 }
808 data_short_err = 0;
809 }
810
811 if (!m_streamLock.tryLock())
812 continue;
813
814 if (!m_listenerLock.tryLock())
815 continue;
816
817 for (auto sit = m_streamDataList.cbegin();
818 sit != m_streamDataList.cend(); ++sit)
819 {
820 remainder = sit.key()->ProcessData
821 (reinterpret_cast<const uint8_t *>
822 (buffer.constData()), buffer.size());
823 }
824
825 m_listenerLock.unlock();
826
827 if (m_replay)
828 {
829 m_replayBuffer += buffer.left(len - remainder);
830 if (m_replayBuffer.size() > (50 * PACKET_SIZE))
831 {
832 m_replayBuffer.remove(0, len - remainder);
833 LOG(VB_RECORD, LOG_WARNING, LOC +
834 QString("Replay size truncated to %1 bytes")
835 .arg(m_replayBuffer.size()));
836 }
837 }
838
839 m_streamLock.unlock();
840
841 if (remainder == 0)
842 {
843 buffer.clear();
844 good_data = (len != 0U);
845 }
846 else if (len > remainder) // leftover bytes
847 {
848 buffer.remove(0, len - remainder);
849 good_data = (len != 0U);
850 }
851 else if (len == remainder)
852 {
853 good_data = false;
854 }
855
856 if (good_data)
857 {
858 if (data_proc_err)
859 {
860 if (data_proc_err > 1)
861 {
862 LOG(VB_RECORD, LOG_WARNING, LOC +
863 QString("Failed to process the data received %1 times.")
864 .arg(data_proc_err));
865 }
866 data_proc_err = 0;
867 }
868 }
869 else
870 {
871 if (data_proc_err++ == 0)
872 {
873 LOG(VB_RECORD, LOG_WARNING, LOC +
874 "Failed to process the data received");
875 }
876 }
877 }
878
879 LOG(VB_RECORD, LOG_INFO, LOC + "run(): " +
880 QString("%1 shutdown").arg(m_bError ? "Error" : "Normal"));
881
883 SetRunning(false, true, false);
884
885 LOG(VB_RECORD, LOG_INFO, LOC + "run(): " + "end");
886
887 RunEpilog();
888}
889
891{
892 QString result;
893
894 if (ProcessCommand("APIVersion?", result, 10s))
895 {
896 QStringList tokens = result.split(':', Qt::SkipEmptyParts);
897 if (tokens.size() > 1)
898 m_apiVersion = tokens[1].toUInt();
899 m_apiVersion = std::min(m_apiVersion, static_cast<int>(MAX_API_VERSION));
900 if (m_apiVersion < 1)
901 {
902 LOG(VB_RECORD, LOG_ERR, LOC +
903 QString("Bad response to 'APIVersion?' - '%1'. "
904 "Expecting 1, 2 or 3").arg(result));
905 m_apiVersion = 1;
906 }
907
908 ProcessCommand(QString("APIVersion:%1").arg(m_apiVersion), result);
909 return true;
910 }
911
912 return false;
913}
914
916{
917 if (m_apiVersion > 1)
918 {
919 QString result;
920
921 if (ProcessCommand("Description?", result))
922 m_loc = result.mid(3);
923 else
924 m_loc = m_device;
925 }
926
927 return m_loc;
928}
929
931{
932 {
933 QMutexLocker locker(&m_ioLock);
934
935 if (m_io)
936 {
937 LOG(VB_RECORD, LOG_WARNING, LOC + "OpenApp: already open!");
938 return true;
939 }
940
941 m_io = new ExternIO(m_app, m_args);
942
943 if (m_io == nullptr)
944 {
945 LOG(VB_GENERAL, LOG_ERR, LOC + "ExternIO failed: " + ENO);
946 m_bError = true;
947 }
948 else
949 {
950 LOG(VB_RECORD, LOG_INFO, LOC + QString("Spawn '%1'").arg(m_device));
951 m_io->Run();
952 if (m_io->Error())
953 {
954 LOG(VB_GENERAL, LOG_ERR,
955 "Failed to start External Recorder: " + m_io->ErrorString());
956 delete m_io;
957 m_io = nullptr;
958 m_bError = true;
959 return false;
960 }
961 }
962 }
963
964 QString result;
965
966 if (!SetAPIVersion())
967 {
968 // Try again using API version 2
969 m_apiVersion = 2;
970 if (!SetAPIVersion())
971 m_apiVersion = 1;
972 }
973
974 if (!IsAppOpen())
975 {
976 LOG(VB_RECORD, LOG_ERR, LOC + "Application is not responding.");
977 m_bError = true;
978 return false;
979 }
980
982
983 // Gather capabilities
984 if (!ProcessCommand("HasTuner?", result))
985 {
986 LOG(VB_RECORD, LOG_ERR, LOC +
987 QString("Bad response to 'HasTuner?' - '%1'").arg(result));
988 m_bError = true;
989 return false;
990 }
991 m_hasTuner = result.startsWith("OK:Yes");
992
993 if (!ProcessCommand("HasPictureAttributes?", result))
994 {
995 LOG(VB_RECORD, LOG_ERR, LOC +
996 QString("Bad response to 'HasPictureAttributes?' - '%1'")
997 .arg(result));
998 m_bError = true;
999 return false;
1000 }
1001 m_hasPictureAttributes = result.startsWith("OK:Yes");
1002
1003 /* Operate in "poll" or "xon/xoff" mode */
1004 m_pollMode = ProcessCommand("FlowControl?", result) &&
1005 result.startsWith("OK:Poll");
1006
1007 LOG(VB_RECORD, LOG_INFO, LOC + "App opened successfully");
1008 LOG(VB_RECORD, LOG_INFO, LOC +
1009 QString("Capabilities: tuner(%1) "
1010 "Picture attributes(%2) "
1011 "Flow control(%3)")
1012 .arg(m_hasTuner ? "yes" : "no",
1013 m_hasPictureAttributes ? "yes" : "no",
1014 m_pollMode ? "Polling" : "XON/XOFF")
1015 );
1016
1017 /* Let the external app know how many bytes will read without blocking */
1018 ProcessCommand(QString("BlockSize:%1").arg(PACKET_SIZE), result);
1019
1020 return true;
1021}
1022
1024{
1025 if (m_io == nullptr)
1026 {
1027 LOG(VB_RECORD, LOG_WARNING, LOC +
1028 "WARNING: Unable to communicate with external app.");
1029 return false;
1030 }
1031
1032 QString result;
1033 return ProcessCommand("Version?", result, 10s);
1034}
1035
1037{
1038 if (m_tsOpen)
1039 return true;
1040
1041 QString result;
1042
1043 if (!ProcessCommand("IsOpen?", result))
1044 return false;
1045
1046 m_tsOpen = true;
1047 return m_tsOpen;
1048}
1049
1051{
1052 m_ioLock.lock();
1053 if (m_io)
1054 {
1055 QString result;
1056
1057 LOG(VB_RECORD, LOG_INFO, LOC + "CloseRecorder");
1058 m_ioLock.unlock();
1059 ProcessCommand("CloseRecorder", result, 10s);
1060 m_ioLock.lock();
1061
1062 if (!result.startsWith("OK"))
1063 {
1064 LOG(VB_RECORD, LOG_INFO, LOC +
1065 "CloseRecorder failed, sending kill.");
1066
1067 QString full_command = QString("%1").arg(m_args.join(" "));
1068
1069 if (!ExternIO::KillIfRunning(full_command))
1070 {
1071 // Give it one more chance.
1072 std::this_thread::sleep_for(50ms);
1073 if (!ExternIO::KillIfRunning(full_command))
1074 {
1075 LOG(VB_GENERAL, LOG_ERR,
1076 QString("Unable to kill existing '%1'.")
1077 .arg(full_command));
1078 return;
1079 }
1080 }
1081 }
1082 delete m_io;
1083 m_io = nullptr;
1084 }
1085 m_ioLock.unlock();
1086}
1087
1089{
1090 bool streaming = (StreamingCount() > 0);
1091
1092 LOG(VB_RECORD, LOG_WARNING, LOC + "Restarting stream.");
1093 m_damaged = true;
1094
1095 if (streaming)
1096 StopStreaming();
1097
1098 std::this_thread::sleep_for(1s);
1099
1100 if (streaming)
1102
1103 return true;
1104}
1105
1107{
1108 if (m_replay)
1109 {
1110 QString result;
1111
1112 // Let the external app know that we could be busy for a little while
1113 if (!m_pollMode)
1114 {
1115 ProcessCommand(QString("XOFF"), result);
1116 m_xon = false;
1117 }
1118
1119 /* If the input is not a 'broadcast' it may only have one
1120 * copy of the SPS right at the beginning of the stream,
1121 * so make sure we don't miss it!
1122 */
1123 QMutexLocker listen_lock(&m_listenerLock);
1124
1125 if (!m_streamDataList.empty())
1126 {
1127 for (auto sit = m_streamDataList.cbegin();
1128 sit != m_streamDataList.cend(); ++sit)
1129 {
1130 sit.key()->ProcessData(reinterpret_cast<const uint8_t *>
1131 (m_replayBuffer.constData()),
1132 m_replayBuffer.size());
1133 }
1134 }
1135 LOG(VB_RECORD, LOG_INFO, LOC + QString("Replayed %1 bytes")
1136 .arg(m_replayBuffer.size()));
1137 m_replayBuffer.clear();
1138 m_replay = false;
1139
1140 // Let the external app know that we are ready
1141 if (!m_pollMode)
1142 {
1143 if (ProcessCommand(QString("XON"), result))
1144 m_xon = true;
1145 }
1146 }
1147}
1148
1150{
1151 QString result;
1152
1153 QMutexLocker locker(&m_streamLock);
1154
1156
1157 LOG(VB_RECORD, LOG_INFO, LOC +
1158 QString("StartStreaming with %1 current listeners")
1159 .arg(StreamingCount()));
1160
1161 if (!IsAppOpen())
1162 {
1163 LOG(VB_GENERAL, LOG_ERR, LOC + "External Recorder not started.");
1164 return false;
1165 }
1166
1167 if (StreamingCount() == 0)
1168 {
1169 if (!ProcessCommand("StartStreaming", result, 15s))
1170 {
1171 LogLevel_t level = LOG_ERR;
1172 if (result.startsWith("warn", Qt::CaseInsensitive))
1173 level = LOG_WARNING;
1174 else
1175 m_bError = true;
1176
1177 LOG(VB_GENERAL, level, LOC + QString("StartStreaming failed: '%1'")
1178 .arg(result));
1179
1180 return false;
1181 }
1182 LOG(VB_RECORD, LOG_INFO, LOC + "Streaming started");
1183 }
1184 else
1185 {
1186 LOG(VB_RECORD, LOG_INFO, LOC + "Already streaming");
1187 }
1188 m_recording = recording;
1189
1190 m_streamingCnt.ref();
1191
1192 LOG(VB_RECORD, LOG_INFO, LOC +
1193 QString("StartStreaming %1 listeners")
1194 .arg(StreamingCount()));
1195
1196 return true;
1197}
1198
1200{
1201 QMutexLocker locker(&m_streamLock);
1202
1203 LOG(VB_RECORD, LOG_INFO, LOC +
1204 QString("StopStreaming %1 listeners")
1205 .arg(StreamingCount()));
1206
1207 if (StreamingCount() == 0)
1208 {
1209 LOG(VB_RECORD, LOG_INFO, LOC +
1210 "StopStreaming requested, but we are not streaming!");
1211 return true;
1212 }
1213
1214 if (m_streamingCnt.deref())
1215 {
1216 LOG(VB_RECORD, LOG_INFO, LOC +
1217 QString("StopStreaming delayed, still have %1 listeners")
1218 .arg(StreamingCount()));
1219 return true;
1220 }
1221
1222 LOG(VB_RECORD, LOG_INFO, LOC + "StopStreaming");
1223
1224 if (!m_pollMode && m_xon)
1225 {
1226 QString result;
1227 ProcessCommand(QString("XOFF"), result);
1228 m_xon = false;
1229 }
1230
1231 if (!IsAppOpen())
1232 {
1233 LOG(VB_GENERAL, LOG_ERR, LOC + "External Recorder not started.");
1234 return false;
1235 }
1236
1237 QString result;
1238 if (!ProcessCommand("StopStreaming", result, 10s))
1239 {
1240 LogLevel_t level = LOG_ERR;
1241 if (result.startsWith("warn", Qt::CaseInsensitive))
1242 level = LOG_WARNING;
1243 else
1244 m_bError = true;
1245
1246 LOG(VB_GENERAL, level, LOC + QString("StopStreaming: '%1'")
1247 .arg(result));
1248
1249 return false;
1250 }
1251
1252 m_recording = false;
1253 PurgeBuffer();
1254 LOG(VB_RECORD, LOG_INFO, LOC + "Streaming stopped");
1255
1256 return true;
1257}
1258
1260 QString & result,
1261 std::chrono::milliseconds timeout,
1262 uint retry_cnt)
1263{
1264 QMutexLocker locker(&m_processLock);
1265
1266 if (m_apiVersion == 3)
1267 {
1268 QVariantMap vcmd;
1269 QVariantMap vresult;
1270 QByteArray response;
1271 QStringList tokens = cmd.split(':');
1272 vcmd["command"] = tokens[0];
1273 if (tokens.size() > 1)
1274 vcmd["value"] = tokens[1];
1275
1276 LOG(VB_RECORD, LOG_DEBUG, LOC +
1277 QString("Arguments: %1").arg(tokens.join("\n")));
1278
1279 bool r = ProcessJson(vcmd, vresult, response, timeout, retry_cnt);
1280 result = QString("%1:%2").arg(vresult["status"].toString(),
1281 vresult["message"].toString());
1282 return r;
1283 }
1284 if (m_apiVersion == 2)
1285 return ProcessVer2(cmd, result, timeout, retry_cnt);
1286 if (m_apiVersion == 1)
1287 return ProcessVer1(cmd, result, timeout, retry_cnt);
1288
1289 LOG(VB_RECORD, LOG_ERR, LOC +
1290 QString("Invalid API version %1. Expected 1 or 2").arg(m_apiVersion));
1291 return false;
1292}
1293
1294bool ExternalStreamHandler::ProcessVer1(const QString & cmd,
1295 QString & result,
1296 std::chrono::milliseconds timeout,
1297 uint retry_cnt)
1298{
1299 LOG(VB_RECORD, LOG_DEBUG, LOC + QString("ProcessVer1('%1')")
1300 .arg(cmd));
1301
1302 for (uint cnt = 0; cnt < retry_cnt; ++cnt)
1303 {
1304 QMutexLocker locker(&m_ioLock);
1305
1306 if (!m_io)
1307 {
1308 LOG(VB_RECORD, LOG_ERR, LOC + "External I/O not ready!");
1309 return false;
1310 }
1311
1312 QByteArray buf = cmd.toUtf8();
1313 buf += '\n';
1314
1315 if (m_io->Error())
1316 {
1317 LOG(VB_GENERAL, LOG_ERR, LOC + "External Recorder in bad state: " +
1318 m_io->ErrorString());
1319 return false;
1320 }
1321
1322 /* Try to keep in sync, if External app was too slow in responding
1323 * to previous query, consume the response before sending new query */
1324 m_io->GetStatus(0ms);
1325
1326 /* Send new query */
1327 m_io->Write(buf);
1328
1330 while (timer.elapsed() < timeout)
1331 {
1332 result = m_io->GetStatus(timeout);
1333 if (m_io->Error())
1334 {
1335 LOG(VB_GENERAL, LOG_ERR, LOC +
1336 "Failed to read from External Recorder: " +
1337 m_io->ErrorString());
1338 m_bError = true;
1339 return false;
1340 }
1341
1342 // Out-of-band error message
1343 if (result.startsWith("STATUS:ERR") ||
1344 result.startsWith("0:STATUS:ERR"))
1345 {
1346 LOG(VB_RECORD, LOG_ERR, LOC + result);
1347 result.remove(0, result.indexOf(":ERR") + 1);
1348 return false;
1349 }
1350 // STATUS message are "out of band".
1351 // Ignore them while waiting for a responds to a command
1352 if (!result.startsWith("STATUS") && !result.startsWith("0:STATUS"))
1353 break;
1354 LOG(VB_RECORD, LOG_INFO, LOC +
1355 QString("Ignoring response '%1'").arg(result));
1356 }
1357
1358 if (result.size() < 1)
1359 {
1360 LOG(VB_GENERAL, LOG_WARNING, LOC +
1361 QString("External Recorder did not respond to '%1'").arg(cmd));
1362 }
1363 else
1364 {
1365 bool okay = result.startsWith("OK");
1366 if (okay || result.startsWith("WARN") || result.startsWith("ERR"))
1367 {
1368 LogLevel_t level = LOG_INFO;
1369
1370 m_ioErrCnt = 0;
1371 if (!okay)
1372 level = LOG_WARNING;
1373 else if (cmd.startsWith("SendBytes"))
1374 level = LOG_DEBUG;
1375
1376 LOG(VB_RECORD, level,
1377 LOC + QString("ProcessCommand('%1') = '%2' took %3ms %4")
1378 .arg(cmd, result,
1379 QString::number(timer.elapsed().count()),
1380 okay ? "" : "<-- NOTE"));
1381
1382 return okay;
1383 }
1384 LOG(VB_GENERAL, LOG_WARNING, LOC +
1385 QString("External Recorder invalid response to '%1': '%2'")
1386 .arg(cmd, result));
1387 }
1388
1389 if (++m_ioErrCnt > 10)
1390 {
1391 LOG(VB_GENERAL, LOG_ERR, LOC + "Too many I/O errors.");
1392 m_bError = true;
1393 break;
1394 }
1395 }
1396
1397 return false;
1398}
1399
1400bool ExternalStreamHandler::ProcessVer2(const QString & command,
1401 QString & result,
1402 std::chrono::milliseconds timeout,
1403 uint retry_cnt)
1404{
1405 QString status;
1406 QString raw;
1407
1408 for (uint cnt = 0; cnt < retry_cnt; ++cnt)
1409 {
1410 QString cmd = QString("%1:%2").arg(++m_serialNo).arg(command);
1411
1412 LOG(VB_RECORD, LOG_DEBUG, LOC + QString("ProcessVer2('%1') serial(%2)")
1413 .arg(cmd).arg(m_serialNo));
1414
1415 QMutexLocker locker(&m_ioLock);
1416
1417 if (!m_io)
1418 {
1419 LOG(VB_RECORD, LOG_ERR, LOC + "External I/O not ready!");
1420 return false;
1421 }
1422
1423 QByteArray buf = cmd.toUtf8();
1424 buf += '\n';
1425
1426 if (m_io->Error())
1427 {
1428 LOG(VB_GENERAL, LOG_ERR, LOC + "External Recorder in bad state: " +
1429 m_io->ErrorString());
1430 return false;
1431 }
1432
1433 /* Send query */
1434 m_io->Write(buf);
1435
1436 QStringList tokens;
1437
1439 while (timer.elapsed() < timeout)
1440 {
1441 result = m_io->GetStatus(timeout);
1442 if (m_io->Error())
1443 {
1444 LOG(VB_GENERAL, LOG_ERR, LOC +
1445 "Failed to read from External Recorder: " +
1446 m_io->ErrorString());
1447 m_bError = true;
1448 return false;
1449 }
1450
1451 if (!result.isEmpty())
1452 {
1453 raw = result;
1454 tokens = result.split(':', Qt::SkipEmptyParts);
1455
1456 // Look for result with the serial number of this query
1457 if (tokens.size() > 1 && tokens[0].toUInt() >= m_serialNo)
1458 break;
1459
1460 /* Other messages are "out of band" */
1461
1462 // Check for error message missing serial#
1463 if (tokens[0].startsWith("ERR"))
1464 break;
1465
1466 // Remove serial#
1467 tokens.removeFirst();
1468 result = tokens.join(':');
1469 bool err = (tokens.size() > 1 && tokens[1].startsWith("ERR"));
1470 LOG(VB_RECORD, (err ? LOG_WARNING : LOG_INFO), LOC + raw);
1471 if (err)
1472 {
1473 // Remove "STATUS"
1474 tokens.removeFirst();
1475 result = tokens.join(':');
1476 return false;
1477 }
1478 }
1479 }
1480
1481 if (timer.elapsed() >= timeout)
1482 {
1483 LOG(VB_RECORD, LOG_ERR, LOC +
1484 QString("ProcessVer2: Giving up waiting for response for "
1485 "command '%2'").arg(cmd));
1486 }
1487 else if (tokens.size() < 2)
1488 {
1489 LOG(VB_RECORD, LOG_ERR, LOC +
1490 QString("Did not receive a valid response "
1491 "for command '%1', received '%2'").arg(cmd, result));
1492 }
1493 else if (tokens[0].toUInt() > m_serialNo)
1494 {
1495 LOG(VB_RECORD, LOG_ERR, LOC +
1496 QString("ProcessVer2: Looking for serial no %1, "
1497 "but received %2 for command '%2'")
1498 .arg(QString::number(m_serialNo), tokens[0], cmd));
1499 }
1500 else
1501 {
1502 tokens.removeFirst();
1503 status = tokens[0].trimmed();
1504 result = tokens.join(':');
1505
1506 bool okay = (status == "OK");
1507 if (okay || status.startsWith("WARN") || status.startsWith("ERR"))
1508 {
1509 LogLevel_t level = LOG_INFO;
1510
1511 m_ioErrCnt = 0;
1512 if (!okay)
1513 level = LOG_WARNING;
1514 else if (command.startsWith("SendBytes") ||
1515 (command.startsWith("TuneStatus") &&
1516 result == "OK:InProgress"))
1517 level = LOG_DEBUG;
1518
1519 LOG(VB_RECORD, level,
1520 LOC + QString("ProcessV2('%1') = '%2' took %3ms %4")
1521 .arg(cmd, result, QString::number(timer.elapsed().count()),
1522 okay ? "" : "<-- NOTE"));
1523
1524 return okay;
1525 }
1526 LOG(VB_GENERAL, LOG_WARNING, LOC +
1527 QString("External Recorder invalid response to '%1': '%2'")
1528 .arg(cmd, result));
1529 }
1530
1531 if (++m_ioErrCnt > 10)
1532 {
1533 LOG(VB_GENERAL, LOG_ERR, LOC + "Too many I/O errors.");
1534 m_bError = true;
1535 break;
1536 }
1537 }
1538
1539 return false;
1540}
1541
1542bool ExternalStreamHandler::ProcessJson(const QVariantMap & vmsg,
1543 QVariantMap & elements,
1544 QByteArray & response,
1545 std::chrono::milliseconds timeout,
1546 uint retry_cnt)
1547{
1548 for (uint cnt = 0; cnt < retry_cnt; ++cnt)
1549 {
1550 QVariantMap query(vmsg);
1551
1552 uint serial = ++m_serialNo;
1553 query["serial"] = serial;
1554 QString cmd = query["command"].toString();
1555
1556 QJsonDocument qdoc;
1557 qdoc = QJsonDocument::fromVariant(query);
1558 QByteArray cmdbuf = qdoc.toJson(QJsonDocument::Compact);
1559
1560 LOG(VB_RECORD, LOG_DEBUG, LOC +
1561 QString("ProcessJson: %1").arg(QString(cmdbuf)));
1562
1563 if (m_io->Error())
1564 {
1565 LOG(VB_GENERAL, LOG_ERR, LOC + "External Recorder in bad state: " +
1566 m_io->ErrorString());
1567 return false;
1568 }
1569
1570 /* Send query */
1571 m_io->Write(cmdbuf);
1572 m_io->Write("\n");
1573
1575 while (timer.elapsed() < timeout)
1576 {
1577 response = m_io->GetStatus(timeout);
1578 if (m_io->Error())
1579 {
1580 LOG(VB_GENERAL, LOG_ERR, LOC +
1581 "Failed to read from External Recorder: " +
1582 m_io->ErrorString());
1583 m_bError = true;
1584 return false;
1585 }
1586
1587 if (!response.isEmpty())
1588 {
1589 QJsonParseError parseError {};
1590 QJsonDocument doc;
1591
1592 doc = QJsonDocument::fromJson(response, &parseError);
1593
1594 if (parseError.error != QJsonParseError::NoError)
1595 {
1596 LOG(VB_GENERAL, LOG_ERR, LOC +
1597 QString("ExternalRecorder returned invalid JSON message: %1: %2\n%3\nfor\n%4")
1598 .arg(parseError.offset)
1599 .arg(parseError.errorString(),
1600 QString(response),
1601 QString(cmdbuf)));
1602 }
1603 else
1604 {
1605 elements = doc.toVariant().toMap();
1606 if (!elements.contains("serial"))
1607 continue;
1608
1609 serial = elements["serial"].toInt();
1610 if (serial >= m_serialNo)
1611 break;
1612
1613 if (elements.contains("status"))
1614 {
1615 LogLevel_t level { LOG_INFO };
1616
1617 if (elements["status"] == "ERR")
1618 level = LOG_ERR;
1619 else if (elements["status"] == "WARN")
1620 level = LOG_WARNING;
1621
1622 LOG(VB_RECORD, level, LOC + QString("%1: %2")
1623 .arg(elements["status"].toString(),
1624 elements["message"].toString()));
1625 }
1626 }
1627 }
1628 }
1629
1630 if (timer.elapsed() >= timeout)
1631 {
1632 LOG(VB_RECORD, LOG_ERR, LOC +
1633 QString("ProcessJson: Giving up waiting for response for "
1634 "command '%2'").arg(QString(cmdbuf)));
1635
1636 }
1637
1638 if (serial > m_serialNo)
1639 {
1640 LOG(VB_RECORD, LOG_ERR, LOC +
1641 QString("ProcessJson: Looking for serial no %1, "
1642 "but received %2 for command '%2'")
1643 .arg(QString::number(m_serialNo))
1644 .arg(serial)
1645 .arg(QString(cmdbuf)));
1646 }
1647 else if (!elements.contains("status"))
1648 {
1649 LOG(VB_RECORD, LOG_ERR, LOC +
1650 QString("ProcessJson: ExternalRecorder 'status' not found in %1")
1651 .arg(QString(response)));
1652 }
1653 else
1654 {
1655 QString status = elements["status"].toString();
1656 bool okay = (status == "OK");
1657 if (okay || status == "WARN" || status == "ERR")
1658 {
1659 LogLevel_t level = LOG_INFO;
1660
1661 m_ioErrCnt = 0;
1662 if (!okay)
1663 level = LOG_WARNING;
1664 else if (cmd == "SendBytes" ||
1665 (cmd == "TuneStatus?" &&
1666 elements["message"] == "InProgress"))
1667 level = LOG_DEBUG;
1668
1669 LOG(VB_RECORD, level,
1670 LOC + QString("ProcessJson('%1') = %2:%3:%4 took %5ms %6")
1671 .arg(QString(cmdbuf))
1672 .arg(elements["serial"].toInt())
1673 .arg(elements["status"].toString(),
1674 elements["message"].toString(),
1675 QString::number(timer.elapsed().count()),
1676 okay ? "" : "<-- NOTE")
1677 );
1678
1679 return okay;
1680 }
1681 LOG(VB_GENERAL, LOG_WARNING, LOC +
1682 QString("External Recorder invalid response to '%1': '%2'")
1683 .arg(QString(cmdbuf),
1684 QString(response)));
1685 }
1686
1687 if (++m_ioErrCnt > 10)
1688 {
1689 LOG(VB_GENERAL, LOG_ERR, LOC + "Too many I/O errors.");
1690 m_bError = true;
1691 break;
1692 }
1693 }
1694
1695 return false;
1696}
1697
1699{
1700 QByteArray response;
1701 bool err = false;
1702
1703 QMutexLocker locker(&m_ioLock);
1704
1705 if (!m_io)
1706 {
1707 LOG(VB_RECORD, LOG_ERR, LOC + "External I/O not ready!");
1708 return true;
1709 }
1710
1711 if (m_io->Error())
1712 {
1713 LOG(VB_GENERAL, LOG_ERR, "External Recorder in bad state: " +
1714 m_io->ErrorString());
1715 return true;
1716 }
1717
1718 response = m_io->GetStatus(0ms);
1719 while (!response.isEmpty())
1720 {
1721 if (m_apiVersion > 2)
1722 {
1723 QJsonParseError parseError {};
1724 QJsonDocument doc;
1725 QVariantMap elements;
1726
1727 doc = QJsonDocument::fromJson(response, &parseError);
1728
1729 if (parseError.error != QJsonParseError::NoError)
1730 {
1731 LOG(VB_GENERAL, LOG_ERR, LOC +
1732 QString("ExternalRecorder returned invalid JSON message: %1: %2\n%3\n")
1733 .arg(parseError.offset)
1734 .arg(parseError.errorString(), QString(response)));
1735 }
1736 else
1737 {
1738 elements = doc.toVariant().toMap();
1739 if (elements.contains("command") &&
1740 (QString::compare(elements["command"].toString(),
1741 "STATUS",
1742 Qt::CaseInsensitive) == 0))
1743 {
1744 LogLevel_t level { LOG_INFO };
1745 QString status = elements["status"].toString().trimmed();
1746 QString message = elements["message"].toString();
1747 if (status.startsWith("crit", Qt::CaseInsensitive))
1748 {
1749 level = LOG_CRIT;
1750 }
1751 if (status.startsWith("err", Qt::CaseInsensitive))
1752 {
1753 level = LOG_ERR;
1754 }
1755 else if (status.startsWith("warn",
1756 Qt::CaseInsensitive))
1757 {
1758 level = LOG_WARNING;
1759 }
1760 else if (status.startsWith("debug",
1761 Qt::CaseInsensitive))
1762 {
1763 level = LOG_DEBUG;
1764 }
1765 else if (status.startsWith("trace",
1766 Qt::CaseInsensitive))
1767 {
1768 level = LOG_TRACE;
1769 }
1770 else if (status.startsWith("damage",
1771 Qt::CaseInsensitive))
1772 {
1773 level = LOG_WARNING;
1774 if (m_recording)
1775 m_damaged |= true;
1776 }
1777
1778 if (message.trimmed().startsWith("damage",
1779 Qt::CaseInsensitive))
1780 {
1781 if (m_recording)
1782 m_damaged |= true;
1783 }
1784
1785 LOG(VB_RECORD, level,
1786 LOC + QString("%1:%2%3")
1787 .arg(status, message,
1788 m_damaged ? " (Damaged)" : ""));
1789 }
1790 }
1791 }
1792 else
1793 {
1794 QString res = QString(response);
1795 if (m_apiVersion == 2)
1796 {
1797 QStringList tokens = res.split(':', Qt::SkipEmptyParts);
1798 tokens.removeFirst();
1799 res = tokens.join(':');
1800 for (int idx = 1; idx < tokens.size(); ++idx)
1801 {
1802 err |= tokens[idx].startsWith("ERR",
1803 Qt::CaseInsensitive);
1804 if (m_recording)
1805 m_damaged |= tokens[idx].startsWith("damage",
1806 Qt::CaseInsensitive);
1807 }
1808 }
1809 else
1810 {
1811 err |= res.startsWith("STATUS:ERR",
1812 Qt::CaseInsensitive);
1813 if (m_recording)
1814 m_damaged |= res.startsWith("STATUS:DAMAGE",
1815 Qt::CaseInsensitive);
1816 }
1817
1818 LOG(VB_RECORD, (err ? LOG_WARNING : LOG_INFO), LOC + res);
1819 }
1820
1821 response = m_io->GetStatus(0ms);
1822 }
1823
1824 return err;
1825}
1826
1828{
1829 if (m_io)
1830 {
1831 QByteArray buffer;
1832 m_io->Read(buffer, PACKET_SIZE, 1ms);
1833 m_io->GetStatus(1ms);
1834 }
1835}
1836
1838{
1839 // TODO report on buffer overruns, etc.
1840}
#define LOC
QStringList m_args
bool Ready(int fd, std::chrono::milliseconds timeout, const QString &what)
int Write(const QByteArray &buffer)
static bool KillIfRunning(const QString &cmd)
int Read(QByteArray &buffer, int maxlen, std::chrono::milliseconds timeout=2500ms)
static constexpr uint8_t kMaxErrorCnt
QByteArray GetStatus(std::chrono::milliseconds timeout=2500ms)
QTextStream m_status
bool Error(void) const
ExternIO(const QString &app, const QStringList &args)
QString ErrorString(void) const
bool ProcessVer1(const QString &cmd, QString &result, std::chrono::milliseconds timeout, uint retry_cnt)
static ExternalStreamHandler * Get(const QString &devname, int inputid, int majorid)
void run(void) override
Runs the Qt event loop unless we have a QRunnable, in which case we run the runnable run instead.
bool StartStreaming(bool recording)
ExternalStreamHandler(const QString &path, int inputid, int majorid)
bool ProcessVer2(const QString &command, QString &result, std::chrono::milliseconds timeout, uint retry_cnt)
static QMap< int, uint > s_handlersRefCnt
bool ProcessCommand(const QString &cmd, QString &result, std::chrono::milliseconds timeout=4s, uint retry_cnt=3)
static QMap< int, ExternalStreamHandler * > s_handlers
bool ProcessJson(const QVariantMap &vmsg, QVariantMap &elements, QByteArray &response, std::chrono::milliseconds timeout=4s, uint retry_cnt=3)
void PriorityEvent(int fd) override
static void Return(ExternalStreamHandler *&ref, int inputid)
void RunProlog(void)
Sets up a thread, call this if you reimplement run().
Definition: mthread.cpp:180
void RunEpilog(void)
Cleans up a thread's resources, call this if you reimplement run().
Definition: mthread.cpp:193
void setObjectName(const QString &name)
Definition: mthread.cpp:222
A QElapsedTimer based timer to replace use of QTime as a timer.
Definition: mythtimer.h:14
std::chrono::milliseconds restart(void)
Returns milliseconds elapsed since last start() or restart() and resets the count.
Definition: mythtimer.cpp:62
std::chrono::milliseconds elapsed(void)
Returns milliseconds elapsed since last start() or restart()
Definition: mythtimer.cpp:91
bool isRunning(void) const
Returns true if start() or restart() has been called at least once since construction and since any c...
Definition: mythtimer.cpp:135
void stop(void)
Stops timer, next call to isRunning() will return false and any calls to elapsed() or restart() will ...
Definition: mythtimer.cpp:78
@ kStartRunning
Definition: mythtimer.h:17
void start(void)
starts measuring elapsed time.
Definition: mythtimer.cpp:47
StreamDataList m_streamDataList
QString m_device
volatile bool m_runningDesired
volatile bool m_bError
bool RemoveAllPIDFilters(void)
void SetRunning(bool running, bool using_buffering, bool using_section_reader)
bool UpdateFiltersFromStreamData(void)
QRecursiveMutex m_listenerLock
#define O_NONBLOCK
Definition: compat.h:142
unsigned int uint
Definition: compat.h:60
#define close
Definition: compat.h:28
#define WEXITSTATUS(w)
Definition: compat.h:136
@ GENERIC_EXIT_DAEMONIZING_ERROR
Error daemonizing or execl.
Definition: exitcodes.h:31
@ GENERIC_EXIT_PIPE_FAILURE
Error creating I/O pipes.
Definition: exitcodes.h:29
QString logPropagateArgs
Definition: logging.cpp:86
#define ENO
This can be appended to the LOG args with "+".
Definition: mythlogging.h:74
#define LOG(_MASK_, _LEVEL_, _QSTRING_)
Definition: mythlogging.h:39
Convenience inline random number generator functions.
QString toString(const QDateTime &raw_dt, uint format)
Returns formatted string representing the time.
Definition: mythdate.cpp:93
def read(device=None, features=[])
Definition: disc.py:35
def error(message)
Definition: smolt.py:409
def write(text, progress=True)
Definition: mythburn.py:306