77 uint32_t chunk_size = 64 * 1024;
78 uint32_t window_bytes = 4 * 1024 * 1024;
79 uint32_t progress_interval = 256 * 1024;
80 uint32_t transfer_timeout_secs = 60;
81 uint32_t worker_threads = 4;
82 uint32_t disk_threads = 4;
83 bool verify_integrity =
true;
84 std::string temp_directory =
".";
87 enum class Status { Pending, Active, Paused, Completed, Failed, Cancelled };
102 bool is_directory =
false;
112 uint64_t bytes_transferred = 0;
113 uint64_t total_bytes = 0;
114 uint32_t files_completed = 0;
115 uint32_t total_files = 0;
117 double transfer_rate_bps = 0.0;
118 double average_rate_bps = 0.0;
119 std::chrono::milliseconds elapsed{0};
120 std::chrono::milliseconds estimated_time_remaining{0};
124 if (total_bytes == 0)
return status == Status::Completed ? 100.0 : 0.0;
125 return static_cast<double>(bytes_transferred) /
static_cast<double>(total_bytes) * 100.0;
131 uint64_t bytes_sent = 0, bytes_received = 0;
132 uint64_t completed = 0, failed = 0;
137 using CompleteHandler = std::function<void(uint64_t
id,
bool success,
const std::string& path)>;
154 void accept(
const PeerId& from, uint64_t
id,
const std::string& dest_path);
179 using clock = std::chrono::steady_clock;
180 clock::time_point start{};
181 clock::time_point mark{};
182 uint64_t mark_bytes = 0;
183 double rate_bps = 0.0;
185 void sample(uint64_t bytes, clock::time_point now) {
186 constexpr int64_t kMinMs = 250, kMaxMs = 2000;
187 constexpr double kAlpha = 0.4;
188 if (start == clock::time_point{}) { start = mark = now; mark_bytes = bytes;
return; }
189 const int64_t dt = std::chrono::duration_cast<std::chrono::milliseconds>(now - mark).count();
190 if (dt < kMinMs)
return;
191 if (dt > kMaxMs) { mark = now; mark_bytes = bytes;
return; }
192 const double inst =
static_cast<double>(bytes - mark_bytes) * 1000.0 /
static_cast<double>(dt);
193 rate_bps = rate_bps <= 0.0 ? inst : kAlpha * inst + (1.0 - kAlpha) * rate_bps;
194 mark = now; mark_bytes = bytes;
197 void fill(Progress& p, uint64_t bytes, uint64_t total, clock::time_point now)
const {
198 if (start == clock::time_point{})
return;
199 const auto ms = std::chrono::duration_cast<std::chrono::milliseconds>(now - start);
201 p.transfer_rate_bps = rate_bps;
203 p.average_rate_bps =
static_cast<double>(bytes) * 1000.0 /
static_cast<double>(ms.count());
204 if (rate_bps > 1.0 && total > bytes)
205 p.estimated_time_remaining = std::chrono::milliseconds(
206 static_cast<int64_t
>(
static_cast<double>(total - bytes) / rate_bps * 1000.0));
217 bool is_directory =
false;
218 std::vector<FileEntry> files;
219 std::vector<std::string> sources;
220 uint64_t total_bytes = 0;
223 std::condition_variable cv;
225 uint64_t cur_offset = 0;
226 uint64_t bytes_done = 0;
228 uint32_t files_done = 0;
229 Status status = Status::Pending;
230 bool worker_active =
false;
231 bool finished =
false;
232 sha256_context_t hash{};
233 std::chrono::steady_clock::time_point last_activity{};
237 struct IncomingFile {
238 std::string relative_path;
240 std::string final_path;
241 std::string temp_path;
242 uint64_t enqueued = 0;
243 uint64_t received = 0;
244 bool temp_created =
false;
245 bool finalized =
false;
255 bool is_file_end =
false;
256 uint8_t sha[SHA256_HASH_SIZE]{};
263 bool is_directory =
false;
264 std::string dest_root;
265 std::vector<IncomingFile> files;
268 size_t recv_file = 0;
269 uint64_t bytes_done = 0;
270 uint64_t last_ack = 0;
271 uint32_t files_done = 0;
272 Status status = Status::Pending;
273 bool finished =
false;
274 std::chrono::steady_clock::time_point last_activity{};
278 std::queue<WriteJob> wq;
279 uint64_t queued_bytes = 0;
280 bool scheduled =
false;
282 size_t out_idx = SIZE_MAX;
283 int hashing_file = -1;
284 sha256_context_t hash{};
290 void on_message(
const Peer& peer, ByteView payload);
291 void handle_offer(
const PeerId& from, uint64_t
id,
bool is_dir, uint64_t total,
292 std::string name, std::vector<FileEntry> files);
293 void handle_chunk(
const PeerId& from, uint64_t
id, uint32_t fidx, uint64_t offset,
295 void handle_file_end(
const PeerId& from, uint64_t
id, uint32_t fidx,
const uint8_t* sha);
298 uint64_t start_send(std::shared_ptr<Outgoing> t);
299 void queue_send(uint64_t
id);
301 void run_send(
const std::shared_ptr<Outgoing>& t);
308 void schedule_writer(
const std::shared_ptr<Incoming>& t);
309 void disk_worker_loop();
310 void drain_writes(
const std::shared_ptr<Incoming>& t);
311 void process_data(
const std::shared_ptr<Incoming>& t, WriteJob& job);
312 void process_file_end(
const std::shared_ptr<Incoming>& t, WriteJob& job);
315 void maintenance_loop();
316 void finish_outgoing(
const std::shared_ptr<Outgoing>& t,
bool success);
317 void finish_incoming(
const std::shared_ptr<Incoming>& t,
bool success,
const std::string& error);
318 void emit_progress(
const std::shared_ptr<Outgoing>& t);
319 void emit_progress(
const std::shared_ptr<Incoming>& t);
321 std::shared_ptr<Outgoing> find_outgoing(uint64_t
id)
const;
322 std::shared_ptr<Incoming> find_incoming(
const PeerId& peer, uint64_t
id)
const;
324 void send_to(
const PeerId& peer,
const Bytes& msg) {
325 if (network_) network_->send(peer, MessageType::FileChunk, ByteView(msg));
327 void send_simple(
const PeerId& peer, uint8_t op, uint64_t
id);
328 void send_complete(
const PeerId& peer, uint64_t
id,
bool ok);
330 PeerNetwork* network_ =
nullptr;
332 std::atomic<uint64_t> next_id_{1};
333 std::atomic<bool> running_{
false};
335 OfferHandler offer_handler_;
336 ProgressHandler progress_handler_;
337 CompleteHandler complete_handler_;
339 mutable std::mutex mutex_;
340 std::unordered_map<uint64_t, std::shared_ptr<Outgoing>> outgoing_;
341 std::unordered_map<PeerId, std::unordered_map<uint64_t, std::shared_ptr<Incoming>>,
342 PeerId::Hash> incoming_;
345 std::vector<std::thread> workers_;
346 std::mutex queue_mutex_;
347 std::condition_variable queue_cv_;
348 std::queue<uint64_t> send_queue_;
353 std::vector<std::thread> disk_workers_;
354 std::mutex disk_mutex_;
355 std::condition_variable disk_cv_;
356 std::queue<std::shared_ptr<Incoming>> disk_ready_;
359 std::thread maintenance_thread_;
360 std::mutex maintenance_mutex_;
361 std::condition_variable maintenance_cv_;
363 mutable std::mutex stats_mutex_;
One file inside a transfer (a single-file transfer has exactly one).
std::string relative_path
POSIX path relative to the transfer root.