/src/rocksdb/file/file_util.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/file_util.h" |
7 | | |
8 | | #include <algorithm> |
9 | | #include <string> |
10 | | |
11 | | #include "file/random_access_file_reader.h" |
12 | | #include "file/sequence_file_reader.h" |
13 | | #include "file/sst_file_manager_impl.h" |
14 | | #include "file/writable_file_writer.h" |
15 | | #include "rocksdb/env.h" |
16 | | #include "rocksdb/statistics.h" |
17 | | |
18 | | namespace ROCKSDB_NAMESPACE { |
19 | | |
20 | | // Utility function to copy a file up to a specified length |
21 | | IOStatus CopyFile(FileSystem* fs, const std::string& source, |
22 | | Temperature src_temp_hint, |
23 | | std::unique_ptr<WritableFileWriter>& dest_writer, |
24 | | uint64_t size, bool use_fsync, |
25 | | const std::shared_ptr<IOTracer>& io_tracer, |
26 | | uint64_t max_read_buffer_size, |
27 | | const std::optional<IOOptions>& readIOOptions, |
28 | 0 | const std::optional<IOOptions>& writeIOOptions) { |
29 | 0 | FileOptions soptions; |
30 | 0 | IOStatus io_s; |
31 | 0 | std::unique_ptr<SequentialFileReader> src_reader; |
32 | 0 | const IOOptions opts; |
33 | |
|
34 | 0 | { |
35 | 0 | soptions.temperature = src_temp_hint; |
36 | 0 | std::unique_ptr<FSSequentialFile> srcfile; |
37 | 0 | io_s = fs->NewSequentialFile(source, soptions, &srcfile, nullptr); |
38 | 0 | if (!io_s.ok()) { |
39 | 0 | return io_s; |
40 | 0 | } |
41 | | |
42 | 0 | if (size == 0) { |
43 | | // default argument means copy everything |
44 | 0 | io_s = |
45 | 0 | fs->GetFileSize(source, readIOOptions.value_or(opts), &size, nullptr); |
46 | 0 | if (!io_s.ok()) { |
47 | 0 | return io_s; |
48 | 0 | } |
49 | 0 | } |
50 | 0 | src_reader.reset( |
51 | 0 | new SequentialFileReader(std::move(srcfile), source, io_tracer)); |
52 | 0 | } |
53 | | |
54 | 0 | const size_t read_buffer_size = std::max( |
55 | 0 | static_cast<size_t>(4096), static_cast<size_t>(max_read_buffer_size)); |
56 | 0 | std::unique_ptr<char[]> buffer; |
57 | 0 | buffer.reset(new char[read_buffer_size]); |
58 | |
|
59 | 0 | Env::IOPriority read_rate_limiter_priority = Env::IO_TOTAL; |
60 | 0 | if (readIOOptions.has_value()) { |
61 | 0 | read_rate_limiter_priority = readIOOptions.value().rate_limiter_priority; |
62 | 0 | } |
63 | 0 | Slice slice; |
64 | 0 | while (size > 0) { |
65 | 0 | size_t bytes_to_read = std::min(static_cast<size_t>(read_buffer_size), |
66 | 0 | static_cast<size_t>(size)); |
67 | | // TODO: rate limit copy file |
68 | 0 | io_s = status_to_io_status(src_reader->Read( |
69 | 0 | bytes_to_read, &slice, buffer.get(), |
70 | 0 | read_rate_limiter_priority /* rate_limiter_priority */)); |
71 | 0 | if (!io_s.ok()) { |
72 | 0 | return io_s; |
73 | 0 | } |
74 | 0 | if (slice.size() == 0) { |
75 | 0 | return IOStatus::Corruption( |
76 | 0 | "File smaller than expected for copy: " + source + " expecting " + |
77 | 0 | std::to_string(size) + " more bytes after " + |
78 | 0 | std::to_string(dest_writer->GetFileSize())); |
79 | 0 | } |
80 | | |
81 | 0 | io_s = dest_writer->Append(slice, writeIOOptions.value_or(opts)); |
82 | 0 | if (!io_s.ok()) { |
83 | 0 | return io_s; |
84 | 0 | } |
85 | 0 | size -= slice.size(); |
86 | 0 | } |
87 | 0 | return use_fsync ? dest_writer->Fsync(writeIOOptions.value_or(opts)) |
88 | 0 | : dest_writer->Sync(writeIOOptions.value_or(opts)); |
89 | 0 | } |
90 | | |
91 | | IOStatus CopyFile(FileSystem* fs, const std::string& source, |
92 | | Temperature src_temp_hint, const std::string& destination, |
93 | | Temperature dst_temp, uint64_t size, bool use_fsync, |
94 | | const std::shared_ptr<IOTracer>& io_tracer, |
95 | | uint64_t max_read_buffer_size, |
96 | | const std::optional<IOOptions>& readIOOptions, |
97 | 0 | const std::optional<IOOptions>& writeIOOptions) { |
98 | 0 | FileOptions options; |
99 | 0 | IOStatus io_s; |
100 | 0 | std::unique_ptr<WritableFileWriter> dest_writer; |
101 | |
|
102 | 0 | { |
103 | 0 | options.temperature = dst_temp; |
104 | 0 | std::unique_ptr<FSWritableFile> destfile; |
105 | 0 | io_s = fs->NewWritableFile(destination, options, &destfile, nullptr); |
106 | 0 | if (!io_s.ok()) { |
107 | 0 | return io_s; |
108 | 0 | } |
109 | | |
110 | | // TODO: pass in Histograms if the destination file is sst or blob |
111 | 0 | dest_writer.reset( |
112 | 0 | new WritableFileWriter(std::move(destfile), destination, options)); |
113 | 0 | } |
114 | | |
115 | 0 | return CopyFile(fs, source, src_temp_hint, dest_writer, size, use_fsync, |
116 | 0 | io_tracer, max_read_buffer_size, readIOOptions, |
117 | 0 | writeIOOptions); |
118 | 0 | } |
119 | | |
120 | | // Utility function to create a file with the provided contents |
121 | | IOStatus CreateFile(FileSystem* fs, const std::string& destination, |
122 | 0 | const std::string& contents, bool use_fsync) { |
123 | 0 | const EnvOptions soptions; |
124 | 0 | IOStatus io_s; |
125 | 0 | std::unique_ptr<WritableFileWriter> dest_writer; |
126 | 0 | const IOOptions opts; |
127 | |
|
128 | 0 | std::unique_ptr<FSWritableFile> destfile; |
129 | 0 | io_s = fs->NewWritableFile(destination, soptions, &destfile, nullptr); |
130 | 0 | if (!io_s.ok()) { |
131 | 0 | return io_s; |
132 | 0 | } |
133 | | // TODO: pass in Histograms if the destination file is sst or blob |
134 | 0 | dest_writer.reset( |
135 | 0 | new WritableFileWriter(std::move(destfile), destination, soptions)); |
136 | 0 | io_s = dest_writer->Append(Slice(contents), opts); |
137 | 0 | if (!io_s.ok()) { |
138 | 0 | return io_s; |
139 | 0 | } |
140 | 0 | return use_fsync ? dest_writer->Fsync(opts) : dest_writer->Sync(opts); |
141 | 0 | } |
142 | | |
143 | | Status DeleteDBFile(const ImmutableDBOptions* db_options, |
144 | | const std::string& fname, const std::string& dir_to_sync, |
145 | 57.0k | const bool force_bg, const bool force_fg) { |
146 | 57.0k | SstFileManagerImpl* sfm = static_cast_with_check<SstFileManagerImpl>( |
147 | 57.0k | db_options->sst_file_manager.get()); |
148 | 57.0k | if (sfm && !force_fg) { |
149 | 57.0k | return sfm->ScheduleFileDeletion(fname, dir_to_sync, force_bg); |
150 | 57.0k | } else { |
151 | 0 | return db_options->env->DeleteFile(fname); |
152 | 0 | } |
153 | 57.0k | } |
154 | | |
155 | | Status DeleteUnaccountedDBFile(const ImmutableDBOptions* db_options, |
156 | | const std::string& fname, |
157 | | const std::string& dir_to_sync, |
158 | | const bool force_bg, const bool force_fg, |
159 | 22.8k | std::optional<int32_t> bucket) { |
160 | 22.8k | SstFileManagerImpl* sfm = static_cast_with_check<SstFileManagerImpl>( |
161 | 22.8k | db_options->sst_file_manager.get()); |
162 | 22.8k | if (sfm && !force_fg) { |
163 | 22.8k | return sfm->ScheduleUnaccountedFileDeletion(fname, dir_to_sync, force_bg, |
164 | 22.8k | bucket); |
165 | 22.8k | } else { |
166 | 0 | return db_options->env->DeleteFile(fname); |
167 | 0 | } |
168 | 22.8k | } |
169 | | |
170 | | // requested_checksum_func_name brings the function name of the checksum |
171 | | // generator in checksum_factory. Empty string is permitted, in which case the |
172 | | // name of the generator created by the factory is unchecked. When |
173 | | // `requested_checksum_func_name` is non-empty, however, the created generator's |
174 | | // name must match it, otherwise an `InvalidArgument` error is returned. |
175 | | IOStatus GenerateOneFileChecksum( |
176 | | FileSystem* fs, const std::string& file_path, |
177 | | FileChecksumGenFactory* checksum_factory, |
178 | | const std::string& requested_checksum_func_name, std::string* file_checksum, |
179 | | std::string* file_checksum_func_name, |
180 | | size_t verify_checksums_readahead_size, bool /*allow_mmap_reads*/, |
181 | | std::shared_ptr<IOTracer>& io_tracer, RateLimiter* rate_limiter, |
182 | | const ReadOptions& read_options, Statistics* stats, SystemClock* clock, |
183 | 0 | const FileOptions& file_options) { |
184 | 0 | if (checksum_factory == nullptr) { |
185 | 0 | return IOStatus::InvalidArgument("Checksum factory is invalid"); |
186 | 0 | } |
187 | 0 | assert(file_checksum != nullptr); |
188 | 0 | assert(file_checksum_func_name != nullptr); |
189 | |
|
190 | 0 | FileChecksumGenContext gen_context; |
191 | 0 | gen_context.requested_checksum_func_name = requested_checksum_func_name; |
192 | 0 | gen_context.file_name = file_path; |
193 | 0 | std::unique_ptr<FileChecksumGenerator> checksum_generator = |
194 | 0 | checksum_factory->CreateFileChecksumGenerator(gen_context); |
195 | 0 | if (checksum_generator == nullptr) { |
196 | 0 | std::string msg = |
197 | 0 | "Cannot get the file checksum generator based on the requested " |
198 | 0 | "checksum function name: " + |
199 | 0 | requested_checksum_func_name + |
200 | 0 | " from checksum factory: " + checksum_factory->Name(); |
201 | 0 | return IOStatus::InvalidArgument(msg); |
202 | 0 | } else { |
203 | | // For backward compatibility and use in file ingestion clients where there |
204 | | // is no stored checksum function name, `requested_checksum_func_name` can |
205 | | // be empty. If we give the requested checksum function name, we expect it |
206 | | // is the same name of the checksum generator. |
207 | 0 | if (!requested_checksum_func_name.empty() && |
208 | 0 | checksum_generator->Name() != requested_checksum_func_name) { |
209 | 0 | std::string msg = "Expected file checksum generator named '" + |
210 | 0 | requested_checksum_func_name + |
211 | 0 | "', while the factory created one " |
212 | 0 | "named '" + |
213 | 0 | checksum_generator->Name() + "'"; |
214 | 0 | return IOStatus::InvalidArgument(msg); |
215 | 0 | } |
216 | 0 | } |
217 | | |
218 | 0 | uint64_t size; |
219 | 0 | IOStatus io_s; |
220 | 0 | std::unique_ptr<RandomAccessFileReader> reader; |
221 | 0 | { |
222 | 0 | std::unique_ptr<FSRandomAccessFile> r_file; |
223 | 0 | FileOptions fopts = file_options; |
224 | 0 | if (fopts.file_checksum.empty()) { |
225 | | // No expected checksum is known -- this is a from-scratch computation. |
226 | 0 | fopts.file_checksum_func_name = kNoFileChecksumFuncName; |
227 | 0 | } |
228 | 0 | io_s = fs->NewRandomAccessFile(file_path, fopts, &r_file, nullptr); |
229 | 0 | if (!io_s.ok()) { |
230 | 0 | return io_s; |
231 | 0 | } |
232 | 0 | io_s = fs->GetFileSize(file_path, IOOptions(), &size, nullptr); |
233 | 0 | if (!io_s.ok()) { |
234 | 0 | return io_s; |
235 | 0 | } |
236 | 0 | reader.reset(new RandomAccessFileReader( |
237 | 0 | std::move(r_file), file_path, clock, io_tracer, stats, |
238 | 0 | Histograms::SST_READ_MICROS, nullptr, rate_limiter)); |
239 | 0 | } |
240 | | |
241 | | // Found that 256 KB readahead size provides the best performance, based on |
242 | | // experiments, for auto readahead. Experiment data is in PR #3282. |
243 | 0 | size_t default_max_read_ahead_size = 256 * 1024; |
244 | 0 | size_t readahead_size = (verify_checksums_readahead_size != 0) |
245 | 0 | ? verify_checksums_readahead_size |
246 | 0 | : default_max_read_ahead_size; |
247 | 0 | std::unique_ptr<char[]> buf; |
248 | 0 | if (reader->use_direct_io()) { |
249 | 0 | size_t alignment = reader->file()->GetRequiredBufferAlignment(); |
250 | 0 | readahead_size = (readahead_size + alignment - 1) & ~(alignment - 1); |
251 | 0 | } |
252 | 0 | buf.reset(new char[readahead_size]); |
253 | |
|
254 | 0 | Slice slice; |
255 | 0 | uint64_t offset = 0; |
256 | 0 | IOOptions opts; |
257 | 0 | IODebugContext dbg; |
258 | 0 | io_s = reader->PrepareIOOptions(read_options, opts, &dbg); |
259 | 0 | if (!io_s.ok()) { |
260 | 0 | return io_s; |
261 | 0 | } |
262 | 0 | while (size > 0) { |
263 | 0 | size_t bytes_to_read = |
264 | 0 | static_cast<size_t>(std::min(uint64_t{readahead_size}, size)); |
265 | 0 | io_s = reader->Read(opts, offset, bytes_to_read, &slice, buf.get(), nullptr, |
266 | 0 | &dbg); |
267 | 0 | if (!io_s.ok()) { |
268 | 0 | return IOStatus::Corruption("file read failed with error: " + |
269 | 0 | io_s.ToString()); |
270 | 0 | } |
271 | 0 | if (slice.size() == 0) { |
272 | 0 | return IOStatus::Corruption( |
273 | 0 | "File smaller than expected for checksum: " + file_path + |
274 | 0 | " expecting " + std::to_string(size) + " more bytes after " + |
275 | 0 | std::to_string(offset)); |
276 | 0 | } |
277 | 0 | checksum_generator->Update(slice.data(), slice.size()); |
278 | 0 | size -= slice.size(); |
279 | 0 | offset += slice.size(); |
280 | |
|
281 | 0 | TEST_SYNC_POINT("GenerateOneFileChecksum::Chunk:0"); |
282 | 0 | } |
283 | 0 | checksum_generator->Finalize(); |
284 | 0 | *file_checksum = checksum_generator->GetChecksum(); |
285 | 0 | *file_checksum_func_name = checksum_generator->Name(); |
286 | 0 | return IOStatus::OK(); |
287 | 0 | } |
288 | | |
289 | 0 | Status DestroyDir(Env* env, const std::string& dir) { |
290 | 0 | Status s; |
291 | 0 | if (env->FileExists(dir).IsNotFound()) { |
292 | 0 | return s; |
293 | 0 | } |
294 | 0 | std::vector<std::string> files_in_dir; |
295 | 0 | s = env->GetChildren(dir, &files_in_dir); |
296 | 0 | if (s.ok()) { |
297 | 0 | for (auto& file_in_dir : files_in_dir) { |
298 | 0 | std::string path = dir + "/" + file_in_dir; |
299 | 0 | bool is_dir = false; |
300 | 0 | s = env->IsDirectory(path, &is_dir); |
301 | 0 | if (s.ok()) { |
302 | 0 | if (is_dir) { |
303 | 0 | s = DestroyDir(env, path); |
304 | 0 | } else { |
305 | 0 | s = env->DeleteFile(path); |
306 | 0 | } |
307 | 0 | } else if (s.IsNotSupported()) { |
308 | 0 | s = Status::OK(); |
309 | 0 | } |
310 | 0 | if (!s.ok()) { |
311 | | // IsDirectory, etc. might not report NotFound |
312 | 0 | if (s.IsNotFound() || env->FileExists(path).IsNotFound()) { |
313 | | // Allow files to be deleted externally |
314 | 0 | s = Status::OK(); |
315 | 0 | } else { |
316 | 0 | break; |
317 | 0 | } |
318 | 0 | } |
319 | 0 | } |
320 | 0 | } |
321 | |
|
322 | 0 | if (s.ok()) { |
323 | 0 | s = env->DeleteDir(dir); |
324 | | // DeleteDir might or might not report NotFound |
325 | 0 | if (!s.ok() && (s.IsNotFound() || env->FileExists(dir).IsNotFound())) { |
326 | | // Allow to be deleted externally |
327 | 0 | s = Status::OK(); |
328 | 0 | } |
329 | 0 | } |
330 | 0 | return s; |
331 | 0 | } |
332 | | |
333 | | } // namespace ROCKSDB_NAMESPACE |