/src/rocksdb/file/sst_file_manager_impl.cc
Line | Count | Source |
1 | | // Copyright (c) 2011-present, Facebook, Inc. All rights reserved. |
2 | | // This source code is licensed under both the GPLv2 (found in the |
3 | | // COPYING file in the root directory) and Apache 2.0 License |
4 | | // (found in the LICENSE.Apache file in the root directory). |
5 | | |
6 | | #include "file/sst_file_manager_impl.h" |
7 | | |
8 | | #include <cinttypes> |
9 | | #include <vector> |
10 | | |
11 | | #include "db/db_impl/db_impl.h" |
12 | | #include "logging/logging.h" |
13 | | #include "port/port.h" |
14 | | #include "rocksdb/env.h" |
15 | | #include "rocksdb/sst_file_manager.h" |
16 | | #include "test_util/sync_point.h" |
17 | | #include "util/mutexlock.h" |
18 | | |
19 | | namespace ROCKSDB_NAMESPACE { |
20 | | |
21 | | SstFileManagerImpl::SstFileManagerImpl( |
22 | | const std::shared_ptr<SystemClock>& clock, |
23 | | const std::shared_ptr<FileSystem>& fs, |
24 | | const std::shared_ptr<Logger>& logger, int64_t rate_bytes_per_sec, |
25 | | double max_trash_db_ratio, uint64_t bytes_max_delete_chunk) |
26 | 61.7k | : clock_(clock), |
27 | 61.7k | fs_(fs), |
28 | 61.7k | logger_(logger), |
29 | 61.7k | total_files_size_(0), |
30 | 61.7k | compaction_buffer_size_(0), |
31 | 61.7k | cur_compactions_reserved_size_(0), |
32 | 61.7k | max_allowed_space_(0), |
33 | 61.7k | delete_scheduler_(clock_.get(), fs_.get(), rate_bytes_per_sec, |
34 | 61.7k | logger.get(), this, max_trash_db_ratio, |
35 | 61.7k | bytes_max_delete_chunk), |
36 | 61.7k | cv_(&mu_), |
37 | 61.7k | closing_(false), |
38 | 61.7k | bg_thread_(nullptr), |
39 | 61.7k | reserved_disk_buffer_(0), |
40 | 61.7k | free_space_trigger_(0), |
41 | 61.7k | cur_instance_(nullptr) {} |
42 | | |
43 | 61.7k | SstFileManagerImpl::~SstFileManagerImpl() { |
44 | 61.7k | Close(); |
45 | 61.7k | bg_err_.PermitUncheckedError(); |
46 | 61.7k | } |
47 | | |
48 | 115k | void SstFileManagerImpl::Close() { |
49 | 115k | { |
50 | 115k | MutexLock l(&mu_); |
51 | 115k | if (closing_) { |
52 | 53.3k | return; |
53 | 53.3k | } |
54 | 61.7k | closing_ = true; |
55 | 61.7k | cv_.SignalAll(); |
56 | 61.7k | } |
57 | 61.7k | if (bg_thread_) { |
58 | 0 | bg_thread_->join(); |
59 | 0 | } |
60 | 61.7k | } |
61 | | |
62 | 4.94k | Status SstFileManagerImpl::OnAddFile(const std::string& file_path) { |
63 | 4.94k | uint64_t file_size; |
64 | 4.94k | Status s = fs_->GetFileSize(file_path, IOOptions(), &file_size, nullptr); |
65 | 4.94k | if (s.ok()) { |
66 | 4.94k | MutexLock l(&mu_); |
67 | 4.94k | OnAddFileImpl(file_path, file_size); |
68 | 4.94k | } |
69 | 4.94k | TEST_SYNC_POINT_CALLBACK("SstFileManagerImpl::OnAddFile", |
70 | 4.94k | const_cast<std::string*>(&file_path)); |
71 | 4.94k | return s; |
72 | 4.94k | } |
73 | | |
74 | | Status SstFileManagerImpl::OnAddFile(const std::string& file_path, |
75 | 119k | uint64_t file_size) { |
76 | 119k | MutexLock l(&mu_); |
77 | 119k | OnAddFileImpl(file_path, file_size); |
78 | 119k | TEST_SYNC_POINT_CALLBACK("SstFileManagerImpl::OnAddFile", |
79 | 119k | const_cast<std::string*>(&file_path)); |
80 | 119k | return Status::OK(); |
81 | 119k | } |
82 | | |
83 | 57.0k | Status SstFileManagerImpl::OnDeleteFile(const std::string& file_path) { |
84 | 57.0k | { |
85 | 57.0k | MutexLock l(&mu_); |
86 | 57.0k | OnDeleteFileImpl(file_path); |
87 | 57.0k | } |
88 | 57.0k | TEST_SYNC_POINT_CALLBACK("SstFileManagerImpl::OnDeleteFile", |
89 | 57.0k | const_cast<std::string*>(&file_path)); |
90 | 57.0k | return Status::OK(); |
91 | 57.0k | } |
92 | | |
93 | 9.07k | void SstFileManagerImpl::OnCompactionCompletion(Compaction* c) { |
94 | 9.07k | MutexLock l(&mu_); |
95 | 9.07k | uint64_t size_added_by_compaction = 0; |
96 | 21.2k | for (size_t i = 0; i < c->num_input_levels(); i++) { |
97 | 41.8k | for (size_t j = 0; j < c->num_input_files(i); j++) { |
98 | 29.6k | FileMetaData* filemeta = c->input(i, j); |
99 | 29.6k | size_added_by_compaction += filemeta->fd.GetFileSize(); |
100 | 29.6k | } |
101 | 12.1k | } |
102 | 9.07k | assert(cur_compactions_reserved_size_ >= size_added_by_compaction); |
103 | 9.07k | cur_compactions_reserved_size_ -= size_added_by_compaction; |
104 | 9.07k | } |
105 | | |
106 | | Status SstFileManagerImpl::OnMoveFile(const std::string& old_path, |
107 | | const std::string& new_path, |
108 | 0 | uint64_t* file_size) { |
109 | 0 | { |
110 | 0 | MutexLock l(&mu_); |
111 | 0 | if (file_size != nullptr) { |
112 | 0 | *file_size = tracked_files_[old_path]; |
113 | 0 | } |
114 | 0 | OnAddFileImpl(new_path, tracked_files_[old_path]); |
115 | 0 | OnDeleteFileImpl(old_path); |
116 | 0 | } |
117 | 0 | TEST_SYNC_POINT("SstFileManagerImpl::OnMoveFile"); |
118 | 0 | return Status::OK(); |
119 | 0 | } |
120 | | |
121 | 0 | Status SstFileManagerImpl::OnUntrackFile(const std::string& file_path) { |
122 | 0 | { |
123 | 0 | MutexLock l(&mu_); |
124 | 0 | OnDeleteFileImpl(file_path); |
125 | 0 | } |
126 | 0 | TEST_SYNC_POINT_CALLBACK("SstFileManagerImpl::OnUntrackFile", |
127 | 0 | const_cast<std::string*>(&file_path)); |
128 | 0 | return Status::OK(); |
129 | 0 | } |
130 | | |
131 | 0 | void SstFileManagerImpl::SetMaxAllowedSpaceUsage(uint64_t max_allowed_space) { |
132 | 0 | MutexLock l(&mu_); |
133 | 0 | max_allowed_space_ = max_allowed_space; |
134 | 0 | } |
135 | | |
136 | | void SstFileManagerImpl::SetCompactionBufferSize( |
137 | 0 | uint64_t compaction_buffer_size) { |
138 | 0 | MutexLock l(&mu_); |
139 | 0 | compaction_buffer_size_ = compaction_buffer_size; |
140 | 0 | } |
141 | | |
142 | 4.94k | bool SstFileManagerImpl::IsMaxAllowedSpaceReached() { |
143 | 4.94k | MutexLock l(&mu_); |
144 | 4.94k | if (max_allowed_space_ <= 0) { |
145 | 4.94k | return false; |
146 | 4.94k | } |
147 | 0 | return total_files_size_ >= max_allowed_space_; |
148 | 4.94k | } |
149 | | |
150 | 0 | bool SstFileManagerImpl::IsMaxAllowedSpaceReachedIncludingCompactions() { |
151 | 0 | MutexLock l(&mu_); |
152 | 0 | if (max_allowed_space_ <= 0) { |
153 | 0 | return false; |
154 | 0 | } |
155 | 0 | return total_files_size_ + cur_compactions_reserved_size_ >= |
156 | 0 | max_allowed_space_; |
157 | 0 | } |
158 | | |
159 | | bool SstFileManagerImpl::EnoughRoomForCompaction( |
160 | | ColumnFamilyData* cfd, const std::vector<CompactionInputFiles>& inputs, |
161 | 9.07k | const Status& bg_error) { |
162 | 9.07k | MutexLock l(&mu_); |
163 | 9.07k | uint64_t size_added_by_compaction = 0; |
164 | | // First check if we even have the space to do the compaction |
165 | 21.2k | for (size_t i = 0; i < inputs.size(); i++) { |
166 | 41.8k | for (size_t j = 0; j < inputs[i].size(); j++) { |
167 | 29.6k | FileMetaData* filemeta = inputs[i][j]; |
168 | 29.6k | size_added_by_compaction += filemeta->fd.GetFileSize(); |
169 | 29.6k | } |
170 | 12.1k | } |
171 | | |
172 | | // Update cur_compactions_reserved_size_ so concurrent compaction |
173 | | // don't max out space |
174 | 9.07k | size_t needed_headroom = cur_compactions_reserved_size_ + |
175 | 9.07k | size_added_by_compaction + compaction_buffer_size_; |
176 | 9.07k | if (max_allowed_space_ != 0 && |
177 | 0 | (needed_headroom + total_files_size_ > max_allowed_space_)) { |
178 | 0 | return false; |
179 | 0 | } |
180 | | |
181 | | // Implement more aggressive checks only if this DB instance has already |
182 | | // seen a NoSpace() error. This is tin order to contain a single potentially |
183 | | // misbehaving DB instance and prevent it from slowing down compactions of |
184 | | // other DB instances |
185 | 9.07k | if (bg_error.IsNoSpace() && CheckFreeSpace()) { |
186 | 0 | auto fn = |
187 | 0 | TableFileName(cfd->ioptions().cf_paths, inputs[0][0]->fd.GetNumber(), |
188 | 0 | inputs[0][0]->fd.GetPathId()); |
189 | 0 | uint64_t free_space = 0; |
190 | 0 | Status s = fs_->GetFreeSpace(fn, IOOptions(), &free_space, nullptr); |
191 | 0 | s.PermitUncheckedError(); // TODO: Check the status |
192 | | // needed_headroom is based on current size reserved by compactions, |
193 | | // minus any files created by running compactions as they would count |
194 | | // against the reserved size. If user didn't specify any compaction |
195 | | // buffer, add reserved_disk_buffer_ that's calculated by default so the |
196 | | // compaction doesn't end up leaving nothing for logs and flush SSTs |
197 | 0 | if (compaction_buffer_size_ == 0) { |
198 | 0 | needed_headroom += reserved_disk_buffer_; |
199 | 0 | } |
200 | 0 | if (free_space < needed_headroom + size_added_by_compaction) { |
201 | | // We hit the condition of not enough disk space |
202 | 0 | ROCKS_LOG_ERROR(logger_, |
203 | 0 | "free space [%" PRIu64 |
204 | 0 | " bytes] is less than " |
205 | 0 | "needed headroom [%" ROCKSDB_PRIszt " bytes]\n", |
206 | 0 | free_space, needed_headroom); |
207 | 0 | return false; |
208 | 0 | } |
209 | 0 | } |
210 | | |
211 | 9.07k | cur_compactions_reserved_size_ += size_added_by_compaction; |
212 | | // Take a snapshot of cur_compactions_reserved_size_ for when we encounter |
213 | | // a NoSpace error. |
214 | 9.07k | free_space_trigger_ = cur_compactions_reserved_size_; |
215 | 9.07k | return true; |
216 | 9.07k | } |
217 | | |
218 | 0 | uint64_t SstFileManagerImpl::GetCompactionsReservedSize() { |
219 | 0 | MutexLock l(&mu_); |
220 | 0 | return cur_compactions_reserved_size_; |
221 | 0 | } |
222 | | |
223 | 57.0k | uint64_t SstFileManagerImpl::GetTotalSize() { |
224 | 57.0k | MutexLock l(&mu_); |
225 | 57.0k | return total_files_size_; |
226 | 57.0k | } |
227 | | |
228 | | std::unordered_map<std::string, uint64_t> |
229 | 0 | SstFileManagerImpl::GetTrackedFiles() { |
230 | 0 | MutexLock l(&mu_); |
231 | 0 | return tracked_files_; |
232 | 0 | } |
233 | | |
234 | 53.3k | int64_t SstFileManagerImpl::GetDeleteRateBytesPerSecond() { |
235 | 53.3k | return delete_scheduler_.GetRateBytesPerSecond(); |
236 | 53.3k | } |
237 | | |
238 | 0 | void SstFileManagerImpl::SetDeleteRateBytesPerSecond(int64_t delete_rate) { |
239 | 0 | return delete_scheduler_.SetRateBytesPerSecond(delete_rate); |
240 | 0 | } |
241 | | |
242 | 0 | double SstFileManagerImpl::GetMaxTrashDBRatio() { |
243 | 0 | return delete_scheduler_.GetMaxTrashDBRatio(); |
244 | 0 | } |
245 | | |
246 | 0 | void SstFileManagerImpl::SetMaxTrashDBRatio(double r) { |
247 | 0 | return delete_scheduler_.SetMaxTrashDBRatio(r); |
248 | 0 | } |
249 | | |
250 | 0 | uint64_t SstFileManagerImpl::GetTotalTrashSize() { |
251 | 0 | return delete_scheduler_.GetTotalTrashSize(); |
252 | 0 | } |
253 | | |
254 | | void SstFileManagerImpl::ReserveDiskBuffer(uint64_t size, |
255 | 53.3k | const std::string& path) { |
256 | 53.3k | MutexLock l(&mu_); |
257 | | |
258 | 53.3k | reserved_disk_buffer_ += size; |
259 | 53.3k | if (path_.empty()) { |
260 | 53.3k | path_ = path; |
261 | 53.3k | } |
262 | 53.3k | } |
263 | | |
264 | 0 | void SstFileManagerImpl::ClearError() { |
265 | 0 | while (true) { |
266 | 0 | MutexLock l(&mu_); |
267 | |
|
268 | 0 | if (error_handler_list_.empty() || closing_) { |
269 | 0 | return; |
270 | 0 | } |
271 | | |
272 | 0 | uint64_t free_space = 0; |
273 | 0 | Status s = fs_->GetFreeSpace(path_, IOOptions(), &free_space, nullptr); |
274 | 0 | free_space = max_allowed_space_ > 0 |
275 | 0 | ? std::min(max_allowed_space_, free_space) |
276 | 0 | : free_space; |
277 | 0 | if (s.ok()) { |
278 | | // In case of multi-DB instances, some of them may have experienced a |
279 | | // soft error and some a hard error. In the SstFileManagerImpl, a hard |
280 | | // error will basically override previously reported soft errors. Once |
281 | | // we clear the hard error, we don't keep track of previous errors for |
282 | | // now |
283 | 0 | if (bg_err_.severity() == Status::Severity::kHardError) { |
284 | 0 | if (free_space < reserved_disk_buffer_) { |
285 | 0 | ROCKS_LOG_ERROR(logger_, |
286 | 0 | "free space [%" PRIu64 |
287 | 0 | " bytes] is less than " |
288 | 0 | "required disk buffer [%" PRIu64 " bytes]\n", |
289 | 0 | free_space, reserved_disk_buffer_); |
290 | 0 | ROCKS_LOG_ERROR(logger_, "Cannot clear hard error\n"); |
291 | 0 | s = Status::NoSpace(); |
292 | 0 | } |
293 | 0 | } else if (bg_err_.severity() == Status::Severity::kSoftError) { |
294 | 0 | if (free_space < free_space_trigger_) { |
295 | 0 | ROCKS_LOG_WARN(logger_, |
296 | 0 | "free space [%" PRIu64 |
297 | 0 | " bytes] is less than " |
298 | 0 | "free space for compaction trigger [%" PRIu64 |
299 | 0 | " bytes]\n", |
300 | 0 | free_space, free_space_trigger_); |
301 | 0 | ROCKS_LOG_WARN(logger_, "Cannot clear soft error\n"); |
302 | 0 | s = Status::NoSpace(); |
303 | 0 | } |
304 | 0 | } |
305 | 0 | } |
306 | | |
307 | | // Someone could have called CancelErrorRecovery() and the list could have |
308 | | // become empty, so check again here |
309 | 0 | if (s.ok()) { |
310 | 0 | assert(!error_handler_list_.empty()); |
311 | 0 | auto error_handler = error_handler_list_.front(); |
312 | | // Since we will release the mutex, set cur_instance_ to signal to the |
313 | | // shutdown thread, if it calls // CancelErrorRecovery() the meantime, |
314 | | // to indicate that this DB instance is busy. The DB instance is |
315 | | // guaranteed to not be deleted before RecoverFromBGError() returns, |
316 | | // since the ErrorHandler::recovery_in_prog_ flag would be true |
317 | 0 | cur_instance_ = error_handler; |
318 | 0 | mu_.Unlock(); |
319 | 0 | s = error_handler->RecoverFromBGError(); |
320 | 0 | TEST_SYNC_POINT("SstFileManagerImpl::ErrorCleared"); |
321 | 0 | mu_.Lock(); |
322 | | // The DB instance might have been deleted while we were |
323 | | // waiting for the mutex, so check cur_instance_ to make sure its |
324 | | // still non-null |
325 | 0 | if (cur_instance_) { |
326 | | // Check for error again, since the instance may have recovered but |
327 | | // immediately got another error. If that's the case, and the new |
328 | | // error is also a NoSpace() non-fatal error, leave the instance in |
329 | | // the list |
330 | 0 | Status err = cur_instance_->GetBGError(); |
331 | 0 | if (s.ok() && err.subcode() == IOStatus::SubCode::kNoSpace && |
332 | 0 | err.severity() < Status::Severity::kFatalError) { |
333 | 0 | s = err; |
334 | 0 | } |
335 | 0 | cur_instance_ = nullptr; |
336 | 0 | } |
337 | |
|
338 | 0 | if (s.ok() || s.IsShutdownInProgress() || |
339 | 0 | (!s.ok() && s.severity() >= Status::Severity::kFatalError)) { |
340 | | // If shutdown is in progress, abandon this handler instance |
341 | | // and continue with the others |
342 | 0 | error_handler_list_.pop_front(); |
343 | 0 | } |
344 | 0 | } |
345 | |
|
346 | 0 | if (!error_handler_list_.empty()) { |
347 | | // If there are more instances to be recovered, reschedule after 5 |
348 | | // seconds |
349 | 0 | int64_t wait_until = clock_->NowMicros() + 5000000; |
350 | 0 | cv_.TimedWait(wait_until); |
351 | 0 | } |
352 | | |
353 | | // Check again for error_handler_list_ empty, as a DB instance shutdown |
354 | | // could have removed it from the queue while we were in timed wait |
355 | 0 | if (error_handler_list_.empty()) { |
356 | 0 | ROCKS_LOG_INFO(logger_, "Clearing error\n"); |
357 | 0 | bg_err_ = Status::OK(); |
358 | 0 | return; |
359 | 0 | } |
360 | 0 | } |
361 | 0 | } |
362 | | |
363 | | void SstFileManagerImpl::StartErrorRecovery(ErrorHandler* handler, |
364 | 0 | Status bg_error) { |
365 | 0 | MutexLock l(&mu_); |
366 | 0 | if (bg_error.severity() == Status::Severity::kSoftError) { |
367 | 0 | if (bg_err_.ok()) { |
368 | | // Setting bg_err_ basically means we're in degraded mode |
369 | | // Assume that all pending compactions will fail similarly. The trigger |
370 | | // for clearing this condition is set to current compaction reserved |
371 | | // size, so we stop checking disk space available in |
372 | | // EnoughRoomForCompaction once this much free space is available |
373 | 0 | bg_err_ = bg_error; |
374 | 0 | } |
375 | 0 | } else if (bg_error.severity() == Status::Severity::kHardError) { |
376 | 0 | bg_err_ = bg_error; |
377 | 0 | } else { |
378 | 0 | assert(false); |
379 | 0 | } |
380 | | |
381 | | // If this is the first instance of this error, kick of a thread to poll |
382 | | // and recover from this condition |
383 | 0 | if (error_handler_list_.empty()) { |
384 | 0 | error_handler_list_.push_back(handler); |
385 | | // Release lock before calling join. Its ok to do so because |
386 | | // error_handler_list_ is now non-empty, so no other invocation of this |
387 | | // function will execute this piece of code |
388 | 0 | mu_.Unlock(); |
389 | 0 | if (bg_thread_) { |
390 | 0 | bg_thread_->join(); |
391 | 0 | } |
392 | | // Start a new thread. The previous one would have exited. |
393 | 0 | bg_thread_.reset(new port::Thread(&SstFileManagerImpl::ClearError, this)); |
394 | 0 | mu_.Lock(); |
395 | 0 | } else { |
396 | | // Check if this DB instance is already in the list |
397 | 0 | for (auto iter = error_handler_list_.begin(); |
398 | 0 | iter != error_handler_list_.end(); ++iter) { |
399 | 0 | if ((*iter) == handler) { |
400 | 0 | return; |
401 | 0 | } |
402 | 0 | } |
403 | 0 | error_handler_list_.push_back(handler); |
404 | 0 | } |
405 | 0 | } |
406 | | |
407 | 53.3k | bool SstFileManagerImpl::CancelErrorRecovery(ErrorHandler* handler) { |
408 | 53.3k | MutexLock l(&mu_); |
409 | | |
410 | 53.3k | if (cur_instance_ == handler) { |
411 | | // This instance is currently busy attempting to recover |
412 | | // Nullify it so the recovery thread doesn't attempt to access it again |
413 | 0 | cur_instance_ = nullptr; |
414 | 0 | return false; |
415 | 0 | } |
416 | | |
417 | 53.3k | for (auto iter = error_handler_list_.begin(); |
418 | 53.3k | iter != error_handler_list_.end(); ++iter) { |
419 | 0 | if ((*iter) == handler) { |
420 | 0 | error_handler_list_.erase(iter); |
421 | 0 | return true; |
422 | 0 | } |
423 | 0 | } |
424 | 53.3k | return false; |
425 | 53.3k | } |
426 | | |
427 | | Status SstFileManagerImpl::ScheduleFileDeletion(const std::string& file_path, |
428 | | const std::string& path_to_sync, |
429 | 57.0k | const bool force_bg) { |
430 | 57.0k | TEST_SYNC_POINT_CALLBACK("SstFileManagerImpl::ScheduleFileDeletion", |
431 | 57.0k | const_cast<std::string*>(&file_path)); |
432 | 57.0k | return delete_scheduler_.DeleteFile(file_path, path_to_sync, force_bg); |
433 | 57.0k | } |
434 | | |
435 | | Status SstFileManagerImpl::ScheduleUnaccountedFileDeletion( |
436 | | const std::string& file_path, const std::string& dir_to_sync, |
437 | 22.8k | const bool force_bg, std::optional<int32_t> bucket) { |
438 | 22.8k | TEST_SYNC_POINT_CALLBACK( |
439 | 22.8k | "SstFileManagerImpl::ScheduleUnaccountedFileDeletion", |
440 | 22.8k | const_cast<std::string*>(&file_path)); |
441 | 22.8k | return delete_scheduler_.DeleteUnaccountedFile(file_path, dir_to_sync, |
442 | 22.8k | force_bg, bucket); |
443 | 22.8k | } |
444 | | |
445 | 0 | void SstFileManagerImpl::WaitForEmptyTrash() { |
446 | 0 | delete_scheduler_.WaitForEmptyTrash(); |
447 | 0 | } |
448 | | |
449 | 0 | std::optional<int32_t> SstFileManagerImpl::NewTrashBucket() { |
450 | 0 | return delete_scheduler_.NewTrashBucket(); |
451 | 0 | } |
452 | | |
453 | 0 | void SstFileManagerImpl::WaitForEmptyTrashBucket(int32_t bucket) { |
454 | 0 | delete_scheduler_.WaitForEmptyTrashBucket(bucket); |
455 | 0 | } |
456 | | |
457 | | void SstFileManagerImpl::OnAddFileImpl(const std::string& file_path, |
458 | 124k | uint64_t file_size) { |
459 | 124k | auto tracked_file = tracked_files_.find(file_path); |
460 | 124k | if (tracked_file != tracked_files_.end()) { |
461 | | // File was added before, we will just update the size |
462 | 0 | total_files_size_ -= tracked_file->second; |
463 | 0 | total_files_size_ += file_size; |
464 | 124k | } else { |
465 | 124k | total_files_size_ += file_size; |
466 | 124k | } |
467 | 124k | tracked_files_[file_path] = file_size; |
468 | 124k | } |
469 | | |
470 | 57.0k | void SstFileManagerImpl::OnDeleteFileImpl(const std::string& file_path) { |
471 | 57.0k | auto tracked_file = tracked_files_.find(file_path); |
472 | 57.0k | if (tracked_file == tracked_files_.end()) { |
473 | | // File is not tracked |
474 | 47.2k | return; |
475 | 47.2k | } |
476 | | |
477 | 9.73k | total_files_size_ -= tracked_file->second; |
478 | 9.73k | tracked_files_.erase(tracked_file); |
479 | 9.73k | } |
480 | | |
481 | | SstFileManager* NewSstFileManager(Env* env, std::shared_ptr<Logger> info_log, |
482 | | std::string trash_dir, |
483 | | int64_t rate_bytes_per_sec, |
484 | | bool delete_existing_trash, Status* status, |
485 | | double max_trash_db_ratio, |
486 | 61.7k | uint64_t bytes_max_delete_chunk) { |
487 | 61.7k | const auto& fs = env->GetFileSystem(); |
488 | 61.7k | return NewSstFileManager(env, fs, info_log, trash_dir, rate_bytes_per_sec, |
489 | 61.7k | delete_existing_trash, status, max_trash_db_ratio, |
490 | 61.7k | bytes_max_delete_chunk); |
491 | 61.7k | } |
492 | | |
493 | | SstFileManager* NewSstFileManager(Env* env, std::shared_ptr<FileSystem> fs, |
494 | | std::shared_ptr<Logger> info_log, |
495 | | const std::string& trash_dir, |
496 | | int64_t rate_bytes_per_sec, |
497 | | bool delete_existing_trash, Status* status, |
498 | | double max_trash_db_ratio, |
499 | 61.7k | uint64_t bytes_max_delete_chunk) { |
500 | 61.7k | const auto& clock = env->GetSystemClock(); |
501 | 61.7k | SstFileManagerImpl* res = |
502 | 61.7k | new SstFileManagerImpl(clock, fs, info_log, rate_bytes_per_sec, |
503 | 61.7k | max_trash_db_ratio, bytes_max_delete_chunk); |
504 | | |
505 | | // trash_dir is deprecated and not needed anymore, but if user passed it |
506 | | // we will still remove files in it. |
507 | 61.7k | Status s = Status::OK(); |
508 | 61.7k | if (delete_existing_trash && trash_dir != "") { |
509 | 0 | std::vector<std::string> files_in_trash; |
510 | 0 | s = fs->GetChildren(trash_dir, IOOptions(), &files_in_trash, nullptr); |
511 | 0 | if (s.ok()) { |
512 | 0 | for (const std::string& trash_file : files_in_trash) { |
513 | 0 | std::string path_in_trash = trash_dir + "/" + trash_file; |
514 | 0 | res->OnAddFile(path_in_trash); |
515 | 0 | Status file_delete = |
516 | 0 | res->ScheduleFileDeletion(path_in_trash, trash_dir); |
517 | 0 | if (s.ok() && !file_delete.ok()) { |
518 | 0 | s = file_delete; |
519 | 0 | } |
520 | 0 | } |
521 | 0 | } |
522 | 0 | } |
523 | | |
524 | 61.7k | if (status) { |
525 | 0 | *status = s; |
526 | 61.7k | } else { |
527 | | // No one passed us a Status, so they must not care about the error... |
528 | 61.7k | s.PermitUncheckedError(); |
529 | 61.7k | } |
530 | | |
531 | 61.7k | return res; |
532 | 61.7k | } |
533 | | |
534 | | } // namespace ROCKSDB_NAMESPACE |