Repository navigation
Expand file tree
/
Copy pathdedicated_file_loader.h
More file actions
219 lines (178 loc) · 5.1 KB
/
Copy pathdedicated_file_loader.h
File metadata and controls
219 lines (178 loc) · 5.1 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
/*
This file is part of Telegram Desktop,
the official desktop application for the Telegram messaging service.
For license and copyright information please follow this link:
https://github.com/telegramdesktop/tdesktop/blob/master/LEGAL
*/
#pragma once
#include "mtproto/mtp_instance.h"
namespace Main {
class Session;
} // namespace Main
namespace MTP {
class WeakInstance final : private QObject {
public:
explicit WeakInstance(base::weak_ptr<Main::Session> session);
template <typename Request>
void send(
const Request &request,
Fn<void(const typename Request::ResponseType &result)> done,
Fn<void(const Error &error)> fail,
ShiftedDcId dcId = 0);
[[nodiscard]] base::weak_ptr<Main::Session> session() const;
[[nodiscard]] bool valid() const;
[[nodiscard]] Instance *instance() const;
~WeakInstance();
private:
void die();
bool removeRequest(mtpRequestId requestId);
void reportUnavailable(Fn<void(const Error &error)> callback);
base::weak_ptr<Main::Session> _session;
Instance *_instance = nullptr;
std::map<mtpRequestId, Fn<void(const Error &)>> _requests;
rpl::lifetime _lifetime;
};
class AbstractDedicatedLoader : public base::has_weak_ptr {
public:
AbstractDedicatedLoader(const QString &filepath, int chunkSize);
static constexpr auto kChunkSize = 128 * 1024;
static constexpr auto kMaxFileSize = 256 * 1024 * 1024;
struct Progress {
int64 already = 0;
int64 size = 0;
bool percent = false;
inline bool operator<(const Progress &other) const {
return (already < other.already)
|| (already == other.already && size < other.size);
}
inline bool operator==(const Progress &other) const {
return (already == other.already) && (size == other.size);
}
};
void start();
void wipeFolder();
void wipeOutput();
int64 alreadySize() const;
int64 totalSize() const;
bool preferPercent() const;
rpl::producer<Progress> progress() const;
rpl::producer<QString> ready() const;
rpl::producer<> failed() const;
rpl::lifetime &lifetime();
virtual ~AbstractDedicatedLoader() = default;
protected:
void threadSafeFailed();
void threadSafeProgress(Progress progress);
void threadSafeReady();
// Single threaded.
void writeChunk(bytes::const_span data, int64 totalSize);
private:
virtual void startLoading() = 0;
bool validateOutput();
QString _filepath;
int _chunkSize = 0;
QFile _output;
int64 _alreadySize = 0;
int64 _totalSize = 0;
bool _preferPercent = false;
mutable QMutex _sizesMutex;
rpl::event_stream<Progress> _progress;
rpl::event_stream<QString> _ready;
rpl::event_stream<> _failed;
rpl::lifetime _lifetime;
};
class DedicatedLoader : public AbstractDedicatedLoader {
public:
struct Location {
QString username;
// When channelId is set the channel is used directly with the
// cached accessHash instead of resolving username.
uint64 channelId = 0;
uint64 accessHash = 0;
int32 postId = 0;
};
struct File {
QString name;
int64 size = 0;
DcId dcId = 0;
MTPInputFileLocation location;
};
DedicatedLoader(
base::weak_ptr<Main::Session> session,
const QString &folder,
const File &file);
private:
struct Request {
int64 offset = 0;
QByteArray bytes;
};
void startLoading() override;
void sendRequest();
void gotPart(int offset, const MTPupload_File &result);
Fn<void(const Error &)> failHandler();
static constexpr auto kRequestsCount = 2;
static constexpr auto kNextRequestDelay = crl::time(20);
std::deque<Request> _requests;
int64 _size = 0;
int64 _offset = 0;
DcId _dcId = 0;
MTPInputFileLocation _location;
WeakInstance _mtp;
};
void ResolveChannel(
not_null<MTP::WeakInstance*> mtp,
const QString &username,
Fn<void(const MTPInputChannel &channel)> done,
Fn<void()> fail);
// With a non-zero messageId only the message with that exact id counts,
// the server may answer a getMessages request with a different message.
std::optional<MTPMessage> GetMessagesElement(
const MTPmessages_Messages &list,
int messageId = 0);
void StartDedicatedLoader(
not_null<MTP::WeakInstance*> mtp,
const DedicatedLoader::Location &location,
const QString &folder,
Fn<void(std::unique_ptr<DedicatedLoader>)> ready);
template <typename Request>
void WeakInstance::send(
const Request &request,
Fn<void(const typename Request::ResponseType &result)> done,
Fn<void(const Error &error)> fail,
MTP::ShiftedDcId dcId) {
using Result = typename Request::ResponseType;
if (!valid()) {
reportUnavailable(fail);
return;
}
const auto onDone = crl::guard((QObject*)this, [=](
const Response &response) {
auto result = Result();
auto from = response.reply.constData();
if (!result.read(from, from + response.reply.size())) {
return false;
}
if (removeRequest(response.requestId)) {
done(result);
}
return true;
});
const auto onFail = crl::guard((QObject*)this, [=](
const Error &error,
const Response &response) {
if (MTP::IsDefaultHandledError(error)) {
return false;
}
if (removeRequest(response.requestId)) {
fail(error);
}
return true;
});
const auto requestId = _instance->send(
request,
std::move(onDone),
std::move(onFail),
dcId);
_requests.emplace(requestId, fail);
}
} // namespace MTP