Directory: | cvmfs/ |
---|---|
File: | cvmfs/ingestion/task.h |
Date: | 2025-02-09 02:34:19 |
Exec | Total | Coverage | |
---|---|---|---|
Lines: | 43 | 44 | 97.7% |
Branches: | 17 | 24 | 70.8% |
Line | Branch | Exec | Source |
---|---|---|---|
1 | /** | ||
2 | * This file is part of the CernVM File System. | ||
3 | */ | ||
4 | |||
5 | #ifndef CVMFS_INGESTION_TASK_H_ | ||
6 | #define CVMFS_INGESTION_TASK_H_ | ||
7 | |||
8 | #include <errno.h> | ||
9 | #include <pthread.h> | ||
10 | #include <unistd.h> | ||
11 | |||
12 | #include <cassert> | ||
13 | #include <vector> | ||
14 | |||
15 | #include "util/exception.h" | ||
16 | #include "util/single_copy.h" | ||
17 | #include "util/tube.h" | ||
18 | |||
19 | /** | ||
20 | * Forward declaration of TubeConsumerGroup so that it can be used as a friend | ||
21 | * class to TubeConsumer. | ||
22 | */ | ||
23 | template<typename ItemT> | ||
24 | class TubeConsumerGroup; | ||
25 | |||
26 | |||
27 | /** | ||
28 | * Base class for threads that processes items from a tube one by one. Concrete | ||
29 | * implementations overwrite the Process() method. | ||
30 | */ | ||
31 | template <class ItemT> | ||
32 | class TubeConsumer : SingleCopy { | ||
33 | friend class TubeConsumerGroup<ItemT>; | ||
34 | |||
35 | public: | ||
36 | 13076 | virtual ~TubeConsumer() { } | |
37 | |||
38 | protected: | ||
39 | 12940 | explicit TubeConsumer(Tube<ItemT> *tube) : tube_(tube) { } | |
40 | virtual void Process(ItemT *item) = 0; | ||
41 | 11964 | virtual void OnTerminate() { } | |
42 | |||
43 | Tube<ItemT> *tube_; | ||
44 | |||
45 | private: | ||
46 | 12172 | static void *MainConsumer(void *data) { | |
47 | 12172 | TubeConsumer<ItemT> *consumer = | |
48 | reinterpret_cast<TubeConsumer<ItemT> *>(data); | ||
49 | |||
50 | 6518832 | while (true) { | |
51 | 6531004 | ItemT *item = consumer->tube_->PopFront(); | |
52 |
2/2✓ Branch 1 taken 6117 times.
✓ Branch 2 taken 3561651 times.
|
6473428 | if (item->IsQuitBeacon()) { |
53 |
1/2✓ Branch 0 taken 6117 times.
✗ Branch 1 not taken.
|
12098 | delete item; |
54 | 12082 | break; | |
55 | } | ||
56 | 6452090 | consumer->Process(item); | |
57 | } | ||
58 | 12082 | consumer->OnTerminate(); | |
59 | 12080 | return NULL; | |
60 | } | ||
61 | }; | ||
62 | |||
63 | |||
64 | template <class ItemT> | ||
65 | class TubeConsumerGroup : SingleCopy { | ||
66 | public: | ||
67 | 684 | TubeConsumerGroup() : is_active_(false) { } | |
68 | |||
69 | 684 | ~TubeConsumerGroup() { | |
70 |
2/2✓ Branch 1 taken 6538 times.
✓ Branch 2 taken 398 times.
|
13624 | for (unsigned i = 0; i < consumers_.size(); ++i) |
71 |
1/2✓ Branch 1 taken 6538 times.
✗ Branch 2 not taken.
|
12940 | delete consumers_[i]; |
72 | 684 | } | |
73 | |||
74 | 12940 | void TakeConsumer(TubeConsumer<ItemT> *consumer) { | |
75 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 6538 times.
|
12940 | assert(!is_active_); |
76 | 12940 | consumers_.push_back(consumer); | |
77 | 12940 | } | |
78 | |||
79 | 660 | void Spawn() { | |
80 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 386 times.
|
660 | assert(!is_active_); |
81 | 660 | unsigned N = consumers_.size(); | |
82 | 660 | threads_.resize(N); | |
83 |
2/2✓ Branch 0 taken 6154 times.
✓ Branch 1 taken 386 times.
|
12832 | for (unsigned i = 0; i < N; ++i) { |
84 | 12172 | int retval = pthread_create( | |
85 | 12172 | &threads_[i], NULL, TubeConsumer<ItemT>::MainConsumer, consumers_[i]); | |
86 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 6154 times.
|
12172 | if (retval != 0) { |
87 | ✗ | PANIC(kLogStderr, "failed to create new thread (error: %d, pid: %d)", | |
88 | errno, getpid()); | ||
89 | } | ||
90 | } | ||
91 | 660 | is_active_ = true; | |
92 | 660 | } | |
93 | |||
94 | 660 | void Terminate() { | |
95 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 386 times.
|
660 | assert(is_active_); |
96 | 660 | unsigned N = consumers_.size(); | |
97 |
2/2✓ Branch 0 taken 6154 times.
✓ Branch 1 taken 386 times.
|
12832 | for (unsigned i = 0; i < N; ++i) { |
98 | 12172 | consumers_[i]->tube_->EnqueueBack(ItemT::CreateQuitBeacon()); | |
99 | } | ||
100 |
2/2✓ Branch 0 taken 6154 times.
✓ Branch 1 taken 386 times.
|
12832 | for (unsigned i = 0; i < N; ++i) { |
101 | 12172 | int retval = pthread_join(threads_[i], NULL); | |
102 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 6154 times.
|
12172 | assert(retval == 0); |
103 | } | ||
104 | 660 | is_active_ = false; | |
105 | 660 | } | |
106 | |||
107 | 112 | bool is_active() { return is_active_; } | |
108 | |||
109 | private: | ||
110 | bool is_active_; | ||
111 | std::vector<TubeConsumer<ItemT> *> consumers_; | ||
112 | std::vector<pthread_t> threads_; | ||
113 | }; | ||
114 | |||
115 | #endif // CVMFS_INGESTION_TASK_H_ | ||
116 |