fix: keep mpv relay routes isolated
This commit is contained in:
@@ -9,6 +9,7 @@
|
|||||||
#include <QNetworkReply>
|
#include <QNetworkReply>
|
||||||
#include <QNetworkRequest>
|
#include <QNetworkRequest>
|
||||||
#include <QPointer>
|
#include <QPointer>
|
||||||
|
#include <QStringList>
|
||||||
#include <QTcpServer>
|
#include <QTcpServer>
|
||||||
#include <QTcpSocket>
|
#include <QTcpSocket>
|
||||||
#include <QTimer>
|
#include <QTimer>
|
||||||
@@ -55,7 +56,6 @@ QUrl MpvHttpStreamRelay::prepare(const QUrl &targetUrl, const QString &serverId,
|
|||||||
return {};
|
return {};
|
||||||
}
|
}
|
||||||
|
|
||||||
stop();
|
|
||||||
if (!m_server->isListening() && !m_server->listen(QHostAddress::LocalHost, 0))
|
if (!m_server->isListening() && !m_server->listen(QHostAddress::LocalHost, 0))
|
||||||
{
|
{
|
||||||
qWarning() << "[MpvHttpStreamRelay] failed to listen"
|
qWarning() << "[MpvHttpStreamRelay] failed to listen"
|
||||||
@@ -66,9 +66,12 @@ QUrl MpvHttpStreamRelay::prepare(const QUrl &targetUrl, const QString &serverId,
|
|||||||
m_targetUrl = targetUrl;
|
m_targetUrl = targetUrl;
|
||||||
m_serverId = serverId;
|
m_serverId = serverId;
|
||||||
m_streamToken = QUuid::createUuid().toString(QUuid::WithoutBraces);
|
m_streamToken = QUuid::createUuid().toString(QUuid::WithoutBraces);
|
||||||
m_network->setProxy(proxy);
|
m_routes.insert(m_streamToken, RouteState{targetUrl, serverId, proxy});
|
||||||
m_bytesRelayedSinceLastTick = 0;
|
m_bytesRelayedSinceLastTick = 0;
|
||||||
|
if (!m_speedTimer->isActive())
|
||||||
|
{
|
||||||
m_speedTimer->start();
|
m_speedTimer->start();
|
||||||
|
}
|
||||||
|
|
||||||
QUrl localUrl;
|
QUrl localUrl;
|
||||||
localUrl.setScheme(QStringLiteral("http"));
|
localUrl.setScheme(QStringLiteral("http"));
|
||||||
@@ -102,6 +105,7 @@ void MpvHttpStreamRelay::stop()
|
|||||||
Q_EMIT upstreamSpeedChanged(0);
|
Q_EMIT upstreamSpeedChanged(0);
|
||||||
}
|
}
|
||||||
m_connections.clear();
|
m_connections.clear();
|
||||||
|
m_routes.clear();
|
||||||
m_targetUrl.clear();
|
m_targetUrl.clear();
|
||||||
m_serverId.clear();
|
m_serverId.clear();
|
||||||
m_streamToken.clear();
|
m_streamToken.clear();
|
||||||
@@ -176,7 +180,7 @@ void MpvHttpStreamRelay::onSocketDisconnected(QTcpSocket *socket)
|
|||||||
|
|
||||||
void MpvHttpStreamRelay::processRequest(QTcpSocket *socket, const QByteArray &requestData)
|
void MpvHttpStreamRelay::processRequest(QTcpSocket *socket, const QByteArray &requestData)
|
||||||
{
|
{
|
||||||
if (m_targetUrl.isEmpty() || m_streamToken.isEmpty())
|
if (m_routes.isEmpty())
|
||||||
{
|
{
|
||||||
writeError(socket, 503, "Relay target is not ready");
|
writeError(socket, 503, "Relay target is not ready");
|
||||||
return;
|
return;
|
||||||
@@ -198,14 +202,24 @@ void MpvHttpStreamRelay::processRequest(QTcpSocket *socket, const QByteArray &re
|
|||||||
|
|
||||||
const QByteArray method = requestParts.at(0).toUpper();
|
const QByteArray method = requestParts.at(0).toUpper();
|
||||||
const QByteArray path = requestParts.at(1);
|
const QByteArray path = requestParts.at(1);
|
||||||
const QByteArray expectedPrefix = QByteArray("/") + m_streamToken.toUtf8() + QByteArray("/");
|
QByteArray requestPath = path;
|
||||||
if (!(method == "GET" || method == "HEAD") || !path.startsWith(expectedPrefix))
|
const int queryStart = requestPath.indexOf('?');
|
||||||
|
if (queryStart >= 0)
|
||||||
|
{
|
||||||
|
requestPath.truncate(queryStart);
|
||||||
|
}
|
||||||
|
const QList<QByteArray> pathParts = requestPath.split('/');
|
||||||
|
const QByteArray tokenBytes = pathParts.size() >= 2 ? pathParts.at(1) : QByteArray();
|
||||||
|
const QString routeToken = QString::fromUtf8(tokenBytes);
|
||||||
|
const auto routeIt = m_routes.constFind(routeToken);
|
||||||
|
if (!(method == "GET" || method == "HEAD") || routeToken.isEmpty() || routeIt == m_routes.constEnd())
|
||||||
{
|
{
|
||||||
writeError(socket, 404, "Not found");
|
writeError(socket, 404, "Not found");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
QNetworkRequest request(m_targetUrl);
|
const RouteState route = routeIt.value();
|
||||||
|
QNetworkRequest request(route.targetUrl);
|
||||||
request.setAttribute(QNetworkRequest::RedirectPolicyAttribute, QNetworkRequest::NoLessSafeRedirectPolicy);
|
request.setAttribute(QNetworkRequest::RedirectPolicyAttribute, QNetworkRequest::NoLessSafeRedirectPolicy);
|
||||||
|
|
||||||
for (int i = 1; i < lines.size(); ++i)
|
for (int i = 1; i < lines.size(); ++i)
|
||||||
@@ -239,12 +253,16 @@ void MpvHttpStreamRelay::processRequest(QTcpSocket *socket, const QByteArray &re
|
|||||||
{
|
{
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
it->routeToken = routeToken;
|
||||||
|
it->targetUrl = route.targetUrl;
|
||||||
|
it->serverId = route.serverId;
|
||||||
it->headOnly = method == "HEAD";
|
it->headOnly = method == "HEAD";
|
||||||
|
m_network->setProxy(route.proxy);
|
||||||
it->reply = it->headOnly ? m_network->head(request) : m_network->get(request);
|
it->reply = it->headOnly ? m_network->head(request) : m_network->get(request);
|
||||||
it->reply->setReadBufferSize(kReplyReadBufferBytes);
|
it->reply->setReadBufferSize(kReplyReadBufferBytes);
|
||||||
|
|
||||||
qDebug() << "[MpvHttpStreamRelay] request"
|
qDebug() << "[MpvHttpStreamRelay] request"
|
||||||
<< "| method:" << method << "| target:" << LogRedactionUtils::url(m_targetUrl)
|
<< "| method:" << method << "| target:" << LogRedactionUtils::url(route.targetUrl)
|
||||||
<< "| range:" << request.rawHeader("Range");
|
<< "| range:" << request.rawHeader("Range");
|
||||||
|
|
||||||
QNetworkReply *reply = it->reply;
|
QNetworkReply *reply = it->reply;
|
||||||
@@ -293,7 +311,7 @@ void MpvHttpStreamRelay::processRequest(QTcpSocket *socket, const QByteArray &re
|
|||||||
reply->attribute(QNetworkRequest::HttpStatusCodeAttribute).toInt() == 0)
|
reply->attribute(QNetworkRequest::HttpStatusCodeAttribute).toInt() == 0)
|
||||||
{
|
{
|
||||||
qWarning() << "[MpvHttpStreamRelay] upstream failed"
|
qWarning() << "[MpvHttpStreamRelay] upstream failed"
|
||||||
<< "| target:" << LogRedactionUtils::url(m_targetUrl)
|
<< "| target:" << LogRedactionUtils::url(it->targetUrl)
|
||||||
<< "| error:" << reply->errorString();
|
<< "| error:" << reply->errorString();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -475,6 +493,7 @@ void MpvHttpStreamRelay::closeConnection(QTcpSocket *socket)
|
|||||||
socket->disconnectFromHost();
|
socket->disconnectFromHost();
|
||||||
}
|
}
|
||||||
socket->deleteLater();
|
socket->deleteLater();
|
||||||
|
cleanupUnusedRoutes();
|
||||||
}
|
}
|
||||||
|
|
||||||
void MpvHttpStreamRelay::recordRelayedBytes(qint64 bytes)
|
void MpvHttpStreamRelay::recordRelayedBytes(qint64 bytes)
|
||||||
@@ -485,6 +504,31 @@ void MpvHttpStreamRelay::recordRelayedBytes(qint64 bytes)
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
void MpvHttpStreamRelay::cleanupUnusedRoutes()
|
||||||
|
{
|
||||||
|
QStringList activeTokens;
|
||||||
|
activeTokens.reserve(m_connections.size());
|
||||||
|
for (auto it = m_connections.cbegin(); it != m_connections.cend(); ++it)
|
||||||
|
{
|
||||||
|
if (!it->routeToken.isEmpty())
|
||||||
|
{
|
||||||
|
activeTokens.append(it->routeToken);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
for (auto it = m_routes.begin(); it != m_routes.end();)
|
||||||
|
{
|
||||||
|
if (it.key() != m_streamToken && !activeTokens.contains(it.key()))
|
||||||
|
{
|
||||||
|
it = m_routes.erase(it);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
++it;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
QByteArray MpvHttpStreamRelay::reasonPhrase(int statusCode)
|
QByteArray MpvHttpStreamRelay::reasonPhrase(int statusCode)
|
||||||
{
|
{
|
||||||
switch (statusCode)
|
switch (statusCode)
|
||||||
|
|||||||
@@ -29,11 +29,20 @@ private:
|
|||||||
struct ConnectionState {
|
struct ConnectionState {
|
||||||
QByteArray buffer;
|
QByteArray buffer;
|
||||||
QNetworkReply *reply = nullptr;
|
QNetworkReply *reply = nullptr;
|
||||||
|
QString routeToken;
|
||||||
|
QUrl targetUrl;
|
||||||
|
QString serverId;
|
||||||
bool headersSent = false;
|
bool headersSent = false;
|
||||||
bool headOnly = false;
|
bool headOnly = false;
|
||||||
bool upstreamFinished = false;
|
bool upstreamFinished = false;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
struct RouteState {
|
||||||
|
QUrl targetUrl;
|
||||||
|
QString serverId;
|
||||||
|
QNetworkProxy proxy;
|
||||||
|
};
|
||||||
|
|
||||||
void onNewConnection();
|
void onNewConnection();
|
||||||
void onSocketReadyRead(QTcpSocket *socket);
|
void onSocketReadyRead(QTcpSocket *socket);
|
||||||
void onSocketDisconnected(QTcpSocket *socket);
|
void onSocketDisconnected(QTcpSocket *socket);
|
||||||
@@ -43,6 +52,7 @@ private:
|
|||||||
void writeError(QTcpSocket *socket, int statusCode, const QByteArray &message);
|
void writeError(QTcpSocket *socket, int statusCode, const QByteArray &message);
|
||||||
void closeConnection(QTcpSocket *socket);
|
void closeConnection(QTcpSocket *socket);
|
||||||
void recordRelayedBytes(qint64 bytes);
|
void recordRelayedBytes(qint64 bytes);
|
||||||
|
void cleanupUnusedRoutes();
|
||||||
|
|
||||||
static QByteArray reasonPhrase(int statusCode);
|
static QByteArray reasonPhrase(int statusCode);
|
||||||
static bool isHopByHopHeader(QByteArray name);
|
static bool isHopByHopHeader(QByteArray name);
|
||||||
@@ -51,6 +61,7 @@ private:
|
|||||||
QNetworkAccessManager *m_network = nullptr;
|
QNetworkAccessManager *m_network = nullptr;
|
||||||
QTimer *m_speedTimer = nullptr;
|
QTimer *m_speedTimer = nullptr;
|
||||||
QHash<QTcpSocket *, ConnectionState> m_connections;
|
QHash<QTcpSocket *, ConnectionState> m_connections;
|
||||||
|
QHash<QString, RouteState> m_routes;
|
||||||
QUrl m_targetUrl;
|
QUrl m_targetUrl;
|
||||||
QString m_serverId;
|
QString m_serverId;
|
||||||
QString m_streamToken;
|
QString m_streamToken;
|
||||||
|
|||||||
Reference in New Issue
Block a user