Coverage Report

Created: 2026-03-23 10:27

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/agent/task_worker_pool.h
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
#pragma once
19
20
#include <atomic>
21
#include <condition_variable>
22
#include <deque>
23
#include <functional>
24
#include <memory>
25
#include <mutex>
26
#include <string>
27
#include <string_view>
28
29
#include "common/status.h"
30
31
namespace doris {
32
33
class ExecEnv;
34
class StorageEngine;
35
class CloudStorageEngine;
36
class Thread;
37
class ThreadPool;
38
class TReportRequest;
39
class TTabletInfo;
40
class TAgentTaskRequest;
41
class ClusterInfo;
42
43
class TaskWorkerPoolIf {
44
public:
45
80
    virtual ~TaskWorkerPoolIf() = default;
46
47
    virtual Status submit_task(const TAgentTaskRequest& task) = 0;
48
};
49
50
class TaskWorkerPool : public TaskWorkerPoolIf {
51
public:
52
    TaskWorkerPool(std::string_view name, int worker_count,
53
                   std::function<void(const TAgentTaskRequest&)> callback,
54
                   std::function<void(const TAgentTaskRequest&)> pre_submit_callback = nullptr);
55
56
    ~TaskWorkerPool() override;
57
58
    void stop();
59
60
    Status submit_task(const TAgentTaskRequest& task) override;
61
62
protected:
63
    std::atomic_bool _stopped {false};
64
    std::unique_ptr<ThreadPool> _thread_pool;
65
    std::function<void(const TAgentTaskRequest&)> _callback;
66
    std::function<void(const TAgentTaskRequest&)> _pre_submit_callback;
67
};
68
69
class PublishVersionWorkerPool final : public TaskWorkerPool {
70
public:
71
    PublishVersionWorkerPool(StorageEngine& engine);
72
73
    ~PublishVersionWorkerPool() override;
74
75
private:
76
    void publish_version_callback(const TAgentTaskRequest& task);
77
78
    StorageEngine& _engine;
79
};
80
81
class PriorTaskWorkerPool final : public TaskWorkerPoolIf {
82
public:
83
    PriorTaskWorkerPool(const std::string& name, int normal_worker_count,
84
                        int high_prior_worker_count,
85
                        std::function<void(const TAgentTaskRequest& task)> callback);
86
87
    ~PriorTaskWorkerPool() override;
88
89
    void stop();
90
91
    Status submit_task(const TAgentTaskRequest& task) override;
92
93
    Status submit_high_prior_and_cancel_low(TAgentTaskRequest& task);
94
95
private:
96
    void normal_loop();
97
98
    void high_prior_loop();
99
100
    bool _stopped {false};
101
102
    std::mutex _mtx;
103
    std::condition_variable _normal_condv;
104
    std::deque<std::unique_ptr<TAgentTaskRequest>> _normal_queue;
105
    std::condition_variable _high_prior_condv;
106
    std::deque<std::unique_ptr<TAgentTaskRequest>> _high_prior_queue;
107
108
    std::vector<std::shared_ptr<Thread>> _workers;
109
110
    std::function<void(const TAgentTaskRequest&)> _callback;
111
};
112
113
class ReportWorker {
114
public:
115
    ReportWorker(std::string name, const ClusterInfo* cluster_info, int report_interval_s,
116
                 std::function<void()> callback);
117
118
5
    ~ReportWorker();
119
120
    std::string_view name() const { return _name; }
121
122
    // Notify the worker to report immediately
123
    void notify();
124
125
    void stop();
126
127
private:
128
    std::string _name;
129
    std::shared_ptr<Thread> _thread;
130
131
    std::mutex _mtx;
132
    std::condition_variable _condv;
133
    bool _stopped {false};
134
    bool _signal {false};
135
};
136
137
void alter_inverted_index_callback(StorageEngine& engine, const TAgentTaskRequest& req);
138
139
void alter_cloud_index_callback(CloudStorageEngine& engine, const TAgentTaskRequest& req);
140
141
void check_consistency_callback(StorageEngine& engine, const TAgentTaskRequest& req);
142
143
void upload_callback(StorageEngine& engine, ExecEnv* env, const TAgentTaskRequest& req);
144
145
void download_callback(StorageEngine& engine, ExecEnv* env, const TAgentTaskRequest& req);
146
147
void download_callback(CloudStorageEngine& engine, ExecEnv* env, const TAgentTaskRequest& req);
148
149
void make_snapshot_callback(StorageEngine& engine, const TAgentTaskRequest& req);
150
151
void release_snapshot_callback(StorageEngine& engine, const TAgentTaskRequest& req);
152
153
void release_snapshot_callback(CloudStorageEngine& engine, const TAgentTaskRequest& req);
154
155
void move_dir_callback(StorageEngine& engine, ExecEnv* env, const TAgentTaskRequest& req);
156
157
void move_dir_callback(CloudStorageEngine& engine, ExecEnv* env, const TAgentTaskRequest& req);
158
159
void submit_table_compaction_callback(StorageEngine& engine, const TAgentTaskRequest& req);
160
161
void push_storage_policy_callback(StorageEngine& engine, const TAgentTaskRequest& req);
162
163
void push_index_policy_callback(const TAgentTaskRequest& req);
164
165
void push_cooldown_conf_callback(StorageEngine& engine, const TAgentTaskRequest& req);
166
167
void create_tablet_callback(StorageEngine& engine, const TAgentTaskRequest& req);
168
169
void drop_tablet_callback(StorageEngine& engine, const TAgentTaskRequest& req);
170
171
void drop_tablet_callback(CloudStorageEngine& engine, const TAgentTaskRequest& req);
172
173
void clear_transaction_task_callback(StorageEngine& engine, const TAgentTaskRequest& req);
174
175
void push_callback(StorageEngine& engine, const TAgentTaskRequest& req);
176
177
void cloud_push_callback(CloudStorageEngine& engine, const TAgentTaskRequest& req);
178
179
void update_tablet_meta_callback(StorageEngine& engine, const TAgentTaskRequest& req);
180
181
void alter_tablet_callback(StorageEngine& engine, const TAgentTaskRequest& req);
182
183
void alter_cloud_tablet_callback(CloudStorageEngine& engine, const TAgentTaskRequest& req);
184
185
void set_alter_version_before_enqueue(CloudStorageEngine& engine, const TAgentTaskRequest& req);
186
187
void clone_callback(StorageEngine& engine, const ClusterInfo* cluster_info,
188
                    const TAgentTaskRequest& req);
189
190
void storage_medium_migrate_callback(StorageEngine& engine, const TAgentTaskRequest& req);
191
192
void gc_binlog_callback(StorageEngine& engine, const TAgentTaskRequest& req);
193
194
void clean_trash_callback(StorageEngine& engine, const TAgentTaskRequest& req);
195
196
void clean_udf_cache_callback(const TAgentTaskRequest& req);
197
198
void visible_version_callback(StorageEngine& engine, const TAgentTaskRequest& req);
199
200
void report_task_callback(const ClusterInfo* cluster_info);
201
202
void report_disk_callback(StorageEngine& engine, const ClusterInfo* cluster_info);
203
204
void report_disk_callback(CloudStorageEngine& engine, const ClusterInfo* cluster_info);
205
206
void report_tablet_callback(StorageEngine& engine, const ClusterInfo* cluster_info);
207
208
void report_tablet_callback(CloudStorageEngine& engine, const ClusterInfo* cluster_info);
209
210
void calc_delete_bitmap_callback(CloudStorageEngine& engine, const TAgentTaskRequest& req);
211
212
void make_cloud_committed_rs_visible_callback(CloudStorageEngine& engine,
213
                                              const TAgentTaskRequest& req);
214
215
void report_index_policy_callback(const ClusterInfo* cluster_info);
216
217
} // namespace doris