Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion meson.build
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
# Meson module fs and its functions like fs.hash_file require atleast meson 0.59.0

project('rt-mbs-transport-function', 'c', 'cpp',
version : '1.3.0',
version : '1.3.1',
license : '5G-MAG Public',
meson_version : '>= 1.4.0',
default_options : [
Expand Down
12 changes: 12 additions & 0 deletions src/mbstf/Curl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@
* https://drive.google.com/file/d/1cinCiA778IErENZ3JN52VFW-1ffHpx7Z/view
*/

#include <unistd.h>

#include <chrono>
#include <iostream>
#include <map>
Expand Down Expand Up @@ -230,6 +232,16 @@ int Curl::getResponseCode() const
return m_statusCode;
}

void Curl::abortFetch()
{
if (!m_curl) return;
curl_socket_t sockfd;
CURLcode res = curl_easy_getinfo(m_curl, CURLINFO_ACTIVESOCKET, &sockfd);
if (res == CURLE_OK) {
close(sockfd);
}
}

bool Curl::extractProtocolAndStatusCode(std::string_view &header_line)
{
header_line.remove_prefix(5); // Skip "HTTP/"
Expand Down
2 changes: 2 additions & 0 deletions src/mbstf/Curl.hh
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,8 @@ public:
unsigned long getAge() const;
int getResponseCode() const;

void abortFetch();

Curl &setUserAgent(const std::string &user_agent);

private:
Expand Down
3 changes: 1 addition & 2 deletions src/mbstf/DistributionSession.cc
Original file line number Diff line number Diff line change
Expand Up @@ -388,8 +388,7 @@ bool DistributionSession::processEvent(Open5GSEvent &event)
ogs_assert(true == Open5GSSBIServer::sendResponse(stream, *response));
} else {
std::ostringstream err;
err << "Distribution Session [" << ptr_resource1 << "], method [" << method
<< "] is not allowed for a Distribution Session";
err << "Distribution Sessions method [" << method << "] is not allowed for a Distribution Session";
ogs_error("%s", err.str().c_str());
ogs_assert(true == NfServer::sendError(stream, OGS_SBI_HTTP_STATUS_MEHTOD_NOT_ALLOWED,
2, message, app_meta, api, std::nullopt, err.str()));
Expand Down
6 changes: 3 additions & 3 deletions src/mbstf/ObjectCarouselPackager.cc
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,7 @@ ObjectCarouselPackager::ObjectCarouselPackager(ObjectStore &object_store, Object
,m_maxStreams(0)
{
if (tunnel_address) {
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::address::from_string(tunnel_address.value()), tunnel_port);
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::make_address(tunnel_address.value()), tunnel_port);
}
startWorker();
startScheduler();
Expand All @@ -146,7 +146,7 @@ ObjectCarouselPackager::ObjectCarouselPackager(ObjectStore &object_store, Object
,m_maxStreams(0)
{
if (tunnel_address) {
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::address::from_string(tunnel_address.value()), tunnel_port);
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::make_address(tunnel_address.value()), tunnel_port);
}
startWorker();
startScheduler();
Expand All @@ -168,7 +168,7 @@ ObjectCarouselPackager::ObjectCarouselPackager(ObjectStore &object_store, Object
,m_maxStreams(0)
{
if (tunnel_address) {
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::address::from_string(tunnel_address.value()), tunnel_port);
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::make_address(tunnel_address.value()), tunnel_port);
}
startWorker();
startScheduler();
Expand Down
6 changes: 3 additions & 3 deletions src/mbstf/ObjectListPackager.cc
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,7 @@ ObjectListPackager::ObjectListPackager(ObjectStore &object_store, ObjectControll
{
sortListByPolicy();
if (tunnel_address) {
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::address::from_string(tunnel_address.value()), tunnel_port);
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::make_address(tunnel_address.value()), tunnel_port);
}
startWorker();
}
Expand All @@ -106,7 +106,7 @@ ObjectListPackager::ObjectListPackager(ObjectStore &object_store, ObjectControll
{
sortListByPolicy();
if (tunnel_address) {
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::address::from_string(tunnel_address.value()), tunnel_port);
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::make_address(tunnel_address.value()), tunnel_port);
}
startWorker();
}
Expand All @@ -120,7 +120,7 @@ ObjectListPackager::ObjectListPackager(ObjectStore &object_store, ObjectControll
,m_tunnelEndpoint()
{
if (tunnel_address) {
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::address::from_string(tunnel_address.value()), tunnel_port);
m_tunnelEndpoint = boost::asio::ip::udp::endpoint(boost::asio::ip::make_address(tunnel_address.value()), tunnel_port);
}
startWorker();
}
Expand Down
7 changes: 6 additions & 1 deletion src/mbstf/ObjectManifestHandler.cc
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,7 @@ bool ObjectManifestHandler::update(const std::shared_ptr<ObjectStore::Object> &n

const auto &new_objects = new_manifest.getObjects();
const auto &old_objects = m_objectManifest.getObjects();
ObjectManifest::ObjectsType to_del;
for (auto old_it = old_objects.begin(); old_it != old_objects.end(); old_it++) {
if (!*old_it || !old_it->value()) continue;
auto new_it = std::find_if(new_objects.begin(), new_objects.end(), [&old_it](const auto &new_obj) -> bool {
Expand All @@ -190,14 +191,18 @@ bool ObjectManifestHandler::update(const std::shared_ptr<ObjectStore::Object> &n
object_store.removeObject(metadata->objectId());
}
m_objectMetadataCache.erase(old_it->value().get());
m_objectManifest.removeObjects(*old_it);
to_del.push_back(*old_it);
} else {
/* object same or update */
(*old_it->value()) = std::move(*new_it->value());
new_manifest.removeObjects(*new_it);
}
}

for (const auto &obj : to_del) {
m_objectManifest.removeObjects(obj);
}

auto now = datetime_type::clock::now();
for (const auto &obj : new_objects) {
if (!obj || !obj.value()) continue;
Expand Down
4 changes: 2 additions & 2 deletions src/mbstf/ObjectPackager.hh
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@

#include <netinet/in.h>

#include <boost/asio/io_service.hpp>
#include <boost/asio/io_context.hpp>

#include "common.hh"
#include "Event.hh"
Expand Down Expand Up @@ -174,7 +174,7 @@ protected:

std::shared_ptr<std::recursive_mutex> m_transmitterMutex;
std::shared_ptr<LibFlute::Transmitter> m_transmitter;
boost::asio::io_service m_io;
boost::asio::io_context m_io;
uint32_t m_queuedToi;
bool m_queued;
std::atomic_bool m_deactivating;
Expand Down
5 changes: 4 additions & 1 deletion src/mbstf/PullObjectIngester.cc
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,10 @@ std::string PullObjectIngester::PullIngestFailedEvent::reprString() const {

/***** PullObjectIngester methods *****/

PullObjectIngester::~PullObjectIngester() {abort();}
PullObjectIngester::~PullObjectIngester() {
if (m_curl) m_curl->abortFetch();
abort();
}

bool PullObjectIngester::fetch(const std::string &object_id, const std::optional<time_type> &download_deadline)
{
Expand Down
2 changes: 1 addition & 1 deletion src/mbstf/meson.build
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ fiveg_api_release = get_option('fiveg_api_release')

fs = import('fs')

boost_dep = dependency('boost', required: true)
boost_dep = dependency('boost', version: '>=1.66.0', required: true)
uuid_dep = dependency('uuid', required: true)
zlib_dep = dependency('zlib', required: true)
libmpdpp_dep = dependency('mpd++', fallback: ['libmpdpp'])
Expand Down
2 changes: 1 addition & 1 deletion subprojects/rt-common-shared
Submodule rt-common-shared updated 84 files
+57 −0 lib/http_server/CaseInsensitiveTraits.hh
+108 −0 lib/http_server/DirectoryIndexHandler.cc
+54 −0 lib/http_server/DirectoryIndexHandler.hh
+134 −0 lib/http_server/DocrootFile.cc
+59 −0 lib/http_server/DocrootFile.hh
+135 −0 lib/http_server/DocrootHTTPRequestHandler.cc
+84 −0 lib/http_server/DocrootHTTPRequestHandler.hh
+101 −0 lib/http_server/HTTPRequest.cc
+70 −0 lib/http_server/HTTPRequest.hh
+54 −0 lib/http_server/HTTPRequestHandler.hh
+329 −0 lib/http_server/HTTPResponse.cc
+93 −0 lib/http_server/HTTPResponse.hh
+146 −0 lib/http_server/HTTPServer.cc
+80 −0 lib/http_server/HTTPServer.hh
+100 −0 lib/http_server/MimeTypeMap.cc
+53 −0 lib/http_server/MimeTypeMap.hh
+97 −0 lib/http_server/PathDelegatorHTTPRequestHandler.cc
+69 −0 lib/http_server/PathDelegatorHTTPRequestHandler.hh
+380 −0 lib/http_server/SockAddr.cc
+114 −0 lib/http_server/SockAddr.hh
+32 −0 lib/http_server/common.hh
+59 −0 lib/http_server/hash.hh
+67 −0 lib/http_server/meson.build
+35 −0 lib/http_server/version.h.in
+17 −0 lib/meson.build
+31 −0 lib/rtsdp/include/meson.build
+93 −0 lib/rtsdp/include/rtsdp/Attribute.hh
+87 −0 lib/rtsdp/include/rtsdp/ConnectionInformation.hh
+123 −0 lib/rtsdp/include/rtsdp/MediaDescription.hh
+105 −0 lib/rtsdp/include/rtsdp/Originator.hh
+84 −0 lib/rtsdp/include/rtsdp/RepeatField.hh
+85 −0 lib/rtsdp/include/rtsdp/RepeatTime.hh
+178 −0 lib/rtsdp/include/rtsdp/SDP.hh
+85 −0 lib/rtsdp/include/rtsdp/TimeZoneAdjustment.hh
+103 −0 lib/rtsdp/include/rtsdp/TimingInformation.hh
+32 −0 lib/rtsdp/include/rtsdp/common.hh
+18 −0 lib/rtsdp/meson.build
+242 −0 lib/rtsdp/src/Address.cc
+104 −0 lib/rtsdp/src/Address.hh
+99 −0 lib/rtsdp/src/Attribute.cc
+98 −0 lib/rtsdp/src/BandwidthInfo.cc
+91 −0 lib/rtsdp/src/BandwidthInfo.hh
+121 −0 lib/rtsdp/src/ConnectionInfo.cc
+88 −0 lib/rtsdp/src/ConnectionInfo.hh
+127 −0 lib/rtsdp/src/ConnectionInformation.cc
+126 −0 lib/rtsdp/src/ConnectionInformationImpl.cc
+98 −0 lib/rtsdp/src/ConnectionInformationImpl.hh
+53 −0 lib/rtsdp/src/EmailAddress.cc
+89 −0 lib/rtsdp/src/EmailAddress.hh
+191 −0 lib/rtsdp/src/MediaDesc.cc
+137 −0 lib/rtsdp/src/MediaDesc.hh
+89 −0 lib/rtsdp/src/MediaDescription.cc
+328 −0 lib/rtsdp/src/MediaDescriptionImpl.cc
+128 −0 lib/rtsdp/src/MediaDescriptionImpl.hh
+132 −0 lib/rtsdp/src/MediaName.cc
+85 −0 lib/rtsdp/src/MediaName.hh
+130 −0 lib/rtsdp/src/Origin.cc
+103 −0 lib/rtsdp/src/Origin.hh
+118 −0 lib/rtsdp/src/Originator.cc
+236 −0 lib/rtsdp/src/OriginatorImpl.cc
+113 −0 lib/rtsdp/src/OriginatorImpl.hh
+48 −0 lib/rtsdp/src/PhoneNumber.cc
+89 −0 lib/rtsdp/src/PhoneNumber.hh
+126 −0 lib/rtsdp/src/RepeatField.cc
+135 −0 lib/rtsdp/src/RepeatTime.cc
+96 −0 lib/rtsdp/src/SDP.cc
+591 −0 lib/rtsdp/src/SDPImpl.cc
+144 −0 lib/rtsdp/src/SDPImpl.hh
+270 −0 lib/rtsdp/src/SessionDescriptionProtocol.cc
+182 −0 lib/rtsdp/src/SessionDescriptionProtocol.hh
+99 −0 lib/rtsdp/src/TimeActive.cc
+97 −0 lib/rtsdp/src/TimeActive.hh
+114 −0 lib/rtsdp/src/TimeZoneAdjustment.cc
+115 −0 lib/rtsdp/src/TimingInfo.cc
+95 −0 lib/rtsdp/src/TimingInfo.hh
+91 −0 lib/rtsdp/src/TimingInformation.cc
+193 −0 lib/rtsdp/src/TimingInformationImpl.cc
+109 −0 lib/rtsdp/src/TimingInformationImpl.hh
+82 −0 lib/rtsdp/src/meson.build
+35 −0 lib/rtsdp/src/version.h.in
+18 −0 lib/rtsdp/tests/meson.build
+193 −0 lib/rtsdp/tests/test_sdp_library.cc
+34 −0 meson.build
+1 −0 open5gs-tools/openapi-generator-templates/cpp-restbed-server/model-source.mustache
Loading