Combine multiple upload jobs into a single KCompositeJob so only 1 notification will be shown
Summary: Combine multiple upload jobs for files into a single KCompositeJob so only 1 notification will be shown Includes changes introduced in D16279 Test Plan: 1. Share of multiple files is performed using 1 composite job Setup: - Select multiple (big) files in dolphin and share with an Android device Result: - The files will be transferred using 1 CompositeUploadJob and showing only 1 notification 2. Share of file while another share is already running adds job to existing composite job Setup: - Select multiple (big) files in dolphin and share with an Android device - Share an additional file with the same Android device Result: - The files are all transferred using 1 CompositeUploadJob and showing only 1 notification - The notification is updated after adding the last file 3. Other packets are transmitted as usual Setup: - Setup sharing desktop notification with device - Share a big file with an Android device - Generate a desktop notification (eg. sending or receiving an email) Result: - Notification packet is send immediately Reviewers: #kde_connect, nicolasfella, albertvaka Reviewed By: #kde_connect, albertvaka Subscribers: albertvaka, apol, nicolasfella, broulik, kdeconnect Tags: #kde_connect Differential Revision: https://phabricator.kde.org/D17081
This commit is contained in:
parent
c6fc4d92b7
commit
b6c15289f5
11 changed files with 454 additions and 155 deletions
|
@ -147,22 +147,26 @@ int main(int argc, char** argv)
|
|||
}
|
||||
|
||||
if (parser.isSet(QStringLiteral("share"))) {
|
||||
QList<QUrl> urls;
|
||||
QStringList urls;
|
||||
|
||||
QUrl url = QUrl::fromUserInput(parser.value(QStringLiteral("share")), QDir::currentPath());
|
||||
urls.append(url);
|
||||
urls.append(url.toString());
|
||||
|
||||
// Check for more arguments
|
||||
const auto args = parser.positionalArguments();
|
||||
for (const QString& input : args) {
|
||||
QUrl url = QUrl::fromUserInput(input, QDir::currentPath());
|
||||
urls.append(url);
|
||||
urls.append(url.toString());
|
||||
}
|
||||
|
||||
for (const QUrl& url : urls) {
|
||||
QDBusMessage msg = QDBusMessage::createMethodCall(QStringLiteral("org.kde.kdeconnect"), "/modules/kdeconnect/devices/"+device+"/share", QStringLiteral("org.kde.kdeconnect.device.share"), QStringLiteral("shareUrl"));
|
||||
msg.setArguments(QVariantList() << url.toString());
|
||||
blockOnReply(QDBusConnection::sessionBus().asyncCall(msg));
|
||||
QTextStream(stdout) << i18n("Shared %1", url.toString()) << endl;
|
||||
QDBusMessage msg = QDBusMessage::createMethodCall(QStringLiteral("org.kde.kdeconnect"), "/modules/kdeconnect/devices/"+device+"/share",
|
||||
QStringLiteral("org.kde.kdeconnect.device.share"), QStringLiteral("shareUrls"));
|
||||
|
||||
msg.setArguments(QVariantList() << QVariant(urls));
|
||||
blockOnReply(QDBusConnection::sessionBus().asyncCall(msg));
|
||||
|
||||
for (const QString& url : qAsConst(urls)) {
|
||||
QTextStream(stdout) << i18n("Shared %1", url) << endl;
|
||||
}
|
||||
} else if(parser.isSet(QStringLiteral("pair"))) {
|
||||
DeviceDbusInterface dev(device);
|
||||
|
|
|
@ -6,6 +6,7 @@ set(backends_kdeconnect_SRCS
|
|||
backends/lan/lanlinkprovider.cpp
|
||||
backends/lan/landevicelink.cpp
|
||||
backends/lan/lanpairinghandler.cpp
|
||||
backends/lan/compositeuploadjob.cpp
|
||||
backends/lan/uploadjob.cpp
|
||||
backends/lan/socketlinereader.cpp
|
||||
|
||||
|
|
275
core/backends/lan/compositeuploadjob.cpp
Normal file
275
core/backends/lan/compositeuploadjob.cpp
Normal file
|
@ -0,0 +1,275 @@
|
|||
/**
|
||||
* Copyright 2018 Erik Duisters
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or
|
||||
* modify it under the terms of the GNU General Public License as
|
||||
* published by the Free Software Foundation; either version 2 of
|
||||
* the License or (at your option) version 3 or any later version
|
||||
* accepted by the membership of KDE e.V. (or its successor approved
|
||||
* by the membership of KDE e.V.), which shall act as a proxy
|
||||
* defined in Section 14 of version 3 of the license.
|
||||
*
|
||||
* This program is distributed in the hope that it will be useful,
|
||||
* but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
* GNU General Public License for more details.
|
||||
*
|
||||
* You should have received a copy of the GNU General Public License
|
||||
* along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
*/
|
||||
|
||||
#include "compositeuploadjob.h"
|
||||
#include <core_debug.h>
|
||||
#include <KLocalizedString>
|
||||
#include <kio/global.h>
|
||||
#include <KJobTrackerInterface>
|
||||
#include "lanlinkprovider.h"
|
||||
#include <daemon.h>
|
||||
|
||||
CompositeUploadJob::CompositeUploadJob(const QString& deviceId, bool displayNotification)
|
||||
: KCompositeJob()
|
||||
, m_server(new Server(this))
|
||||
, m_socket(nullptr)
|
||||
, m_port(0)
|
||||
, m_deviceId(deviceId)
|
||||
, m_running(false)
|
||||
, m_currentJobNum(1)
|
||||
, m_totalJobs(0)
|
||||
, m_currentJobSendPayloadSize(0)
|
||||
, m_totalSendPayloadSize(0)
|
||||
, m_totalPayloadSize(0)
|
||||
, m_currentJob(nullptr)
|
||||
{
|
||||
setCapabilities(Killable);
|
||||
|
||||
if (displayNotification) {
|
||||
KIO::getJobTracker()->registerJob(this);
|
||||
}
|
||||
}
|
||||
|
||||
bool CompositeUploadJob::isRunning()
|
||||
{
|
||||
return m_running;
|
||||
}
|
||||
|
||||
void CompositeUploadJob::start() {
|
||||
if (m_running) {
|
||||
qCWarning(KDECONNECT_CORE) << "CompositeUploadJob::start() - allready running";
|
||||
return;
|
||||
}
|
||||
|
||||
if (!hasSubjobs()) {
|
||||
qCWarning(KDECONNECT_CORE) << "CompositeUploadJob::start() - there are no subjobs to start";
|
||||
emitResult();
|
||||
return;
|
||||
}
|
||||
|
||||
if (!startListening()) {
|
||||
return;
|
||||
}
|
||||
|
||||
connect(m_server, &QTcpServer::newConnection, this, &CompositeUploadJob::newConnection);
|
||||
|
||||
m_running = true;
|
||||
|
||||
//Give SharePlugin some time to add subjobs
|
||||
QMetaObject::invokeMethod(this, "startNextSubJob", Qt::QueuedConnection);
|
||||
}
|
||||
|
||||
bool CompositeUploadJob::startListening()
|
||||
{
|
||||
m_port = MIN_PORT;
|
||||
while (!m_server->listen(QHostAddress::Any, m_port)) {
|
||||
m_port++;
|
||||
if (m_port > MAX_PORT) { //No ports available?
|
||||
qCWarning(KDECONNECT_CORE) << "CompositeUploadJob::startListening() - Error opening a port in range" << MIN_PORT << "-" << MAX_PORT;
|
||||
m_port = 0;
|
||||
setError(NoPortAvailable);
|
||||
setErrorText(i18n("Couldn't find an available port"));
|
||||
emitResult();
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
qCDebug(KDECONNECT_CORE) << "CompositeUploadJob::startListening() - listening on port: " << m_port;
|
||||
return true;
|
||||
}
|
||||
|
||||
void CompositeUploadJob::startNextSubJob()
|
||||
{
|
||||
m_currentJob = qobject_cast<UploadJob*>(subjobs().at(0));
|
||||
m_currentJobSendPayloadSize = 0;
|
||||
emitDescription(m_currentJob->getNetworkPacket().get<QString>(QStringLiteral("filename")));
|
||||
|
||||
connect(m_currentJob, SIGNAL(processedAmount(KJob*,KJob::Unit,qulonglong)), this, SLOT(slotProcessedAmount(KJob*,KJob::Unit,qulonglong)));
|
||||
//Already done by KCompositeJob
|
||||
//connect(m_currentJob, &KJob::result, this, &CompositeUploadJob::slotResult);
|
||||
|
||||
//TODO: Create a copy of the networkpacket that can be re-injected if sending via lan fails?
|
||||
NetworkPacket np = m_currentJob->getNetworkPacket();
|
||||
np.setPayload(nullptr, np.payloadSize());
|
||||
np.setPayloadTransferInfo({{"port", m_port}});
|
||||
np.set<int>(QStringLiteral("numberOfFiles"), m_totalJobs);
|
||||
np.set<quint64>(QStringLiteral("totalPayloadSize"), m_totalPayloadSize);
|
||||
|
||||
if (Daemon::instance()->getDevice(m_deviceId)->sendPacket(np)) {
|
||||
m_server->resumeAccepting();
|
||||
} else {
|
||||
setError(SendingNetworkPacketFailed);
|
||||
setErrorText(i18n("Failed to send packet to %1", Daemon::instance()->getDevice(m_deviceId)->name()));
|
||||
//TODO: cleanup/resend remaining jobs
|
||||
emitResult();
|
||||
}
|
||||
}
|
||||
|
||||
void CompositeUploadJob::newConnection()
|
||||
{
|
||||
m_server->pauseAccepting();
|
||||
|
||||
m_socket = m_server->nextPendingConnection();
|
||||
|
||||
if (!m_socket) {
|
||||
qCDebug(KDECONNECT_CORE) << "CompositeUploadJob::newConnection() - m_server->nextPendingConnection() returned a nullptr";
|
||||
return;
|
||||
}
|
||||
|
||||
m_currentJob->setSocket(m_socket);
|
||||
|
||||
connect(m_socket, &QSslSocket::disconnected, this, &CompositeUploadJob::socketDisconnected);
|
||||
connect(m_socket, QOverload<QAbstractSocket::SocketError>::of(&QAbstractSocket::error), this, &CompositeUploadJob::socketError);
|
||||
connect(m_socket, QOverload<const QList<QSslError> &>::of(&QSslSocket::sslErrors), this, &CompositeUploadJob::sslError);
|
||||
connect(m_socket, &QSslSocket::encrypted, this, &CompositeUploadJob::encrypted);
|
||||
|
||||
LanLinkProvider::configureSslSocket(m_socket, m_deviceId, true);
|
||||
|
||||
m_socket->startServerEncryption();
|
||||
}
|
||||
|
||||
void CompositeUploadJob::socketDisconnected()
|
||||
{
|
||||
m_socket->close();
|
||||
}
|
||||
|
||||
void CompositeUploadJob::socketError(QAbstractSocket::SocketError error)
|
||||
{
|
||||
Q_UNUSED(error);
|
||||
|
||||
m_socket->close();
|
||||
setError(SocketError);
|
||||
emitResult();
|
||||
//TODO: cleanup jobs?
|
||||
m_running = false;
|
||||
}
|
||||
|
||||
void CompositeUploadJob::sslError(const QList<QSslError>& errors)
|
||||
{
|
||||
Q_UNUSED(errors);
|
||||
|
||||
m_socket->close();
|
||||
setError(SslError);
|
||||
emitResult();
|
||||
//TODO: cleanup jobs?
|
||||
m_running = false;
|
||||
}
|
||||
|
||||
void CompositeUploadJob::encrypted()
|
||||
{
|
||||
if (!m_timer.isValid()) {
|
||||
m_timer.start();
|
||||
}
|
||||
|
||||
m_currentJob->start();
|
||||
}
|
||||
|
||||
bool CompositeUploadJob::addSubjob(KJob* job)
|
||||
{
|
||||
if (UploadJob *uploadJob = qobject_cast<UploadJob*>(job)) {
|
||||
NetworkPacket np = uploadJob->getNetworkPacket();
|
||||
|
||||
m_totalJobs++;
|
||||
|
||||
if (np.payloadSize() >= 0 ) {
|
||||
m_totalPayloadSize += np.payloadSize();
|
||||
setTotalAmount(Bytes, m_totalPayloadSize);
|
||||
}
|
||||
|
||||
QString filename;
|
||||
QString filenameArg = QStringLiteral("filename");
|
||||
|
||||
if (m_currentJob) {
|
||||
filename = m_currentJob->getNetworkPacket().get<QString>(filenameArg);
|
||||
} else {
|
||||
filename = np.get<QString>(filenameArg);
|
||||
}
|
||||
|
||||
emitDescription(filename);
|
||||
|
||||
return KCompositeJob::addSubjob(job);
|
||||
} else {
|
||||
qCDebug(KDECONNECT_CORE) << "CompositeUploadJob::addSubjob() - you can only add UploadJob's, ignoring";
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
bool CompositeUploadJob::doKill()
|
||||
{
|
||||
//TODO: Remove all subjobs?
|
||||
//TODO: cleanup jobs?
|
||||
if (m_running) {
|
||||
m_running = false;
|
||||
|
||||
return m_currentJob->stop();
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
void CompositeUploadJob::slotProcessedAmount(KJob *job, KJob::Unit unit, qulonglong amount) {
|
||||
Q_UNUSED(job);
|
||||
|
||||
m_currentJobSendPayloadSize = amount;
|
||||
|
||||
quint64 uploaded = m_totalSendPayloadSize + m_currentJobSendPayloadSize;
|
||||
setProcessedAmount(unit, uploaded);
|
||||
|
||||
const auto elapsed = m_timer.elapsed();
|
||||
if (elapsed > 0) {
|
||||
emitSpeed((1000 * uploaded) / elapsed);
|
||||
}
|
||||
}
|
||||
|
||||
void CompositeUploadJob::slotResult(KJob *job) {
|
||||
//Copies job error and errorText and emits result if job is in error otherwise removes job from subjob list
|
||||
KCompositeJob::slotResult(job);
|
||||
|
||||
//TODO: cleanup jobs?
|
||||
|
||||
if (error() || !m_running) {
|
||||
return;
|
||||
}
|
||||
|
||||
m_totalSendPayloadSize += m_currentJobSendPayloadSize;
|
||||
|
||||
if (hasSubjobs()) {
|
||||
m_currentJobNum++;
|
||||
startNextSubJob();
|
||||
} else {
|
||||
Q_EMIT description(this, i18n("Finished sending to %1", Daemon::instance()->getDevice(this->m_deviceId)->name()),
|
||||
{ QStringLiteral(""), i18np("Sent 1 file", "Sent %1 files", m_totalJobs) }
|
||||
);
|
||||
emitResult();
|
||||
}
|
||||
}
|
||||
|
||||
void CompositeUploadJob::emitDescription(const QString& currentFileName) {
|
||||
QPair<QString, QString> field2;
|
||||
|
||||
if (m_totalJobs > 1) {
|
||||
field2.first = i18n("Progress");
|
||||
field2.second = i18n("Sending file %1 of %2", m_currentJobNum, m_totalJobs);
|
||||
}
|
||||
|
||||
Q_EMIT description(this, i18n("Sending to %1", Daemon::instance()->getDevice(this->m_deviceId)->name()),
|
||||
{ i18n("File"), currentFileName }, field2
|
||||
);
|
||||
}
|
84
core/backends/lan/compositeuploadjob.h
Normal file
84
core/backends/lan/compositeuploadjob.h
Normal file
|
@ -0,0 +1,84 @@
|
|||
/**
|
||||
* Copyright 2018 Erik Duisters
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or
|
||||
* modify it under the terms of the GNU General Public License as
|
||||
* published by the Free Software Foundation; either version 2 of
|
||||
* the License or (at your option) version 3 or any later version
|
||||
* accepted by the membership of KDE e.V. (or its successor approved
|
||||
* by the membership of KDE e.V.), which shall act as a proxy
|
||||
* defined in Section 14 of version 3 of the license.
|
||||
*
|
||||
* This program is distributed in the hope that it will be useful,
|
||||
* but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
* GNU General Public License for more details.
|
||||
*
|
||||
* You should have received a copy of the GNU General Public License
|
||||
* along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
*/
|
||||
|
||||
#ifndef COMPOSITEUPLOADJOB_H
|
||||
#define COMPOSITEUPLOADJOB_H
|
||||
|
||||
#include "kdeconnectcore_export.h"
|
||||
#include <KCompositeJob>
|
||||
#include "server.h"
|
||||
#include "uploadjob.h"
|
||||
|
||||
class KDECONNECTCORE_EXPORT CompositeUploadJob
|
||||
: public KCompositeJob
|
||||
{
|
||||
Q_OBJECT
|
||||
|
||||
public:
|
||||
explicit CompositeUploadJob(const QString& deviceId, bool displayNotification);
|
||||
|
||||
void start() override;
|
||||
QVariantMap transferInfo();
|
||||
bool isRunning();
|
||||
bool addSubjob(KJob* job) override;
|
||||
|
||||
private:
|
||||
bool startListening();
|
||||
void emitDescription(const QString& currentFileName);
|
||||
|
||||
protected:
|
||||
bool doKill() override;
|
||||
|
||||
private:
|
||||
enum {
|
||||
NoPortAvailable = UserDefinedError,
|
||||
SendingNetworkPacketFailed,
|
||||
SocketError,
|
||||
SslError
|
||||
};
|
||||
|
||||
Server *const m_server;
|
||||
QSslSocket *m_socket;
|
||||
quint16 m_port;
|
||||
const QString& m_deviceId;
|
||||
bool m_running;
|
||||
int m_currentJobNum;
|
||||
int m_totalJobs;
|
||||
quint64 m_currentJobSendPayloadSize;
|
||||
quint64 m_totalSendPayloadSize;
|
||||
quint64 m_totalPayloadSize;
|
||||
UploadJob *m_currentJob;
|
||||
QElapsedTimer m_timer;
|
||||
|
||||
const static quint16 MIN_PORT = 1739;
|
||||
const static quint16 MAX_PORT = 1764;
|
||||
|
||||
private Q_SLOTS:
|
||||
void newConnection();
|
||||
void socketDisconnected();
|
||||
void socketError(QAbstractSocket::SocketError socketError);
|
||||
void sslError(const QList<QSslError>& errors);
|
||||
void encrypted();
|
||||
void slotProcessedAmount(KJob *job, KJob::Unit unit, qulonglong amount);
|
||||
void slotResult(KJob *job) override;
|
||||
void startNextSubJob();
|
||||
};
|
||||
|
||||
#endif //COMPOSITEUPLOADJOB_H
|
|
@ -26,9 +26,7 @@
|
|||
#include "backends/linkprovider.h"
|
||||
#include "socketlinereader.h"
|
||||
#include "lanlinkprovider.h"
|
||||
#include <kio/global.h>
|
||||
#include <KJobTrackerInterface>
|
||||
#include <plugins/share/shareplugin.h>
|
||||
#include "plugins/share/shareplugin.h"
|
||||
|
||||
LanDeviceLink::LanDeviceLink(const QString& deviceId, LinkProvider* parent, QSslSocket* socket, ConnectionStarted connectionSource)
|
||||
: DeviceLink(deviceId, parent)
|
||||
|
@ -85,26 +83,32 @@ QString LanDeviceLink::name()
|
|||
|
||||
bool LanDeviceLink::sendPacket(NetworkPacket& np)
|
||||
{
|
||||
if (np.hasPayload()) {
|
||||
np.setPayloadTransferInfo(sendPayload(np)->transferInfo());
|
||||
if (np.payload()) {
|
||||
if (np.type() == PACKET_TYPE_SHARE_REQUEST && np.payloadSize() >= 0) {
|
||||
if (!m_compositeUploadJob || !m_compositeUploadJob->isRunning()) {
|
||||
m_compositeUploadJob = new CompositeUploadJob(deviceId(), true);
|
||||
}
|
||||
|
||||
m_compositeUploadJob->addSubjob(new UploadJob(np));
|
||||
|
||||
if (!m_compositeUploadJob->isRunning()) {
|
||||
m_compositeUploadJob->start();
|
||||
}
|
||||
} else { //Infinite stream
|
||||
CompositeUploadJob* fireAndForgetJob = new CompositeUploadJob(deviceId(), false);
|
||||
fireAndForgetJob->addSubjob(new UploadJob(np));
|
||||
fireAndForgetJob->start();
|
||||
}
|
||||
|
||||
return true;
|
||||
} else {
|
||||
int written = m_socketLineReader->write(np.serialize());
|
||||
|
||||
//Actually we can't detect if a packet is received or not. We keep TCP
|
||||
//"ESTABLISHED" connections that look legit (return true when we use them),
|
||||
//but that are actually broken (until keepalive detects that they are down).
|
||||
return (written != -1);
|
||||
}
|
||||
|
||||
int written = m_socketLineReader->write(np.serialize());
|
||||
|
||||
//Actually we can't detect if a packet is received or not. We keep TCP
|
||||
//"ESTABLISHED" connections that look legit (return true when we use them),
|
||||
//but that are actually broken (until keepalive detects that they are down).
|
||||
return (written != -1);
|
||||
}
|
||||
|
||||
UploadJob* LanDeviceLink::sendPayload(const NetworkPacket& np)
|
||||
{
|
||||
UploadJob* job = new UploadJob(np.payload(), deviceId());
|
||||
if (np.type() == PACKET_TYPE_SHARE_REQUEST && np.payloadSize() >= 0) {
|
||||
KIO::getJobTracker()->registerJob(job);
|
||||
}
|
||||
job->start();
|
||||
return job;
|
||||
}
|
||||
|
||||
void LanDeviceLink::dataReceived()
|
||||
|
|
|
@ -22,6 +22,7 @@
|
|||
#define LANDEVICELINK_H
|
||||
|
||||
#include <QObject>
|
||||
#include <QPointer>
|
||||
#include <QString>
|
||||
#include <QSslSocket>
|
||||
#include <QSslCertificate>
|
||||
|
@ -29,6 +30,7 @@
|
|||
#include <kdeconnectcore_export.h>
|
||||
#include "backends/devicelink.h"
|
||||
#include "uploadjob.h"
|
||||
#include "compositeuploadjob.h"
|
||||
|
||||
class SocketLineReader;
|
||||
|
||||
|
@ -45,7 +47,6 @@ public:
|
|||
|
||||
QString name() override;
|
||||
bool sendPacket(NetworkPacket& np) override;
|
||||
UploadJob* sendPayload(const NetworkPacket& np);
|
||||
|
||||
void userRequestsPair() override;
|
||||
void userRequestsUnpair() override;
|
||||
|
@ -63,6 +64,7 @@ private:
|
|||
SocketLineReader* m_socketLineReader;
|
||||
ConnectionStarted m_connectionSource;
|
||||
QHostAddress m_hostAddress;
|
||||
QPointer<CompositeUploadJob> m_compositeUploadJob;
|
||||
};
|
||||
|
||||
#endif
|
||||
|
|
|
@ -25,73 +25,41 @@
|
|||
#include "lanlinkprovider.h"
|
||||
#include "kdeconnectconfig.h"
|
||||
#include "core_debug.h"
|
||||
#include <device.h>
|
||||
#include <daemon.h>
|
||||
|
||||
UploadJob::UploadJob(const QSharedPointer<QIODevice>& source, const QString& deviceId)
|
||||
UploadJob::UploadJob(const NetworkPacket& networkPacket)
|
||||
: KJob()
|
||||
, m_input(source)
|
||||
, m_server(new Server(this))
|
||||
, m_networkPacket(networkPacket)
|
||||
, m_input(networkPacket.payload())
|
||||
, m_socket(nullptr)
|
||||
, m_port(0)
|
||||
, m_deviceId(deviceId) // We will use this info if link is on ssl, to send encrypted payload
|
||||
{
|
||||
connect(m_input.data(), &QIODevice::aboutToClose, this, &UploadJob::aboutToClose);
|
||||
setCapabilities(Killable);
|
||||
}
|
||||
|
||||
void UploadJob::setSocket(QSslSocket* socket)
|
||||
{
|
||||
m_socket = socket;
|
||||
m_socket->setParent(this);
|
||||
}
|
||||
|
||||
void UploadJob::start()
|
||||
{
|
||||
m_port = MIN_PORT;
|
||||
while (!m_server->listen(QHostAddress::Any, m_port)) {
|
||||
m_port++;
|
||||
if (m_port > MAX_PORT) { //No ports available?
|
||||
qCWarning(KDECONNECT_CORE) << "Error opening a port in range" << MIN_PORT << "-" << MAX_PORT;
|
||||
m_port = 0;
|
||||
setError(1);
|
||||
setErrorText(i18n("Couldn't find an available port"));
|
||||
emitResult();
|
||||
return;
|
||||
}
|
||||
}
|
||||
connect(m_server, &QTcpServer::newConnection, this, &UploadJob::newConnection);
|
||||
|
||||
Q_EMIT description(this, i18n("Sending file to %1", Daemon::instance()->getDevice(this->m_deviceId)->name()),
|
||||
{ i18nc("File transfer origin", "From"), m_input.staticCast<QFile>().data()->fileName() }
|
||||
);
|
||||
}
|
||||
|
||||
void UploadJob::newConnection()
|
||||
{
|
||||
if (!m_input->open(QIODevice::ReadOnly)) {
|
||||
if (!m_input->open(QIODevice::ReadOnly)) {
|
||||
qCWarning(KDECONNECT_CORE) << "error when opening the input to upload";
|
||||
return; //TODO: Handle error, clean up...
|
||||
}
|
||||
|
||||
m_socket = m_server->nextPendingConnection();
|
||||
m_socket->setParent(this);
|
||||
connect(m_socket, &QSslSocket::disconnected, this, &UploadJob::cleanup);
|
||||
connect(m_socket, SIGNAL(error(QAbstractSocket::SocketError)), this, SLOT(socketFailed(QAbstractSocket::SocketError)));
|
||||
connect(m_socket, SIGNAL(sslErrors(QList<QSslError>)), this, SLOT(sslErrors(QList<QSslError>)));
|
||||
connect(m_socket, &QSslSocket::encrypted, this, &UploadJob::startUploading);
|
||||
if (!m_socket) {
|
||||
qCWarning(KDECONNECT_CORE) << "you must call setSocket() before calling start()";
|
||||
return;
|
||||
}
|
||||
|
||||
LanLinkProvider::configureSslSocket(m_socket, m_deviceId, true);
|
||||
connect(m_input.data(), &QIODevice::aboutToClose, this, &UploadJob::aboutToClose);
|
||||
|
||||
m_socket->startServerEncryption();
|
||||
}
|
||||
|
||||
void UploadJob::startUploading()
|
||||
{
|
||||
bytesUploaded = 0;
|
||||
setProcessedAmount(Bytes, bytesUploaded);
|
||||
setTotalAmount(Bytes, m_input.data()->size());
|
||||
|
||||
connect(m_socket, &QSslSocket::encryptedBytesWritten, this, &UploadJob::encryptedBytesWritten);
|
||||
|
||||
if (!m_timer.isValid()) {
|
||||
m_timer.start();
|
||||
}
|
||||
|
||||
uploadNextPacket();
|
||||
}
|
||||
|
||||
|
@ -102,10 +70,6 @@ void UploadJob::uploadNextPacket()
|
|||
if ( bytesAvailable > 0) {
|
||||
qint64 bytesToSend = qMin(m_input->bytesAvailable(), (qint64)4096);
|
||||
bytesUploading = m_socket->write(m_input->read(bytesToSend));
|
||||
|
||||
if (bytesUploading < 0) {
|
||||
qCWarning(KDECONNECT_CORE) << "error when writing data to upload" << bytesToSend << m_input->bytesAvailable();
|
||||
}
|
||||
}
|
||||
|
||||
if (bytesAvailable <= 0 || bytesUploading < 0) {
|
||||
|
@ -118,59 +82,27 @@ void UploadJob::encryptedBytesWritten(qint64 bytes)
|
|||
{
|
||||
Q_UNUSED(bytes);
|
||||
|
||||
bytesUploaded += bytesUploading;
|
||||
|
||||
if (m_socket->encryptedBytesToWrite() == 0) {
|
||||
bytesUploaded += bytesUploading;
|
||||
setProcessedAmount(Bytes, bytesUploaded);
|
||||
|
||||
const auto elapsed = m_timer.elapsed();
|
||||
if (elapsed > 0) {
|
||||
emitSpeed((1000 * bytesUploaded) / elapsed);
|
||||
}
|
||||
|
||||
uploadNextPacket();
|
||||
}
|
||||
}
|
||||
|
||||
void UploadJob::aboutToClose()
|
||||
{
|
||||
qWarning() << "aboutToClose()";
|
||||
m_socket->disconnectFromHost();
|
||||
}
|
||||
|
||||
void UploadJob::cleanup()
|
||||
{
|
||||
qWarning() << "cleanup()";
|
||||
m_socket->close();
|
||||
emitResult();
|
||||
}
|
||||
|
||||
QVariantMap UploadJob::transferInfo()
|
||||
{
|
||||
Q_ASSERT(m_port != 0);
|
||||
return {{"port", m_port}};
|
||||
}
|
||||
bool UploadJob::stop() {
|
||||
m_input->close();
|
||||
|
||||
void UploadJob::socketFailed(QAbstractSocket::SocketError error)
|
||||
{
|
||||
qWarning() << "socketFailed() " << error;
|
||||
m_socket->close();
|
||||
setError(2);
|
||||
emitResult();
|
||||
}
|
||||
|
||||
void UploadJob::sslErrors(const QList<QSslError>& errors)
|
||||
{
|
||||
qWarning() << "sslErrors() " << errors;
|
||||
setError(1);
|
||||
emitResult();
|
||||
m_socket->close();
|
||||
}
|
||||
|
||||
bool UploadJob::doKill()
|
||||
{
|
||||
if (m_input) {
|
||||
m_input->close();
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
const NetworkPacket UploadJob::getNetworkPacket()
|
||||
{
|
||||
return m_networkPacket;
|
||||
}
|
||||
|
|
|
@ -25,49 +25,37 @@
|
|||
|
||||
#include <QIODevice>
|
||||
#include <QVariantMap>
|
||||
#include <QSharedPointer>
|
||||
#include <QSslSocket>
|
||||
#include "server.h"
|
||||
#include <QElapsedTimer>
|
||||
#include <QTimer>
|
||||
#include <networkpacket.h>
|
||||
|
||||
class KDECONNECTCORE_EXPORT UploadJob
|
||||
: public KJob
|
||||
{
|
||||
Q_OBJECT
|
||||
public:
|
||||
explicit UploadJob(const QSharedPointer<QIODevice>& source, const QString& deviceId);
|
||||
explicit UploadJob(const NetworkPacket& networkPacket);
|
||||
|
||||
void setSocket(QSslSocket* socket);
|
||||
void start() override;
|
||||
|
||||
QVariantMap transferInfo();
|
||||
bool stop();
|
||||
const NetworkPacket getNetworkPacket();
|
||||
|
||||
private:
|
||||
const QSharedPointer<QIODevice> m_input;
|
||||
Server * const m_server;
|
||||
const NetworkPacket m_networkPacket;
|
||||
QSharedPointer<QIODevice> m_input;
|
||||
QSslSocket* m_socket;
|
||||
quint16 m_port;
|
||||
const QString m_deviceId;
|
||||
qint64 bytesUploading;
|
||||
qint64 bytesUploaded;
|
||||
QElapsedTimer m_timer;
|
||||
|
||||
const static quint16 MIN_PORT = 1739;
|
||||
const static quint16 MAX_PORT = 1764;
|
||||
|
||||
protected:
|
||||
bool doKill() override;
|
||||
|
||||
private Q_SLOTS:
|
||||
void startUploading();
|
||||
void uploadNextPacket();
|
||||
void newConnection();
|
||||
void encryptedBytesWritten(qint64 bytes);
|
||||
void aboutToClose();
|
||||
void cleanup();
|
||||
|
||||
void socketFailed(QAbstractSocket::SocketError);
|
||||
void sslErrors(const QList<QSslError>& errors);
|
||||
};
|
||||
|
||||
#endif // UPLOADJOB_H
|
||||
|
|
|
@ -178,6 +178,12 @@ void SharePlugin::shareUrl(const QUrl& url)
|
|||
sendPacket(packet);
|
||||
}
|
||||
|
||||
void SharePlugin::shareUrls(const QStringList& urls) {
|
||||
for(const QString url : urls) {
|
||||
shareUrl(QUrl(url));
|
||||
}
|
||||
}
|
||||
|
||||
void SharePlugin::shareText(const QString& text)
|
||||
{
|
||||
NetworkPacket packet(PACKET_TYPE_SHARE_REQUEST);
|
||||
|
|
|
@ -38,6 +38,7 @@ public:
|
|||
|
||||
///Helper method, QDBus won't recognize QUrl
|
||||
Q_SCRIPTABLE void shareUrl(const QString& url) { shareUrl(QUrl(url)); }
|
||||
Q_SCRIPTABLE void shareUrls(const QStringList& urls);
|
||||
Q_SCRIPTABLE void shareText(const QString& text);
|
||||
Q_SCRIPTABLE void openFile(const QString& file) { openFile(QUrl(file)); }
|
||||
|
||||
|
@ -55,7 +56,6 @@ Q_SIGNALS:
|
|||
private:
|
||||
void shareUrl(const QUrl& url);
|
||||
void openFile(const QUrl& url);
|
||||
|
||||
QUrl destinationDir() const;
|
||||
QUrl getFileDestination(const QString filename) const;
|
||||
};
|
||||
|
|
|
@ -37,6 +37,8 @@
|
|||
#include <backends/pairinghandler.h>
|
||||
#include "kdeconnect-version.h"
|
||||
#include "testdaemon.h"
|
||||
#include <plugins/share/shareplugin.h>
|
||||
#include <backends/lan/compositeuploadjob.h>
|
||||
|
||||
class TestSendFile : public QObject
|
||||
{
|
||||
|
@ -102,17 +104,18 @@ class TestSendFile : public QObject
|
|||
kcc->setDeviceProperty(deviceId, QStringLiteral("certificate"), QString::fromLatin1(kcc->certificate().toPem())); // Using same certificate from kcc, instead of generating
|
||||
|
||||
QSharedPointer<QFile> f(new QFile(aFile));
|
||||
UploadJob* uj = new UploadJob(f, deviceId);
|
||||
QSignalSpy spyUpload(uj, &KJob::result);
|
||||
uj->start();
|
||||
NetworkPacket np(PACKET_TYPE_SHARE_REQUEST);
|
||||
np.setPayload(f, aFile.size());
|
||||
CompositeUploadJob* job = new CompositeUploadJob(deviceId, false);
|
||||
UploadJob* uj = new UploadJob(np);
|
||||
job->addSubjob(uj);
|
||||
|
||||
auto info = uj->transferInfo();
|
||||
info.insert(QStringLiteral("deviceId"), deviceId);
|
||||
info.insert(QStringLiteral("size"), aFile.size());
|
||||
QSignalSpy spyUpload(job, &KJob::result);
|
||||
job->start();
|
||||
|
||||
f->open(QIODevice::ReadWrite);
|
||||
|
||||
FileTransferJob* ft = new FileTransferJob(f, uj->transferInfo()[QStringLiteral("size")].toInt(), QUrl::fromLocalFile(destFile));
|
||||
FileTransferJob* ft = new FileTransferJob(f, aFile.size(), QUrl::fromLocalFile(destFile));
|
||||
|
||||
QSignalSpy spyTransfer(ft, &KJob::result);
|
||||
|
||||
|
|
Loading…
Reference in a new issue