22#include "mythrecordingplayback.h"
23#include "private/debug.h"
24#include "private/ringbuffer.h"
25#include "private/os/threads/latch.h"
26#include "private/builtin.h"
31#define BUFFER_CAPACITY 2
40RecordingPlayback::RecordingPlayback(
EventHandler& handler)
42, m_eventHandler(handler)
43, m_eventSubscriberId(0)
47, m_chunk(MYTH_RECORDING_CHUNK_SIZE)
49 m_buffer.rbuf =
new RingBuffer(BUFFER_CAPACITY);
50 m_buffer.packet =
nullptr;
51 m_buffer.consumed = 0;
53 m_eventSubscriberId = m_eventHandler.CreateSubscription(
this);
54 m_eventHandler.SubscribeForEvent(m_eventSubscriberId, EVENT_UPDATE_FILE_SIZE);
58RecordingPlayback::RecordingPlayback(
const std::string& server,
unsigned port)
60, m_eventHandler(server, port)
61, m_eventSubscriberId(0)
65, m_chunk(MYTH_RECORDING_CHUNK_SIZE)
67 m_buffer.rbuf =
new RingBuffer(BUFFER_CAPACITY);
68 m_buffer.packet =
nullptr;
69 m_buffer.consumed = 0;
72 m_eventSubscriberId = m_eventHandler.CreateSubscription(
this);
73 m_eventHandler.SubscribeForEvent(m_eventSubscriberId, EVENT_UPDATE_FILE_SIZE);
77RecordingPlayback::~RecordingPlayback()
79 if (m_eventSubscriberId)
80 m_eventHandler.RevokeSubscription(m_eventSubscriberId);
83 m_buffer.rbuf->freePacket(m_buffer.packet);
87bool RecordingPlayback::Open()
90 OS::WriteLock lock(*m_latch);
91 if (ProtoPlayback::IsOpen())
93 if (ProtoPlayback::Open())
95 if (!m_eventHandler.IsRunning())
96 m_eventHandler.Start();
102void RecordingPlayback::Close()
105 OS::WriteLock lock(*m_latch);
107 ProtoPlayback::Close();
110bool RecordingPlayback::OpenTransfer(ProgramPtr recording)
113 OS::WriteLock lock(*m_latch);
114 if (!ProtoPlayback::IsOpen())
119 m_transfer.reset(
new ProtoTransfer(m_server, m_port, recording->fileName, recording->recording.storageGroup));
120 if (m_transfer->Open())
122 m_recording.swap(recording);
123 m_recording->fileSize = m_transfer->GetSize();
131void RecordingPlayback::CloseTransfer()
134 OS::WriteLock lock(*m_latch);
138 TransferDone(*m_transfer);
144bool RecordingPlayback::TransferIsOpen()
146 m_latch->lock_shared();
147 ProtoTransferPtr transfer(m_transfer);
148 m_latch->unlock_shared();
150 return ProtoPlayback::TransferIsOpen(*transfer);
154void RecordingPlayback::SetChunk(
unsigned size)
156 if (size < MYTH_RECORDING_CHUNK_MIN)
157 size = MYTH_RECORDING_CHUNK_MIN;
158 else if (size > MYTH_RECORDING_CHUNK_MAX)
159 size = MYTH_RECORDING_CHUNK_MAX;
163int64_t RecordingPlayback::GetSize()
const
165 m_latch->lock_shared();
166 ProtoTransferPtr transfer(m_transfer);
167 m_latch->unlock_shared();
169 return transfer->GetSize();
173int RecordingPlayback::Read(
void* buffer,
unsigned n)
177 if (m_buffer.packet ==
nullptr)
180 m_buffer.packet = m_buffer.rbuf->read();
181 m_buffer.consumed = 0;
186 int s = m_buffer.packet->size - m_buffer.consumed;
187 int r = ((int)n < s ? (int)n : s);
188 memcpy(
static_cast<char*
>(buffer), m_buffer.packet->data + m_buffer.consumed, r);
189 m_buffer.consumed += r;
190 if (m_buffer.consumed >= m_buffer.packet->size)
192 m_buffer.rbuf->freePacket(m_buffer.packet);
193 m_buffer.packet =
nullptr;
199 RingBufferPacket * p = m_buffer.rbuf->newPacket(m_chunk);
200 int r = _read(p->data, m_chunk);
204 m_buffer.rbuf->writePacket(p);
207 m_buffer.rbuf->freePacket(p);
214int RecordingPlayback::_read(
void *buffer,
unsigned n)
216 m_latch->lock_shared();
217 ProtoTransferPtr transfer(m_transfer);
218 m_latch->unlock_shared();
223 int64_t s = transfer->GetRemaining();
229 return TransferRequestBlock(*transfer, buffer, n);
236 return TransferRequestBlock(*transfer, buffer, n);
242int64_t RecordingPlayback::Seek(int64_t offset, WHENCE_t whence)
244 if (whence == WHENCE_CUR)
247 unsigned unread = m_buffer.rbuf->bytesUnread() + (m_buffer.packet ? m_buffer.packet->size - m_buffer.consumed : 0);
250 int64_t p = _seek(offset, whence);
252 return (p >= unread ? p - unread : p);
260 m_buffer.rbuf->freePacket(m_buffer.packet);
261 m_buffer.packet =
nullptr;
263 m_buffer.rbuf->clear();
265 return _seek(offset, whence);
268int64_t RecordingPlayback::_seek(int64_t offset, WHENCE_t whence)
270 m_latch->lock_shared();
271 ProtoTransferPtr transfer(m_transfer);
272 m_latch->unlock_shared();
274 return TransferSeek(*transfer, offset, whence);
278int64_t RecordingPlayback::GetPosition()
const
280 m_latch->lock_shared();
281 ProtoTransferPtr transfer(m_transfer);
282 m_latch->unlock_shared();
286 unsigned unread = m_buffer.rbuf->bytesUnread() + (m_buffer.packet ? m_buffer.packet->size - m_buffer.consumed : 0);
287 return transfer->GetPosition() - unread;
292void RecordingPlayback::HandleBackendMessage(EventMessagePtr msg)
295 m_latch->lock_shared();
296 ProgramPtr recording(m_recording);
297 ProtoTransferPtr transfer(m_transfer);
298 m_latch->unlock_shared();
301 case EVENT_UPDATE_FILE_SIZE:
302 if (msg->subject.size() >= 3 && recording && transfer)
306 if (msg->subject.size() >= 4)
310 if (string_to_uint32(msg->subject[1].c_str(), &chanid)
311 || string_to_time(msg->subject[2].c_str(), &startts)
312 || recording->channel.chanId != chanid
313 || recording->recording.startTs != startts
314 || string_to_int64(msg->subject[3].c_str(), &newsize))
321 if (string_to_uint32(msg->subject[1].c_str(), &recordedid)
322 || recording->recording.recordedId != recordedid
323 || string_to_int64(msg->subject[2].c_str(), &newsize))
328 transfer->SetSize(newsize);
329 recording->fileSize = newsize;
330 DBG(DBG_DEBUG,
"%s: (%d) %s %" PRIi64
"\n", __FUNCTION__,
331 msg->event, recording->fileName.c_str(), newsize);
This is the main namespace that encloses all public classes.