Coverage Report

Created: 2026-08-13 12:06

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
common/cpp/obj-client/obj_storage_client.cpp
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
#include "obj_storage_client.h"
19
20
#include <cpp/sync_point.h>
21
#include <glog/logging.h>
22
23
#include <algorithm>
24
#include <chrono>
25
26
namespace doris {
27
333
std::unique_ptr<ObjStorageListIterator> ObjStorageClient::list_objects(const ObjStoragePath& opts) {
28
333
    return std::make_unique<ObjStorageListIterator>(shared_from_this(), opts);
29
333
}
30
31
ObjStorageResponse ObjStorageClient::list_objects(const ObjStoragePath& opts,
32
332
                                                  std::vector<ObjectMeta>* objects) {
33
332
    objects->clear();
34
332
    auto iter = list_objects(opts);
35
574
    for (;;) {
36
574
        auto result = iter->next();
37
574
        if (!result.object.has_value()) {
38
332
            if (!result.resp.ok()) {
39
1
                objects->clear();
40
1
            }
41
332
            return result.resp;
42
332
        }
43
242
        objects->emplace_back(std::move(*result.object));
44
242
    }
45
332
}
46
47
577
ObjStorageResponse ObjStorageListIterator::has_next() {
48
577
    if (!is_valid_) {
49
0
        return {
50
0
                .status = {TStatusCode::INTERNAL_ERROR, "Iterator is invalid"},
51
0
                .http_code = 0,
52
0
        };
53
0
    }
54
913
    while (next_index_ == objects_.size()) {
55
668
        if (!has_more_) {
56
331
            return {
57
331
                    .status = {TStatusCode::END_OF_FILE, "No more results"},
58
331
                    .http_code = 200,
59
331
            };
60
331
        }
61
337
        auto page = client_->list_objects_page(opts_, continuation_token_);
62
337
        if (!page.resp.ok()) {
63
1
            is_valid_ = false;
64
1
            return page.resp;
65
1
        }
66
336
        objects_ = std::move(page.objects);
67
336
        next_index_ = 0;
68
336
        continuation_token_ = std::move(page.continuation_token);
69
336
        has_more_ = page.has_more;
70
336
    }
71
245
    return ObjStorageResponse::OK();
72
577
}
73
74
577
ObjStorageListResult ObjStorageListIterator::next() {
75
577
    auto response = has_next();
76
577
    if (response.status.code == ObjStorageStatus::END_OF_FILE) {
77
331
        return {.resp = ObjStorageResponse::OK(), .object = {}};
78
331
    }
79
246
    if (!response.ok()) {
80
1
        return {.resp = std::move(response), .object = {}};
81
1
    }
82
245
    return {
83
245
            .resp = ObjStorageResponse::OK(),
84
245
            .object = std::move(objects_[next_index_++]),
85
245
    };
86
246
}
87
88
ObjStorageResponse ObjStorageClient::delete_objects_recursively(
89
0
        const ObjStoragePath& path, const ObjStorageRecursiveDeleteOptions& delete_options) {
90
0
    return delete_objects_recursively_impl(path, delete_options);
91
0
}
92
93
ObjStorageResponse ObjStorageClient::delete_objects_recursively_impl(
94
0
        const ObjStoragePath& path, const ObjStorageRecursiveDeleteOptions& delete_options) {
95
0
    const auto start_time = std::chrono::steady_clock::now();
96
0
    auto list_path = path;
97
0
    if (list_path.prefix.empty()) {
98
0
        list_path.prefix = list_path.key;
99
0
    }
100
0
    auto delete_batch_size = std::max<size_t>(1, capabilities().max_delete_batch);
101
0
    TEST_SYNC_POINT_CALLBACK("ObjStorageClient::delete_objects_recursively_", &delete_batch_size);
102
0
    delete_batch_size = std::max<size_t>(1, delete_batch_size);
103
0
    const auto max_tasks_per_batch = std::max<size_t>(1, delete_options.max_tasks_per_batch);
104
0
    std::vector<std::string> keys;
105
0
    keys.reserve(delete_batch_size);
106
0
    size_t pending_tasks = 0;
107
0
    size_t total_batches = 0;
108
0
    size_t num_deleted = 0;
109
0
    size_t error_count = 0;
110
0
    auto first_error = ObjStorageResponse::OK();
111
112
0
    auto elapsed_milliseconds = [&]() {
113
0
        return std::chrono::duration_cast<std::chrono::milliseconds>(
114
0
                       std::chrono::steady_clock::now() - start_time)
115
0
                .count();
116
0
    };
117
0
    auto finish = [&](ObjStorageResponse response) {
118
0
        LOG(INFO) << "delete objects under " << list_path.bucket << "/" << list_path.prefix
119
0
                  << " finished, ret=" << response.status.code
120
0
                  << ", total_batches=" << total_batches << ", num_deleted=" << num_deleted
121
0
                  << ", error_count=" << error_count << ", cost=" << elapsed_milliseconds()
122
0
                  << " ms";
123
0
        return response;
124
0
    };
125
0
    auto record_error = [&](ObjStorageResponse response) {
126
0
        if (response.ok()) {
127
0
            return;
128
0
        }
129
0
        ++error_count;
130
0
        if (first_error.ok()) {
131
0
            first_error = std::move(response);
132
0
        }
133
0
    };
134
135
0
    auto wait_for_tasks = [&]() {
136
0
        if (pending_tasks == 0) {
137
0
            return ObjStorageResponse::OK();
138
0
        }
139
0
        const auto tasks_in_batch = pending_tasks;
140
0
        pending_tasks = 0;
141
0
        auto response = delete_options.executor ? delete_options.executor->wait()
142
0
                                                : ObjStorageResponse::OK();
143
0
        ++total_batches;
144
0
        LOG(INFO) << "delete objects under " << list_path.bucket << "/" << list_path.prefix
145
0
                  << " batch " << total_batches << " completed"
146
0
                  << ", tasks_in_batch=" << tasks_in_batch << ", total_deleted=" << num_deleted
147
0
                  << ", elapsed=" << elapsed_milliseconds() << " ms";
148
0
        return response;
149
0
    };
150
0
    auto submit_delete_task = [&]() -> bool {
151
0
        ObjStorageDeleteTask task = [this, bucket = path.bucket,
152
0
                                     batch = std::move(keys)]() mutable {
153
0
            return delete_objects(ObjStoragePath {.bucket = std::move(bucket)}, std::move(batch));
154
0
        };
155
0
        keys.clear();
156
0
        keys.reserve(delete_batch_size);
157
158
0
        ObjStorageResponse response;
159
0
        if (delete_options.executor) {
160
0
            response = delete_options.executor->submit(std::move(task));
161
0
        } else {
162
0
            response = task();
163
0
        }
164
0
        ++pending_tasks;
165
0
        if (!response.ok()) {
166
0
            record_error(std::move(response));
167
0
            record_error(wait_for_tasks());
168
0
            return false;
169
0
        }
170
0
        if (pending_tasks == max_tasks_per_batch) {
171
0
            record_error(wait_for_tasks());
172
0
        }
173
        // Match the pre-refactor Recycler behavior: do not scan the next task batch after the
174
        // current batch reports a submit, delete, or wait failure.
175
0
        return first_error.ok();
176
0
    };
177
178
0
    auto iter = list_objects(list_path);
179
0
    for (;;) {
180
0
        auto result = iter->next();
181
0
        if (!result.object.has_value()) {
182
0
            if (result.resp.ok()) {
183
0
                break;
184
0
            }
185
0
            if (!keys.empty()) {
186
0
                submit_delete_task();
187
0
            }
188
0
            record_error(wait_for_tasks());
189
0
            record_error(std::move(result.resp));
190
0
            return finish(std::move(first_error));
191
0
        }
192
0
        auto& object = *result.object;
193
0
        if (delete_options.expiration_time > 0 && object.mtime_s > delete_options.expiration_time) {
194
0
            continue;
195
0
        }
196
0
        ++num_deleted;
197
0
        keys.emplace_back(std::move(object.key));
198
0
        if (keys.size() == delete_batch_size && !submit_delete_task()) {
199
0
            return finish(std::move(first_error));
200
0
        }
201
0
    }
202
0
    if (!keys.empty() && !submit_delete_task()) {
203
0
        return finish(std::move(first_error));
204
0
    }
205
0
    record_error(wait_for_tasks());
206
0
    return finish(std::move(first_error));
207
0
}
208
209
} // namespace doris