26 #define OPENSSL_SUPPRESS_DEPRECATED 1
27 #include <openssl/md5.h>
31 #include "spdlog/spdlog.h"
32 #include "spdlog/fmt/fmt.h"
39 : _meta( std::move(entry) )
40 , _received_at( time(nullptr) )
43 spdlog::debug(
"Creating File from FileEntry");
45 spdlog::debug(
"Allocating buffer");
47 if (_buffer ==
nullptr)
49 throw std::runtime_error(
"Failed to allocate file buffer");
53 this->calculate_partitioning();
54 this->create_blocks();
57 File::File(
const std::shared_ptr<Transmitter::FileDescription> &file_description)
59 , _file_description(file_description)
61 spdlog::debug(
"Creating File from FileDescription");
63 auto length = _file_description->data_length();
64 _buffer = (
char*)malloc(
length);
65 if (_buffer ==
nullptr)
67 throw std::runtime_error(
"No data allocated");
70 memcpy(_buffer, _file_description->data(),
length);
71 _meta = _file_description->file_entry();
77 throw std::runtime_error(
"Unsupported FEC scheme");
82 calculate_partitioning();
88 std::string content_location,
89 std::string content_type,
98 spdlog::debug(
"Creating File from data");
100 spdlog::debug(
"Allocating buffer");
101 _buffer = (
char*)malloc(
length);
102 if (_buffer ==
nullptr)
104 throw std::runtime_error(
"Failed to allocate file buffer");
106 memcpy(_buffer, data,
length);
112 unsigned char md5[MD5_DIGEST_LENGTH];
113 MD5((
const unsigned char*)data,
length, md5);
119 _meta.
content_md5 = base64_encode(md5, MD5_DIGEST_LENGTH);
127 throw std::runtime_error(
"Unsupported FEC scheme");
130 this->calculate_partitioning();
131 this->create_blocks();
136 spdlog::debug(
"Destroying File");
137 if (_own_buffer && _buffer !=
nullptr)
139 spdlog::debug(
"Freeing buffer");
153 if (symbol.source_block_number() >= _source_blocks.size()) {
154 throw std::runtime_error(fmt::format(
"FLUTE: source block number {} out of range (have {} blocks)",
155 symbol.source_block_number(), _source_blocks.size()));
158 SourceBlock& source_block = _source_blocks[ symbol.source_block_number() ];
160 if (symbol.id() >= source_block.symbols.size()) {
161 throw std::runtime_error(fmt::format(
"FLUTE: encoding symbol id {} out of range (block {} has {} symbols)",
162 symbol.id(), symbol.source_block_number(), source_block.symbols.size()));
168 symbol.decode_to(target_symbol.
data, target_symbol.
length);
171 check_source_block_completion(source_block);
172 check_file_completion();
177 auto File::check_source_block_completion( SourceBlock& block ) ->
void
179 block.complete = std::all_of(block.symbols.begin(), block.symbols.end(), [](
const auto& symbol){ return symbol.second.complete; });
182 auto File::check_file_completion() ->
void
184 _complete = std::all_of(_source_blocks.begin(), _source_blocks.end(), [](
const auto& block){ return block.second.complete; });
186 if (_complete && !_meta.content_md5.empty() && _meta.content_encoding.empty()) {
188 unsigned char md5[MD5_DIGEST_LENGTH];
189 MD5((
const unsigned char*)buffer(), length(), md5);
191 auto content_md5 = base64_decode(_meta.content_md5);
192 if (memcmp(md5, content_md5.c_str(), MD5_DIGEST_LENGTH) != 0) {
193 spdlog::debug(
"MD5 mismatch for TOI {}, discarding", _meta.toi);
196 for (
auto& block : _source_blocks) {
197 for (
auto& symbol : block.second.symbols) {
198 symbol.second.complete =
false;
200 block.second.complete =
false;
207 auto File::calculate_partitioning() ->
void
210 _nof_source_symbols = ceil((
double)_meta.fec_oti.transfer_length / (
double)_meta.fec_oti.encoding_symbol_length);
211 _nof_source_blocks = ceil((
double)_nof_source_symbols / (
double)_meta.fec_oti.max_source_block_length);
212 _large_source_block_length = ceil((
double)_nof_source_symbols / (
double)_nof_source_blocks);
213 _small_source_block_length = floor((
double)_nof_source_symbols / (
double)_nof_source_blocks);
214 _nof_large_source_blocks = _nof_source_symbols - _small_source_block_length * _nof_source_blocks;
217 auto File::create_blocks() ->
void
220 auto buffer_ptr = _buffer;
221 size_t remaining_size = _meta.fec_oti.transfer_length;
222 decltype(_nof_large_source_blocks) number = 0;
223 while (remaining_size > 0) {
225 size_t symbol_id = 0;
226 auto block_length = ( number < _nof_large_source_blocks ) ? _large_source_block_length : _small_source_block_length;
228 for (decltype(block_length) i = 0; i < block_length; i++) {
229 auto symbol_length = std::min(remaining_size, (
size_t)_meta.fec_oti.encoding_symbol_length);
230 assert(buffer_ptr + symbol_length <= _buffer + _meta.fec_oti.transfer_length);
232 SourceBlock::Symbol symbol{.data = buffer_ptr, .length = symbol_length, .complete =
false};
233 block.symbols[ symbol_id++ ] = symbol;
235 remaining_size -= symbol_length;
236 buffer_ptr += symbol_length;
238 if (remaining_size <= 0)
break;
240 _source_blocks[number++] = block;
246 int nof_symbols = std::ceil((
float)(max_size - 4) / (
float)_meta.fec_oti.encoding_symbol_length);
248 std::vector<EncodingSymbol> symbols;
250 for (
auto& block : _source_blocks) {
251 if (cnt >= nof_symbols)
break;
253 if (!block.second.complete) {
254 for (
auto& symbol : block.second.symbols) {
255 if (cnt >= nof_symbols)
break;
257 if (!symbol.second.complete && !symbol.second.queued) {
258 symbols.emplace_back(symbol.first, block.first, symbol.second.data, symbol.second.length, _meta.fec_oti.encoding_id);
259 symbol.second.queued =
true;
271 for (
auto& symbol : symbols) {
272 auto block = _source_blocks.find(symbol.source_block_number());
273 if (block != _source_blocks.end()) {
274 auto sym = block->second.symbols.find(symbol.id());
275 if (sym != block->second.symbols.end()) {
276 sym->second.queued =
false;
277 sym->second.complete = success;
279 check_source_block_completion(block->second);
280 check_file_completion();
287 if (!_been_encoded && !_meta.content_encoding.empty()) {
288 if (_meta.content_encoding ==
"gzip" || _meta.content_encoding==
"deflate") {
289 auto decomp_buffer = _buffer;
290 bool own_decomp = _own_buffer;
291 std::shared_ptr<unsigned char[]> comp_buffer(
new unsigned char[16384]);
293 .next_in =
reinterpret_cast<unsigned char*
>(decomp_buffer),
294 .avail_in =
static_cast<uint32_t
>(_meta.content_length),
295 .next_out = comp_buffer.get(),
298 spdlog::debug(
"Compressing contents with {}", _meta.content_encoding);
300 if (deflateInit2(&zs, Z_DEFAULT_COMPRESSION, Z_DEFLATED, 15 | 16, 8, Z_DEFAULT_STRATEGY) == Z_OK) {
302 auto zstate = deflate(&zs, Z_FINISH);
304 while (zstate == Z_OK) {
305 spdlog::debug(
"Part compressed: {} bytes", 16384-zs.avail_out);
306 _buffer =
reinterpret_cast<char*
>(realloc(_buffer, zs.total_out));
307 memcpy(_buffer+last_out, comp_buffer.get(), 16384-zs.avail_out);
308 last_out = zs.total_out;
310 zs.avail_out = 16384;
311 zs.next_out = comp_buffer.get();
312 zstate = deflate(&zs, Z_FINISH);
314 if (zstate==Z_STREAM_END) {
315 if (last_out != zs.total_out) {
316 spdlog::debug(
"Finish compress, last block is {} bytes. Total {} bytes", 16384-zs.avail_out, zs.total_out);
317 _buffer =
reinterpret_cast<char*
>(realloc(_buffer, zs.total_out));
318 memcpy(_buffer+last_out, comp_buffer.get(), 16384-zs.avail_out);
321 _meta.fec_oti.transfer_length = zs.total_out;
323 spdlog::error(
"Error compressing file {}: {}", _meta.toi, zs.msg);
328 if (own_decomp) free(decomp_buffer);
331 spdlog::error(
"Unknown Content-Encoding {}", _meta.content_encoding);
332 throw std::runtime_error(
"Content-Encoding not known");
335 _been_encoded =
true;
336 _been_decoded =
false;
342 if (!_been_decoded && !_meta.content_encoding.empty()) {
343 if (_meta.content_encoding ==
"gzip" || _meta.content_encoding==
"deflate") {
344 auto comp_buffer = _buffer;
345 bool own_comp = _own_buffer;
346 std::shared_ptr<unsigned char[]> decomp_buffer(
new unsigned char[16384]);
348 .next_in =
reinterpret_cast<unsigned char*
>(comp_buffer),
349 .avail_in =
static_cast<uint32_t
>(_meta.fec_oti.transfer_length),
350 .next_out = decomp_buffer.get(),
353 spdlog::debug(
"Decompressing contents with {}", _meta.content_encoding);
355 inflateInit2(&zs, 15 | ((_meta.content_encoding ==
"gzip")?16:0));
357 auto zstate = inflate(&zs, Z_FINISH);
359 while (zstate == Z_OK) {
360 spdlog::debug(
"Part decompressed: {} bytes", 16384-zs.avail_out);
361 _buffer =
reinterpret_cast<char*
>(realloc(_buffer, zs.total_out));
362 memcpy(_buffer+last_out, decomp_buffer.get(), 16384-zs.avail_out);
363 last_out = zs.total_out;
365 zs.avail_out = 16384;
366 zs.next_out = decomp_buffer.get();
367 zstate = inflate(&zs, Z_FINISH);
369 if (zstate==Z_STREAM_END) {
370 if (last_out != zs.total_out) {
371 spdlog::debug(
"Finish decompress, last block is {} bytes. Total {} bytes", 16384-zs.avail_out, zs.total_out);
372 _buffer =
reinterpret_cast<char*
>(realloc(_buffer, zs.total_out));
373 memcpy(_buffer+last_out, decomp_buffer.get(), 16384-zs.avail_out);
376 if (!_meta.content_length) {
377 _meta.content_length = zs.total_out;
378 }
else if (_meta.content_length != zs.total_out) {
379 spdlog::error(
"Decompressed length does not match expected Content-Length ({} != {})", _meta.content_length, zs.total_out);
382 spdlog::error(
"Error decompressing file {}: {}", _meta.toi, zs.msg);
386 if (own_comp) free(comp_buffer);
388 spdlog::error(
"Unknown Content-Encoding {}", _meta.content_encoding);
389 throw std::runtime_error(
"Content-Encoding not known");
392 _been_decoded =
true;
393 _been_encoded =
false;
396 if (!_meta.content_md5.empty()) {
397 unsigned char md5[MD5_DIGEST_LENGTH];
398 MD5((
const unsigned char*)buffer(), length(), md5);
400 auto content_md5 = base64_decode(_meta.content_md5);
401 if (memcmp(md5, content_md5.c_str(), MD5_DIGEST_LENGTH) != 0) {
402 spdlog::debug(
"MD5 mismatch for TOI {}, discarding", _meta.toi);
405 for (
auto& block : _source_blocks) {
406 for (
auto& symbol : block.second.symbols) {
407 symbol.second.complete =
false;
409 block.second.complete =
false;
A class for handling FEC encoding symbols.
const FecOti & fec_oti() const
Get the FEC OTI values.
void put_symbol(const EncodingSymbol &symbol)
Write the data from an encoding symbol into the appropriate place in the buffer.
size_t length() const
Get the data buffer length.
virtual ~File()
Default destructor.
void encode()
Encode the buffer using the Content-Encoding.
File(LibFlute::FileDeliveryTable::FileEntry entry)
Create a file from an FDT entry (used for reception)
void mark_completed(const std::vector< EncodingSymbol > &symbols, bool success)
Mark encoding symbols as completed.
std::vector< EncodingSymbol > get_next_symbols(size_t max_size)
Get the next encoding symbols that fit in max_size bytes.
void decode()
Decode the buffer using the Content-Encoding.
An entry for a file in the FDT.
std::string content_location