CPPMyth
Library to interoperate with MythTV server
Loading...
Searching...
No Matches
mythrecordingplayback.cpp
1/*
2 * Copyright (C) 2014 Jean-Luc Barriere
3 *
4 * This Program is free software; you can redistribute it and/or modify
5 * it under the terms of the GNU General Public License as published by
6 * the Free Software Foundation; either version 2, or (at your option)
7 * any later version.
8 *
9 * This Program is distributed in the hope that it will be useful,
10 * but WITHOUT ANY WARRANTY; without even the implied warranty of
11 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
12 * GNU General Public License for more details.
13 *
14 * You should have received a copy of the GNU General Public License
15 * along with this program; see the file COPYING. If not, write to
16 * the Free Software Foundation, 51 Franklin Street, Fifth Floor, Boston,
17 * MA 02110-1301 USA
18 * http://www.gnu.org/copyleft/gpl.html
19 *
20 */
21
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"
27
28#include <limits>
29#include <cstdio>
30
31#define BUFFER_CAPACITY 2 // 2 chunks
32
33using namespace Myth;
34
39
40RecordingPlayback::RecordingPlayback(EventHandler& handler)
41: ProtoPlayback(handler.GetServer(), handler.GetPort()), EventSubscriber()
42, m_eventHandler(handler)
43, m_eventSubscriberId(0)
44, m_transfer(nullptr)
45, m_recording(nullptr)
46, m_readAhead(false)
47, m_chunk(MYTH_RECORDING_CHUNK_SIZE)
48{
49 m_buffer.rbuf = new RingBuffer(BUFFER_CAPACITY);
50 m_buffer.packet = nullptr;
51 m_buffer.consumed = 0;
52
53 m_eventSubscriberId = m_eventHandler.CreateSubscription(this);
54 m_eventHandler.SubscribeForEvent(m_eventSubscriberId, EVENT_UPDATE_FILE_SIZE);
55 Open();
56}
57
58RecordingPlayback::RecordingPlayback(const std::string& server, unsigned port)
59: ProtoPlayback(server, port), EventSubscriber()
60, m_eventHandler(server, port)
61, m_eventSubscriberId(0)
62, m_transfer(nullptr)
63, m_recording(nullptr)
64, m_readAhead(false)
65, m_chunk(MYTH_RECORDING_CHUNK_SIZE)
66{
67 m_buffer.rbuf = new RingBuffer(BUFFER_CAPACITY);
68 m_buffer.packet = nullptr;
69 m_buffer.consumed = 0;
70
71 // Private handler will be stopped and closed by destructor.
72 m_eventSubscriberId = m_eventHandler.CreateSubscription(this);
73 m_eventHandler.SubscribeForEvent(m_eventSubscriberId, EVENT_UPDATE_FILE_SIZE);
74 Open();
75}
76
77RecordingPlayback::~RecordingPlayback()
78{
79 if (m_eventSubscriberId)
80 m_eventHandler.RevokeSubscription(m_eventSubscriberId);
81 Close();
82 if (m_buffer.packet)
83 m_buffer.rbuf->freePacket(m_buffer.packet);
84 delete m_buffer.rbuf;
85}
86
87bool RecordingPlayback::Open()
88{
89 // Begin critical section
90 OS::WriteLock lock(*m_latch);
91 if (ProtoPlayback::IsOpen())
92 return true;
93 if (ProtoPlayback::Open())
94 {
95 if (!m_eventHandler.IsRunning())
96 m_eventHandler.Start();
97 return true;
98 }
99 return false;
100}
101
102void RecordingPlayback::Close()
103{
104 // Begin critical section
105 OS::WriteLock lock(*m_latch);
106 CloseTransfer();
107 ProtoPlayback::Close();
108}
109
110bool RecordingPlayback::OpenTransfer(ProgramPtr recording)
111{
112 // Begin critical section
113 OS::WriteLock lock(*m_latch);
114 if (!ProtoPlayback::IsOpen())
115 return false;
116 CloseTransfer();
117 if (recording)
118 {
119 m_transfer.reset(new ProtoTransfer(m_server, m_port, recording->fileName, recording->recording.storageGroup));
120 if (m_transfer->Open())
121 {
122 m_recording.swap(recording);
123 m_recording->fileSize = m_transfer->GetSize();
124 return true;
125 }
126 m_transfer.reset();
127 }
128 return false;
129}
130
131void RecordingPlayback::CloseTransfer()
132{
133 // Begin critical section
134 OS::WriteLock lock(*m_latch);
135 m_recording.reset();
136 if (m_transfer)
137 {
138 TransferDone(*m_transfer);
139 m_transfer->Close();
140 m_transfer.reset();
141 }
142}
143
144bool RecordingPlayback::TransferIsOpen()
145{
146 m_latch->lock_shared();
147 ProtoTransferPtr transfer(m_transfer);
148 m_latch->unlock_shared();
149 if (transfer)
150 return ProtoPlayback::TransferIsOpen(*transfer);
151 return false;
152}
153
154void RecordingPlayback::SetChunk(unsigned size)
155{
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;
160 m_chunk = size;
161}
162
163int64_t RecordingPlayback::GetSize() const
164{
165 m_latch->lock_shared();
166 ProtoTransferPtr transfer(m_transfer);
167 m_latch->unlock_shared();
168 if (transfer)
169 return transfer->GetSize();
170 return 0;
171}
172
173int RecordingPlayback::Read(void* buffer, unsigned n)
174{
175 for (;;)
176 {
177 if (m_buffer.packet == nullptr)
178 {
179 // get next available packet
180 m_buffer.packet = m_buffer.rbuf->read();
181 m_buffer.consumed = 0;
182 }
183 // read available packet
184 if (m_buffer.packet)
185 {
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)
191 {
192 m_buffer.rbuf->freePacket(m_buffer.packet);
193 m_buffer.packet = nullptr;
194 }
195 return r;
196 }
197 // no packet available, so read to fill the buffer
198 {
199 RingBufferPacket * p = m_buffer.rbuf->newPacket(m_chunk);
200 int r = _read(p->data, m_chunk);
201 if (r > 0)
202 {
203 p->size = r;
204 m_buffer.rbuf->writePacket(p);
205 continue;
206 }
207 m_buffer.rbuf->freePacket(p);
208 return r;
209 }
210 }
211 return -1;
212}
213
214int RecordingPlayback::_read(void *buffer, unsigned n)
215{
216 m_latch->lock_shared();
217 ProtoTransferPtr transfer(m_transfer);
218 m_latch->unlock_shared();
219 if (transfer)
220 {
221 if (!m_readAhead)
222 {
223 int64_t s = transfer->GetRemaining(); // Acceptable block size
224 if (s > 0)
225 {
226 if (s < (int64_t)n)
227 n = (unsigned)s;
228 // Request block data from transfer socket
229 return TransferRequestBlock(*transfer, buffer, n);
230 }
231 return 0;
232 }
233 else
234 {
235 // Request block data from transfer socket
236 return TransferRequestBlock(*transfer, buffer, n);
237 }
238 }
239 return -1;
240}
241
242int64_t RecordingPlayback::Seek(int64_t offset, WHENCE_t whence)
243{
244 if (whence == WHENCE_CUR)
245 {
246 // Unread bytes: remaining bytes in the ring buffer + remaining bytes in the available packet
247 unsigned unread = m_buffer.rbuf->bytesUnread() + (m_buffer.packet ? m_buffer.packet->size - m_buffer.consumed : 0);
248 if (offset == 0)
249 {
250 int64_t p = _seek(offset, whence);
251 // it must returns the current position of the first byte in buffer
252 return (p >= unread ? p - unread : p);
253 }
254 // rebase to the first position in the buffer
255 offset -= unread;
256 }
257 // clear all buffered data
258 if (m_buffer.packet)
259 {
260 m_buffer.rbuf->freePacket(m_buffer.packet);
261 m_buffer.packet = nullptr;
262 }
263 m_buffer.rbuf->clear();
264
265 return _seek(offset, whence);
266}
267
268int64_t RecordingPlayback::_seek(int64_t offset, WHENCE_t whence)
269{
270 m_latch->lock_shared();
271 ProtoTransferPtr transfer(m_transfer);
272 m_latch->unlock_shared();
273 if (transfer)
274 return TransferSeek(*transfer, offset, whence);
275 return -1;
276}
277
278int64_t RecordingPlayback::GetPosition() const
279{
280 m_latch->lock_shared();
281 ProtoTransferPtr transfer(m_transfer);
282 m_latch->unlock_shared();
283 if (transfer)
284 {
285 // it must returns the current position of first byte in buffer
286 unsigned unread = m_buffer.rbuf->bytesUnread() + (m_buffer.packet ? m_buffer.packet->size - m_buffer.consumed : 0);
287 return transfer->GetPosition() - unread;
288 }
289 return 0;
290}
291
292void RecordingPlayback::HandleBackendMessage(EventMessagePtr msg)
293{
294 // First of all i hold shared resources using copies
295 m_latch->lock_shared();
296 ProgramPtr recording(m_recording);
297 ProtoTransferPtr transfer(m_transfer);
298 m_latch->unlock_shared();
299 switch (msg->event)
300 {
301 case EVENT_UPDATE_FILE_SIZE:
302 if (msg->subject.size() >= 3 && recording && transfer)
303 {
304 int64_t newsize;
305 // Message contains chanid + starttime as recorded key
306 if (msg->subject.size() >= 4)
307 {
308 uint32_t chanid;
309 time_t startts;
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))
315 break;
316 }
317 // Message contains recordedid as key
318 else
319 {
320 uint32_t recordedid;
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))
324 break;
325 }
326 // The file grows. Allow reading ahead
327 m_readAhead = true;
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);
332 }
333 break;
334 //case EVENT_HANDLER_STATUS:
335 // if (msg->subject[0] == EVENTHANDLER_DISCONNECTED)
336 // closeTransfer();
337 // break;
338 default:
339 break;
340 }
341}
This is the main namespace that encloses all public classes.
Definition mythcontrol.h:30