83 uint32_t chunk_size = 64 * 1024;
84 uint32_t window_bytes = 4 * 1024 * 1024;
85 uint32_t progress_interval = 256 * 1024;
86 uint32_t transfer_timeout_secs = 60;
91 uint32_t offer_timeout_secs = 300;
92 uint32_t worker_threads = 4;
93 uint32_t disk_threads = 4;
94 bool verify_integrity =
true;
95 std::string temp_directory =
".";
98 enum class Status { Pending, Active, Paused, Completed, Failed, Cancelled };
113 bool is_directory =
false;
123 uint64_t bytes_transferred = 0;
124 uint64_t total_bytes = 0;
125 uint32_t files_completed = 0;
126 uint32_t total_files = 0;
128 double transfer_rate_bps = 0.0;
129 double average_rate_bps = 0.0;
130 std::chrono::milliseconds elapsed{0};
131 std::chrono::milliseconds estimated_time_remaining{0};
135 if (total_bytes == 0)
return status == Status::Completed ? 100.0 : 0.0;
136 return static_cast<double>(bytes_transferred) /
static_cast<double>(total_bytes) * 100.0;
142 uint64_t bytes_sent = 0, bytes_received = 0;
143 uint64_t completed = 0, failed = 0;
153 std::function<void(
const PeerId& peer, uint64_t
id,
bool success,
const std::string& path)>;
170 void accept(
const PeerId& from, uint64_t
id,
const std::string& dest_path);
196 const std::string& partial_path);
220 using clock = std::chrono::steady_clock;
221 clock::time_point start{};
222 clock::time_point mark{};
223 uint64_t mark_bytes = 0;
224 double rate_bps = 0.0;
226 void sample(uint64_t bytes, clock::time_point now) {
227 constexpr int64_t kMinMs = 250, kMaxMs = 2000;
228 constexpr double kAlpha = 0.4;
229 if (start == clock::time_point{}) { start = mark = now; mark_bytes = bytes;
return; }
230 const int64_t dt = std::chrono::duration_cast<std::chrono::milliseconds>(now - mark).count();
231 if (dt < kMinMs)
return;
232 if (dt > kMaxMs) { mark = now; mark_bytes = bytes;
return; }
233 const double inst =
static_cast<double>(bytes - mark_bytes) * 1000.0 /
static_cast<double>(dt);
234 rate_bps = rate_bps <= 0.0 ? inst : kAlpha * inst + (1.0 - kAlpha) * rate_bps;
235 mark = now; mark_bytes = bytes;
238 void fill(Progress& p, uint64_t bytes, uint64_t total, clock::time_point now)
const {
239 if (start == clock::time_point{})
return;
240 const auto ms = std::chrono::duration_cast<std::chrono::milliseconds>(now - start);
242 p.transfer_rate_bps = rate_bps;
244 p.average_rate_bps =
static_cast<double>(bytes) * 1000.0 /
static_cast<double>(ms.count());
245 if (rate_bps > 1.0 && total > bytes)
246 p.estimated_time_remaining = std::chrono::milliseconds(
247 static_cast<int64_t
>(
static_cast<double>(total - bytes) / rate_bps * 1000.0));
258 bool is_directory =
false;
259 std::vector<FileEntry> files;
260 std::vector<std::string> sources;
261 uint64_t total_bytes = 0;
264 std::condition_variable cv;
266 uint64_t cur_offset = 0;
269 size_t hash_file = SIZE_MAX;
270 uint64_t bytes_done = 0;
272 uint32_t files_done = 0;
273 Status status = Status::Pending;
274 bool worker_active =
false;
275 bool finished =
false;
276 rats_sha256_context_t hash{};
277 std::chrono::steady_clock::time_point last_activity{};
281 struct IncomingFile {
282 std::string relative_path;
284 std::string final_path;
285 std::string temp_path;
286 uint64_t enqueued = 0;
287 uint64_t received = 0;
288 uint64_t resume_offset = 0;
289 bool keep_partial =
false;
290 bool temp_created =
false;
291 bool finalized =
false;
301 bool is_file_end =
false;
302 uint8_t sha[RATS_SHA256_HASH_SIZE]{};
309 bool is_directory =
false;
310 std::string dest_root;
311 std::vector<IncomingFile> files;
314 size_t recv_file = 0;
315 uint64_t bytes_done = 0;
316 uint64_t last_ack = 0;
317 uint32_t files_done = 0;
318 Status status = Status::Pending;
319 bool finished =
false;
320 std::chrono::steady_clock::time_point last_activity{};
324 std::queue<WriteJob> wq;
325 uint64_t queued_bytes = 0;
326 bool scheduled =
false;
328 size_t out_idx = SIZE_MAX;
329 int hashing_file = -1;
330 rats_sha256_context_t hash{};
336 void on_message(
const Peer& peer, ByteView payload);
337 void handle_offer(
const PeerId& from, uint64_t
id,
bool is_dir, uint64_t total,
338 std::string name, std::vector<FileEntry> files);
339 void handle_chunk(
const PeerId& from, uint64_t
id, uint32_t fidx, uint64_t offset,
341 void handle_file_end(
const PeerId& from, uint64_t
id, uint32_t fidx,
const uint8_t* sha);
344 uint64_t start_send(std::shared_ptr<Outgoing> t);
345 void queue_send(uint64_t
id);
347 void run_send(
const std::shared_ptr<Outgoing>& t);
354 void schedule_writer(
const std::shared_ptr<Incoming>& t);
355 void disk_worker_loop();
359 bool begin_hash(
const std::shared_ptr<Incoming>& t, uint32_t fidx,
const std::string& path, uint64_t prefix);
360 void drain_writes(
const std::shared_ptr<Incoming>& t);
361 void process_data(
const std::shared_ptr<Incoming>& t, WriteJob& job);
362 void process_file_end(
const std::shared_ptr<Incoming>& t, WriteJob& job);
366 bool accept_impl(
const PeerId& from, uint64_t
id,
const std::string& dest_path,
367 const std::string& partial_path);
370 void maintenance_loop();
371 void finish_outgoing(
const std::shared_ptr<Outgoing>& t,
bool success);
372 void finish_incoming(
const std::shared_ptr<Incoming>& t,
bool success,
const std::string& error);
373 void emit_progress(
const std::shared_ptr<Outgoing>& t);
374 void emit_progress(
const std::shared_ptr<Incoming>& t);
381 std::shared_ptr<Outgoing> find_outgoing(
const PeerId& peer, uint64_t
id)
const;
384 std::shared_ptr<Outgoing> find_own_outgoing(uint64_t
id)
const;
385 std::shared_ptr<Incoming> find_incoming(
const PeerId& peer, uint64_t
id)
const;
387 void send_to(
const PeerId& peer,
const Bytes& msg) {
388 if (network_) network_->send(peer, MessageType::FileChunk, ByteView(msg));
390 void send_simple(
const PeerId& peer, uint8_t op, uint64_t
id);
391 void send_complete(
const PeerId& peer, uint64_t
id,
bool ok);
393 PeerNetwork* network_ =
nullptr;
395 std::atomic<uint64_t> next_id_{1};
396 std::atomic<bool> running_{
false};
398 OfferHandler offer_handler_;
399 ProgressHandler progress_handler_;
400 CompleteHandler complete_handler_;
402 mutable std::mutex mutex_;
403 std::unordered_map<uint64_t, std::shared_ptr<Outgoing>> outgoing_;
404 std::unordered_map<PeerId, std::unordered_map<uint64_t, std::shared_ptr<Incoming>>,
405 PeerId::Hash> incoming_;
408 std::vector<std::thread> workers_;
409 std::mutex queue_mutex_;
410 std::condition_variable queue_cv_;
411 std::queue<uint64_t> send_queue_;
416 std::vector<std::thread> disk_workers_;
417 std::mutex disk_mutex_;
418 std::condition_variable disk_cv_;
419 std::queue<std::shared_ptr<Incoming>> disk_ready_;
422 std::thread maintenance_thread_;
423 std::mutex maintenance_mutex_;
424 std::condition_variable maintenance_cv_;
426 mutable std::mutex stats_mutex_;