19 #include <netinet/ip.h>
20 #include <netinet/udp.h>
27 #define OPENSSL_SUPPRESS_DEPRECATED 1
28 #include <openssl/md5.h>
29 #include "../utils/base64.h"
41 #include <system_error>
43 #include "spdlog/spdlog.h"
51 static void create_udp_pkt(
char *udp_buffer,
const boost::asio::ip::udp::endpoint &endpoint,
const char *data,
size_t data_len,
52 const boost::asio::ip::address &local_address );
53 static void create_ip_hdr(
char *ip_buffer,
const boost::asio::ip::udp::endpoint &endpoint,
size_t pkt_size,
54 const boost::asio::ip::address &local_address );
55 static uint16_t
calculate_sum(
const uint8_t *buffer,
size_t len );
65 , _file_entry({ .toi=0, .content_location=content_location})
72 _attach_file(filename);
73 _calculate_file_entry();
78 , _file_entry({ .toi=0, .content_location=content_location})
83 , _data_length(data.size())
85 _calculate_file_entry();
90 , _file_entry({ .toi=0, .content_location=content_location})
94 , _data(
reinterpret_cast<const char*
>(data.data()))
95 , _data_length(data.size())
97 _calculate_file_entry();
102 , _file_entry({ .toi=0, .content_location=content_location})
107 , _data_length(data?length:0)
109 _calculate_file_entry();
114 , _file_entry({ .toi=0, .content_location=content_location})
121 _calculate_file_entry();
126 , _file_entry(other._file_entry)
127 , _compression_type(other._compression_type)
128 , _filename(other._filename)
131 , _data_length(other._data_length)
133 if (!_filename.empty()) {
134 if (other._file_handle >= 0) {
135 _file_handle = dup(other._file_handle);
139 _data =
reinterpret_cast<char*
>(mmap(
nullptr, _data_length, PROT_READ, MAP_SHARED, _file_handle, 0));
142 char *
data =
new char[_data_length];
144 memcpy(
data, other._data, _data_length);
150 : _tsi(std::move(other._tsi))
151 , _file_entry(other._file_entry)
152 , _compression_type(other._compression_type)
153 , _filename(std::move(other._filename))
154 , _file_handle(other._file_handle)
156 , _data_length(other._data_length)
158 other._data =
nullptr;
159 other._data_length = 0;
160 other._file_handle = -1;
171 _file_entry = other._file_entry;
172 _compression_type = other._compression_type;
173 _filename = other._filename;
176 _data_length = other._data_length;
178 if (!_filename.empty()) {
179 if (other._file_handle >= 0) {
180 _file_handle = dup(other._file_handle);
184 _data =
reinterpret_cast<char*
>(mmap(
nullptr, _data_length, PROT_READ, MAP_SHARED, _file_handle, 0));
187 char *data =
new char[_data_length];
189 memcpy(data, other._data, _data_length);
198 _tsi = std::move(other._tsi);
199 _file_entry = other._file_entry;
200 _compression_type = other._compression_type;
201 _filename = std::move(other._filename);
202 _file_handle = other._file_handle;
203 other._file_handle = -1;
206 other._data =
nullptr;
207 _data_length = other._data_length;
208 other._data_length = 0;
215 if (_tsi != other._tsi)
return false;
216 if (_compression_type != other._compression_type)
return false;
219 if (_file_entry != other._file_entry)
return false;
223 if (_data_length != other._data_length)
return false;
225 if (_data == other._data)
return true;
226 return memcmp(_data, other._data, _data_length) == 0;
239 void Transmitter::FileDescription::_reset_toi()
246 if (_file_entry.toi != 0) _previous_toi = _file_entry.toi;
253 if (compression != _compression_type) {
254 _compression_type = compression;
255 switch (_compression_type) {
256 case COMPRESSION_GZIP:
257 _file_entry.content_encoding =
"gzip";
259 case COMPRESSION_DEFLATE:
260 _file_entry.content_encoding =
"deflate";
263 _file_entry.content_encoding.clear();
268 _calculate_file_entry();
276 _file_entry.content_location = location;
283 if (filename != _filename) {
285 _attach_file(filename);
288 _calculate_file_entry();
296 if (!data) data_length=0;
297 if (data != _data || _data_length != data_length) {
299 if (_data_length != data_length) {
308 }
else if (data_length) {
310 unsigned char md5[MD5_DIGEST_LENGTH];
311 MD5(
reinterpret_cast<const unsigned char*
>(data), data_length, md5);
312 if (_file_entry.content_md5 != base64_encode(md5,
sizeof(md5))) {
324 _data_length = data_length;
325 _calculate_file_entry();
333 return set_content(data.data(), data.size());
338 return set_content(
reinterpret_cast<const char*
>(data.data()), data.size());
343 _file_entry.content_type = content_type;
349 static bool is_set =
false;
352 std::tm ntp_epoch_tm = {.tm_mday=1, .tm_mon=0, .tm_year=0};
353 ntp_epoch = std::chrono::system_clock::from_time_t(std::mktime(&ntp_epoch_tm));
362 auto diff = std::chrono::duration_cast<std::chrono::seconds>(expiry_time -
_get_ntp_epoch());
363 _file_entry.expires = diff.count();
364 _file_entry.cache_control.cache_expires = _file_entry.expires;
371 auto durn = std::chrono::duration_cast<date_time_type::duration>(std::chrono::seconds(_file_entry.expires));
377 _file_entry.etag = etag;
383 return _file_entry.etag;
388 if (
static_cast<unsigned>(_file_entry.fec_oti.encoding_id) == 0) {
389 _file_entry.fec_oti.encoding_id = fec_oti.
encoding_id;
391 if (!_file_entry.fec_oti.instance_id) {
392 _file_entry.fec_oti.instance_id = fec_oti.
instance_id;
394 if (!_file_entry.fec_oti.transfer_length) {
397 if (!_file_entry.fec_oti.encoding_symbol_length) {
400 if (!_file_entry.fec_oti.max_source_block_length) {
403 if (!_file_entry.fec_oti.max_number_of_encoding_symbols) {
409 void Transmitter::FileDescription::_attach_file(
const std::string &filename)
411 _filename = filename;
412 _file_handle = open(_filename.c_str(), O_RDONLY);
413 if (_file_handle < 0) {
414 throw std::system_error(errno, std::generic_category(),
"Could not open the file");
417 off_t pos = lseek(_file_handle, 0, SEEK_END);
419 throw std::system_error(errno, std::generic_category(),
"Could not find the file length");
421 _data_length =
static_cast<size_t>(pos);
422 lseek(_file_handle, 0, SEEK_SET);
426 _data =
reinterpret_cast<char*
>(mmap(
nullptr, _data_length, PROT_READ, MAP_SHARED, _file_handle, 0));
429 char *data =
new char[_data_length];
431 ssize_t nread = read(_file_handle, data, _data_length);
432 if (nread < 0 ||
static_cast<size_t>(nread) != _data_length) {
433 throw std::system_error(errno, std::generic_category(),
"Could not read the file contents");
440 void Transmitter::FileDescription::_free_file_data()
442 if (!_filename.empty()) {
444 if (_data) munmap(
const_cast<char*
>(_data), _data_length);
445 if (_file_handle >= 0) close(_file_handle);
447 delete[]
const_cast<char*
>(_data);
453 void Transmitter::FileDescription::_calculate_file_entry()
456 _file_entry.content_length = _data_length;
459 _file_entry.fec_oti.transfer_length = _data_length;
462 if (_data && _data_length) {
463 unsigned char md5[MD5_DIGEST_LENGTH];
464 MD5(
reinterpret_cast<const unsigned char*
>(_data), _data_length, md5);
465 _file_entry.content_md5 = base64_encode(md5,
sizeof(md5));
467 _file_entry.content_md5.clear();
476 uint64_t tsi,
unsigned short mtu, uint32_t
rate_limit,
477 boost::asio::io_context& io_context,
478 const std::optional<boost::asio::ip::udp::endpoint> &tunnel_endpoint,
481 : _endpoint(boost::asio::ip::make_address(destination_address), port)
483 , _socket(io_context, _endpoint.protocol())
484 , _io_context(io_context)
485 , _send_timer(io_context)
486 , _fdt_timer(io_context)
491 , _mcast_address(destination_address)
493 , _tunnel_endpoint(tunnel_endpoint)
494 , _tunnel_local_address()
498 _source_address = boost::asio::ip::make_address(
source_address.value());
505 if (_tunnel_endpoint.has_value()) {
509 boost::asio::ip::udp::socket local_socket(_io_context, _tunnel_endpoint.value().protocol());
510 local_socket.connect(_tunnel_endpoint.value());
511 _tunnel_local_address = local_socket.local_endpoint().address();
513 uint32_t max_source_block_length = 64;
515 _socket.set_option(boost::asio::ip::multicast::enable_loopback(
true));
516 _socket.set_option(boost::asio::ip::udp::socket::reuse_address(
true));
518 if (_source_address && !_tunnel_endpoint) {
519 _socket.bind(boost::asio::ip::udp::endpoint(_source_address.value(),0));
524 .encoding_symbol_length = _max_payload,
525 .max_source_block_length = max_source_block_length};
526 _fdt = std::make_unique<FileDeliveryTable>(1, _fec_oti, fdt_namespace);
529 start_fdt_repeat_timer();
538 return udp_tunnel_address(std::optional<boost::asio::ip::udp::endpoint>(new_tunnel_endpoint));
543 return udp_tunnel_address(std::optional<boost::asio::ip::udp::endpoint>(std::move(new_tunnel_endpoint)));
548 return udp_tunnel_address(std::move(std::optional<boost::asio::ip::udp::endpoint>(new_tunnel_endpoint)));
553 if (!!_tunnel_endpoint == !!new_tunnel_endpoint) {
555 if (_tunnel_endpoint) _tunnel_endpoint = new_tunnel_endpoint;
556 }
else if (_tunnel_endpoint) {
560 _tunnel_endpoint = std::nullopt;
563 _tunnel_endpoint = std::move(new_tunnel_endpoint);
568 if (_tunnel_endpoint) {
569 boost::asio::ip::udp::socket local_socket(_io_context, _tunnel_endpoint.value().protocol());
570 local_socket.connect(_tunnel_endpoint.value());
571 _tunnel_local_address = local_socket.local_endpoint().address();
578 return udp_tunnel_address(std::optional<boost::asio::ip::udp::endpoint>(std::nullopt));
583 return endpoint(boost::asio::ip::udp::endpoint(boost::asio::ip::make_address(address), port));
588 _endpoint = destination;
594 _endpoint = std::move(destination);
600 _source_address = source_address;
601 if (_source_address && !_tunnel_endpoint) {
602 _socket.bind(boost::asio::ip::udp::endpoint(_source_address.value(),0));
609 _source_address = std::move(source_address);
610 if (_source_address && !_tunnel_endpoint) {
611 _socket.bind(boost::asio::ip::udp::endpoint(_source_address.value(),0));
621 auto Transmitter::handle_send_to(
const boost::system::error_code& error) ->
void
629 return std::chrono::duration_cast<std::chrono::seconds>(
630 std::chrono::system_clock::now().time_since_epoch()).count() +
635 auto Transmitter::send_fdt() ->
void {
636 if (_fdt->file_entries().empty())
return;
637 _fdt->set_expires(seconds_since_epoch() + _fdt_repeat_interval * 2);
638 auto fdt = _fdt->to_string();
639 auto file = std::make_shared<File>(
644 seconds_since_epoch() + _fdt_repeat_interval * 2,
649 file->set_fdt_instance_id( _fdt->instance_id() );
650 spdlog::debug(
"Sending FDT instance {}:\n{}", _fdt->instance_id(), _fdt->to_string());
652 std::lock_guard<std::mutex> guard(_files_mutex);
653 _files.insert_or_assign(0, file);
660 const std::string& content_location,
661 const std::string& content_type,
664 size_t length) -> uint16_t
668 if (_toi == 0) _toi = 1;
670 auto file = std::make_shared<File>(
679 _fdt->add(file->meta());
682 std::lock_guard<std::mutex> guard(_files_mutex);
683 _files.insert({toi, file});
688 auto Transmitter::send(
const std::shared_ptr<Transmitter::FileDescription> &file_description) -> uint16_t
690 if (file_description->has_tsi() && file_description->tsi() != _tsi) {
692 file_description->toi(0);
693 spdlog::debug(
"Reset TOI for FileDescription");
697 file_description->tsi(_tsi);
698 bool is_resend = (file_description->toi() != 0);
699 if (file_description->toi() == 0) {
700 if (file_description->previous_toi() != 0) {
705 _fdt->remove(file_description->previous_toi());
706 file_description->reset_previous_toi();
708 file_description->toi(_toi);
710 if (_toi == 0) _toi = 1;
711 spdlog::debug(
"Assigned new TOI {}", file_description->toi());
715 file_description->merge_fec_oti(_fec_oti);
717 auto file = std::make_shared<File>(file_description);
719 std::lock_guard<std::mutex> guard(_files_mutex);
724 _files[file_description->toi()] = file;
731 _fdt->remove(file_description->toi());
733 _fdt->add(file->meta());
735 return file_description->toi();
738 auto Transmitter::fdt_send_tick(
const boost::system::error_code& error) ->
void
740 if (error == boost::asio::error::operation_aborted)
return;
743 start_fdt_repeat_timer();
747 auto Transmitter::file_transmitted(uint32_t toi) ->
void
750 std::lock_guard<std::mutex> guard(_files_mutex);
757 if (_completion_cb) {
763 std::lock_guard<std::mutex> guard(_files_mutex);
764 if (_deactivate_when_all_files_sent && _files.empty()) {
765 _complete_deactivation();
770 auto Transmitter::send_next_packet() ->
void
772 uint32_t bytes_queued = 0;
774 if (!_active)
return;
775 std::shared_ptr<File> file;
777 std::lock_guard<std::mutex> guard(_files_mutex);
778 for (
auto& file_m : _files) {
779 auto &next_file = file_m.second;
781 if (next_file && !next_file->complete()) {
788 auto symbols = file->get_next_symbols(_max_payload);
790 if (symbols.size()) {
791 for(
const auto& symbol : symbols) {
792 spdlog::debug(
"sending TOI {} SBN {} ID {}", file->meta().toi, symbol.source_block_number(), symbol.id() );
794 auto packet = std::make_shared<AlcPacket>(_tsi, file->meta().toi, file->meta().fec_oti, symbols, _max_payload, file->fdt_instance_id());
795 bytes_queued += packet->size();
797 boost::asio::ip::udp::endpoint send_endpoint;
798 const char *data =
nullptr;
799 size_t data_size = 0;
800 std::shared_ptr<std::vector<char>> tunnel_data;
801 if (_tunnel_endpoint) {
802 send_endpoint = _tunnel_endpoint.value();
803 data_size = packet->size() + 20 + 8 ;
808 tunnel_data = std::make_shared<std::vector<char>>(data_size);
809 data = tunnel_data->data();
810 auto local_address = _source_address ? _source_address.value() : _tunnel_local_address;
811 create_udp_pkt(
const_cast<char*
>(data) + 20, _endpoint, packet->data(), packet->size(), local_address);
812 create_ip_hdr(
const_cast<char*
>(data), _endpoint, data_size, local_address);
814 send_endpoint = _endpoint;
815 data = packet->data();
816 data_size = packet->size();
818 _socket.async_send_to(
819 boost::asio::buffer(data, data_size),
821 [file, symbols, packet, tunnel_data,
this](
822 const boost::system::error_code& error,
823 std::size_t bytes_transferred)
827 (void)bytes_transferred;
829 spdlog::debug(
"sent_to error: {}", error.message());
831 file->mark_completed(symbols, !error);
832 if (file->complete()) {
833 file_transmitted(file->meta().toi);
841 _send_timer.expires_after(std::chrono::milliseconds(10));
842 _send_timer.async_wait( boost::bind(&Transmitter::send_next_packet,
this));
844 if (_rate_limit == 0) {
845 boost::asio::post(_io_context, boost::bind(&Transmitter::send_next_packet,
this));
847 auto send_duration = ((bytes_queued * 8.0) / (
double)_rate_limit/1000.0) * 1000.0 * 1000.0;
848 spdlog::trace(
"Rate limiter: queued {} bytes, limit {} kbps, next send in {} us",
849 bytes_queued, _rate_limit, send_duration);
850 _send_timer.expires_after(std::chrono::microseconds(
851 static_cast<int>(ceil(send_duration))));
852 _send_timer.async_wait( boost::bind(&Transmitter::send_next_packet,
this));
861 _deactivate_when_all_files_sent =
false;
863 start_fdt_repeat_timer();
871 if (finish_file_transmissions) {
872 std::lock_guard<std::mutex> guard(_files_mutex);
873 if (!_files.empty()) {
874 _deactivate_when_all_files_sent =
true;
878 _complete_deactivation();
882 _complete_deactivation();
886 auto Transmitter::_complete_deactivation() ->
void
888 _deactivate_when_all_files_sent =
false;
891 _send_timer.cancel();
894 auto Transmitter::start_fdt_repeat_timer() ->
void
896 _fdt_timer.expires_after(std::chrono::seconds(_fdt_repeat_interval));
897 _fdt_timer.async_wait( boost::bind(&Transmitter::fdt_send_tick,
this, boost::placeholders::_1));
900 static void create_udp_pkt(
char *udp_buffer,
const boost::asio::ip::udp::endpoint &endpoint,
const char *data,
size_t data_len,
const boost::asio::ip::address &local_address)
902 auto *udp_bytes =
reinterpret_cast<uint8_t*
>(udp_buffer);
903 const auto udp_length =
static_cast<uint16_t
>(data_len + 8);
904 const auto source_address = local_address.to_v4().to_uint();
905 const auto destination_address = endpoint.address().to_v4().to_uint();
911 memcpy(udp_buffer + 8, data, data_len);
913 std::vector<uint8_t> checksum_bytes(12 + udp_length);
916 checksum_bytes[8] = 0;
917 checksum_bytes[9] = endpoint.protocol().protocol();
919 memcpy(checksum_bytes.data() + 12, udp_bytes, udp_length);
924 static void create_ip_hdr(
char *ip_buffer,
const boost::asio::ip::udp::endpoint &endpoint,
size_t pkt_size,
const boost::asio::ip::address &local_address)
926 auto *ip_bytes =
reinterpret_cast<uint8_t*
>(ip_buffer);
928 memset(ip_bytes, 0, 20);
935 ip_bytes[9] = endpoint.protocol().protocol();
946 cksum += (
static_cast<uint32_t
>(buffer[0]) << 8) |
static_cast<uint32_t
>(buffer[1]);
951 cksum +=
static_cast<uint32_t
>(buffer[0]) << 8;
954 while (cksum >> 16) {
955 cksum = (cksum & 0xFFFF) + (cksum >> 16);
958 return static_cast<uint16_t
>(~cksum);
963 buffer[0] =
static_cast<uint8_t
>((value >> 8) & 0xFF);
964 buffer[1] =
static_cast<uint8_t
>(value & 0xFF);
969 buffer[0] =
static_cast<uint8_t
>((value >> 24) & 0xFF);
970 buffer[1] =
static_cast<uint8_t
>((value >> 16) & 0xFF);
971 buffer[2] =
static_cast<uint8_t
>((value >> 8) & 0xFF);
972 buffer[3] =
static_cast<uint8_t
>(value & 0xFF);
FdtNamespace
FDT namespace enumeration.
FileDescription & set_content_location(const std::string &location)
Set Content-Location.
size_t data_length()
Get the length in bytes of the data to be transmitted.
FileDescription & set_expiry_time(const date_time_type &expiry_time)
Change the file expiry time.
const char * data()
Get the data to be transmitted.
FileDescription & set_etag(const std::string &etag)
Set the ETag value for the file.
std::chrono::system_clock::time_point date_time_type
FileDescription & operator=(const FileDescription &other)
Copy operator.
FileDescription & set_content_type(const std::string &content_type)
Change the file content type.
virtual ~FileDescription()
Destructor.
FileDescription & set_content(const std::string &filename)
Change the file contents using a local file.
FileDescription & merge_fec_oti(const FecOti &fec_oti)
Merge the FecOti values.
FileDescription & set_compression(CompressionAlgorithm compression)
Set the compression algorithm.
bool operator==(const FileDescription &other) const
Equality operator.
date_time_type get_expiry_time() const
Get the currently set expiry time.
const std::string & get_etag() const
Get the current ETag value.
uint64_t seconds_since_epoch()
Convenience function to get the current timestamp for expiry calculation.
virtual ~Transmitter()
Default destructor.
const std::optional< boost::asio::ip::udp::endpoint > & udp_tunnel_address() const
Get UDP Tunnel Address.
const std::optional< boost::asio::ip::address > & source_address() const
Get the optional source address for the FLUTE session.
const boost::asio::ip::udp::endpoint & endpoint() const
Get UDP Address for FLUTE session.
uint16_t send(const std::string &content_location, const std::string &content_type, uint32_t expires, char *data, size_t length)
Transmit a file (deprecated).
void enable_ipsec(uint32_t spi, const std::string &aes_key)
Enable IPSEC ESP encryption of FLUTE payloads.
Transmitter(const std::string &destination_address, short port, uint64_t tsi, unsigned short mtu, uint32_t rate_limit, boost::asio::io_context &io_context, const std::optional< boost::asio::ip::udp::endpoint > &tunnel_endpoint=std::nullopt, FdtNamespace fdt_namespace=FileDeliveryTable::FDT_NS_NONE, bool active=true, const std::optional< std::string > &source_address=std::nullopt)
Constructor.
uint32_t rate_limit() const
Get Maximum Bit Rate.
void deactivate(bool finish_file_transmissions=false)
Deactivate the FLUTE session.
void activate()
Activate the FLUTE session.
void enable_esp(uint32_t spi, const std::string &dest_address, Direction direction, const std::string &key)
static uint16_t calculate_sum(const uint8_t *buffer, size_t len)
static void create_udp_pkt(char *udp_buffer, const boost::asio::ip::udp::endpoint &endpoint, const char *data, size_t data_len, const boost::asio::ip::address &local_address)
static void write_uint32_be(uint8_t *buffer, uint32_t value)
static const Transmitter::FileDescription::date_time_type & _get_ntp_epoch()
static void write_uint16_be(uint8_t *buffer, uint16_t value)
static void create_ip_hdr(char *ip_buffer, const boost::asio::ip::udp::endpoint &endpoint, size_t pkt_size, const boost::asio::ip::address &local_address)
uint32_t max_source_block_length
uint32_t max_number_of_encoding_symbols
uint32_t encoding_symbol_length