Skip to content

Commit e0f1a98

Browse files
committed
Add threadpool/circular_queue and example/src ... && Remove mutex
1 parent b29f08b commit e0f1a98

15 files changed

Lines changed: 1267 additions & 129 deletions

examples/CMakeLists.txt

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
cmake_minimum_required(VERSION 3.15)
2+
3+
find_package(Threads REQUIRED)
4+
5+
file(GLOB_RECURSE EXAMPLE_FILES "*.cpp")
6+
foreach(file ${EXAMPLE_FILES})
7+
get_filename_component(name ${file} NAME_WE)
8+
add_executable(${name} ${file})
9+
target_link_libraries(${name} cppbase Threads::Threads)
10+
endforeach()
Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,79 @@
1+
#include <iostream>
2+
#include "cppbase/circular_queue.hpp"
3+
4+
void TestBasicPushPop()
5+
{
6+
std::cout << "=== Test CircularQueue Basic Push/Pop ===" << std::endl;
7+
8+
cppbase::CircularQueue<int> queue(5);
9+
10+
std::cout << "IsEmpty: " << queue.IsEmpty() << std::endl;
11+
std::cout << "IsFull: " << queue.IsFull() << std::endl;
12+
13+
for (int i = 1; i <= 5; ++i)
14+
{
15+
queue.PushBack(i * 10);
16+
}
17+
std::cout << "After pushing 5 elements, IsFull: " << queue.IsFull() << std::endl;
18+
std::cout << "Size: " << queue.GetSize() << std::endl;
19+
20+
std::cout << "Front: " << queue.Front() << std::endl;
21+
queue.PopFront();
22+
std::cout << "After pop, Size: " << queue.GetSize() << std::endl;
23+
24+
queue.PushBack(60);
25+
std::cout << "After overwrite, Front: " << queue.Front() << std::endl;
26+
std::cout << "OverRunCounter: " << queue.GetOverRunCounter() << std::endl;
27+
28+
std::cout << std::endl;
29+
}
30+
31+
void TestCircularOverwrite()
32+
{
33+
std::cout << "=== Test Circular Overwrite ===" << std::endl;
34+
35+
cppbase::CircularQueue<int> queue(3);
36+
37+
queue.PushBack(1);
38+
queue.PushBack(2);
39+
queue.PushBack(3);
40+
41+
std::cout << "After pushing 3 elements, Size: " << queue.GetSize() << std::endl;
42+
std::cout << "OverRunCounter: " << queue.GetOverRunCounter() << std::endl;
43+
44+
std::cout << "Front: " << queue.Front() << std::endl;
45+
queue.PopFront();
46+
47+
std::cout << "After pop, Front: " << queue.Front() << std::endl;
48+
queue.PopFront();
49+
50+
std::cout << "IsEmpty: " << queue.IsEmpty() << std::endl;
51+
52+
std::cout << std::endl;
53+
}
54+
55+
void TestMoveConstructor()
56+
{
57+
std::cout << "=== Test Move Constructor ===" << std::endl;
58+
59+
cppbase::CircularQueue<int> queue1(5);
60+
queue1.PushBack(10);
61+
queue1.PushBack(20);
62+
63+
cppbase::CircularQueue<int> queue2(std::move(queue1));
64+
65+
std::cout << "Queue2 Size: " << queue2.GetSize() << std::endl;
66+
std::cout << "Queue2 Front: " << queue2.Front() << std::endl;
67+
68+
std::cout << std::endl;
69+
}
70+
71+
int main()
72+
{
73+
TestBasicPushPop();
74+
TestCircularOverwrite();
75+
TestMoveConstructor();
76+
77+
std::cout << "All tests completed!" << std::endl;
78+
return 0;
79+
}

examples/mpmc_queue_example.cpp

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
#include <iostream>
2+
#include <thread>
3+
#include <string>
4+
#include <vector>
5+
#include "cppbase/mpmc_blocking_queue.hpp"
6+
7+
void TestMpmcBlockingQueue()
8+
{
9+
std::cout << "=== Test MpmcBlockingQueue ===" << std::endl;
10+
11+
cppbase::MpmcBlockingQueue<std::string> queue(10);
12+
13+
auto producer = [&queue](int id)
14+
{
15+
for (int i = 0; i < 3; ++i)
16+
{
17+
std::string msg = "Producer" + std::to_string(id) + "_Msg" + std::to_string(i);
18+
queue.Enqueue(std::move(msg));
19+
std::cout << "Producer " << id << " enqueued: " << i << std::endl;
20+
}
21+
};
22+
23+
auto consumer = [&queue](int id)
24+
{
25+
for (int i = 0; i < 3; ++i)
26+
{
27+
std::string msg;
28+
queue.DeQueue(msg);
29+
std::cout << "Consumer " << id << " dequeued: " << msg << std::endl;
30+
}
31+
};
32+
33+
std::thread p1(producer, 1);
34+
std::thread p2(producer, 2);
35+
std::thread c1(consumer, 1);
36+
std::thread c2(consumer, 2);
37+
38+
p1.join();
39+
p2.join();
40+
c1.join();
41+
c2.join();
42+
43+
std::cout << std::endl;
44+
}
45+
46+
void TestTryEnqueue()
47+
{
48+
std::cout << "=== Test TryEnqueue ===" << std::endl;
49+
50+
cppbase::MpmcBlockingQueue<int> queue(3);
51+
52+
queue.Enqueue_Immediate(1);
53+
queue.Enqueue_Immediate(2);
54+
queue.Enqueue_Immediate(3);
55+
56+
std::cout << "Queue full, try to enqueue more..." << std::endl;
57+
58+
queue.TryEnqueue(4);
59+
queue.TryEnqueue(5);
60+
61+
std::cout << "DropCount: " << queue.GetDropCount() << std::endl;
62+
std::cout << std::endl;
63+
}
64+
65+
void TestTryDeQueue()
66+
{
67+
std::cout << "=== Test TryDeQueue ===" << std::endl;
68+
69+
cppbase::MpmcBlockingQueue<int> queue(5);
70+
71+
int value;
72+
73+
bool success = queue.TryDeQueue(value, std::chrono::milliseconds(100));
74+
std::cout << "TryDeQueue from empty (timeout 100ms): " << (success ? "success" : "timeout") << std::endl;
75+
76+
queue.Enqueue(42);
77+
success = queue.TryDeQueue(value, std::chrono::milliseconds(100));
78+
std::cout << "TryDeQueue after enqueue: " << value << ", success: " << success << std::endl;
79+
std::cout << std::endl;
80+
}
81+
82+
int main()
83+
{
84+
TestMpmcBlockingQueue();
85+
TestTryEnqueue();
86+
TestTryDeQueue();
87+
88+
std::cout << "All tests passed!" << std::endl;
89+
return 0;
90+
}

examples/threadpool_example.cpp

Lines changed: 177 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,177 @@
1+
#include <iostream>
2+
#include <thread>
3+
#include <chrono>
4+
#include <atomic>
5+
#include "cppbase/threadpool.hpp"
6+
#include "cppbase/err_code.hpp"
7+
8+
void TestBasicUsage()
9+
{
10+
std::cout << "=== Test Basic Usage ===" << std::endl;
11+
12+
cppbase::ThreadPool pool;
13+
int ret = pool.Start(4, 100);
14+
if (ret != cppbase::ErrorCode::EC_OK)
15+
{
16+
std::cout << "Failed to start thread pool, error: " << ret << std::endl;
17+
return;
18+
}
19+
20+
std::cout << "Thread pool started" << std::endl;
21+
22+
for (int i = 0; i < 5; ++i)
23+
{
24+
pool.AddTask([i]()
25+
{
26+
std::cout << "Task " << i << " executed by thread" << std::endl;
27+
});
28+
}
29+
30+
std::this_thread::sleep_for(std::chrono::milliseconds(100));
31+
32+
std::cout << std::endl;
33+
}
34+
35+
void TestSubmitWithFuture()
36+
{
37+
std::cout << "=== Test Submit with Future ===" << std::endl;
38+
39+
cppbase::ThreadPool pool;
40+
pool.Start(4, 100);
41+
42+
auto future1 = pool.Submit([]()
43+
{
44+
std::this_thread::sleep_for(std::chrono::milliseconds(50));
45+
return 42;
46+
});
47+
48+
auto future2 = pool.Submit([]()
49+
{
50+
return 10 + 20;
51+
});
52+
53+
auto future3 = pool.Submit([]()
54+
{
55+
std::string msg = "Hello ThreadPool";
56+
return msg.length();
57+
});
58+
59+
std::cout << "Future1 result: " << future1.get() << std::endl;
60+
std::cout << "Future2 result: " << future2.get() << std::endl;
61+
std::cout << "Future3 result: " << future3.get() << std::endl;
62+
63+
std::cout << std::endl;
64+
}
65+
66+
void TestAddTaskWithTimeout()
67+
{
68+
std::cout << "=== Test AddTask with Timeout ===" << std::endl;
69+
70+
cppbase::ThreadPool pool;
71+
pool.Start(1, 2);
72+
73+
pool.AddTask([]()
74+
{
75+
std::this_thread::sleep_for(std::chrono::milliseconds(500));
76+
std::cout << "Task 1 completed" << std::endl;
77+
});
78+
pool.AddTask([]()
79+
{
80+
std::this_thread::sleep_for(std::chrono::milliseconds(500));
81+
std::cout << "Task 2 completed" << std::endl;
82+
});
83+
84+
std::cout << "Queue full, trying to add with timeout..." << std::endl;
85+
86+
int ret = pool.AddTask([]()
87+
{
88+
std::cout << "Task 3" << std::endl;
89+
}, std::chrono::milliseconds(100));
90+
91+
if (ret == cppbase::ErrorCode::EC_TIMEOUT)
92+
{
93+
std::cout << "AddTask timeout" << std::endl;
94+
}
95+
else if (ret == cppbase::ErrorCode::EC_OK)
96+
{
97+
std::cout << "AddTask success" << std::endl;
98+
}
99+
100+
std::this_thread::sleep_for(std::chrono::seconds(1));
101+
102+
std::cout << std::endl;
103+
}
104+
105+
void TestMultiThreaded()
106+
{
107+
std::cout << "=== Test Multi-threaded ===" << std::endl;
108+
109+
cppbase::ThreadPool pool;
110+
pool.Start(4, 100);
111+
112+
const int numTasks = 20;
113+
std::atomic<int> counter{0};
114+
115+
std::vector<std::thread> threads;
116+
for (int t = 0; t < 4; ++t)
117+
{
118+
threads.emplace_back([&pool, &counter, t, numTasks]()
119+
{
120+
for (int i = 0; i < numTasks; ++i)
121+
{
122+
pool.AddTask([&counter, t, i]()
123+
{
124+
std::this_thread::sleep_for(std::chrono::milliseconds(10));
125+
counter++;
126+
});
127+
}
128+
});
129+
}
130+
131+
for (auto& th : threads)
132+
{
133+
th.join();
134+
}
135+
136+
std::this_thread::sleep_for(std::chrono::milliseconds(500));
137+
138+
std::cout << "Expected: " << numTasks * 4 << ", Actual: " << counter << std::endl;
139+
std::cout << std::endl;
140+
}
141+
142+
void TestExceptionHandling()
143+
{
144+
std::cout << "=== Test Exception Handling ===" << std::endl;
145+
146+
cppbase::ThreadPool pool;
147+
pool.Start(2, 10);
148+
149+
auto future = pool.Submit([]()
150+
{
151+
throw std::runtime_error("Test exception");
152+
return 0;
153+
});
154+
155+
try
156+
{
157+
future.get();
158+
}
159+
catch (const std::exception& e)
160+
{
161+
std::cout << "Caught exception: " << e.what() << std::endl;
162+
}
163+
164+
std::cout << std::endl;
165+
}
166+
167+
int main()
168+
{
169+
TestBasicUsage();
170+
TestSubmitWithFuture();
171+
TestAddTaskWithTimeout();
172+
TestMultiThreaded();
173+
TestExceptionHandling();
174+
175+
std::cout << "All tests completed!" << std::endl;
176+
return 0;
177+
}

0 commit comments

Comments
 (0)