载入中...
搜索中...
未找到
Thread.cpp
浏览该文件的文档.
1#include "thread/Thread.h"
2
3#include "common/AsyncWork.h"
4#include "common/Diagnostic.h"
5#include "common/Exception.h"
6#include "common/Capability.h"
8#include "common/Result.h"
9
10#include <simplesquirrel/simplesquirrel.hpp>
11
12#include <algorithm>
13#include <functional>
14#include <thread>
15
16namespace eve {
17namespace thread {
18
19namespace {
20
21class PoolExecutor final : public eve::caps::IAsyncWorkExecutor {
22public:
23 explicit PoolExecutor(Thread *owner) : pool_(std::min(8, owner->getHardwareConcurrency())) {}
24
25 eve::Result<void> submit(std::function<void()> work) override {
26 if (!work) {
28 eve::DiagnosticCode::InvalidArgument, "IAsyncWorkExecutor: null work"));
29 }
30 try {
31 std::unique_ptr<Task> task(pool_.submit(std::move(work)));
33 } catch (const eve::Exception &ex) {
36 }
37 }
38
39private:
40 // CPU resource decoders are allocation/IO heavy. Bound admission independently
41 // of the general-purpose pool and drain accepted work during destruction.
42 ThreadPool pool_;
43};
44
45} // namespace
46
48
50 executor_ = std::make_unique<PoolExecutor>(this);
51 cap::provide<caps::IAsyncWorkExecutor>(executor_.get());
52}
53
55 cap::revoke<caps::IAsyncWorkExecutor>(executor_.get());
56 executor_.reset();
57 if (defaultPool_)
58 defaultPool_->stop();
59 if (defaultJobSystem_)
60 defaultJobSystem_->stop();
61}
62
64 unsigned n = std::thread::hardware_concurrency();
65 return n == 0 ? 1 : static_cast<int>(n);
66}
67
69 std::lock_guard<std::mutex> lock(mu_);
70 if (!defaultPool_)
71 defaultPool_ = std::make_unique<ThreadPool>(getHardwareConcurrency());
72 return defaultPool_.get();
73}
74
76 if (workerCount < 0)
77 throw eve::Exception("Thread::newThreadPool: workerCount must be >= 0");
78 return new ThreadPool(workerCount);
79}
80
82 std::lock_guard<std::mutex> lock(mu_);
83 if (!defaultJobSystem_)
84 defaultJobSystem_ = std::unique_ptr<JobSystem>(createJobSystem(getHardwareConcurrency()));
85 return defaultJobSystem_.get();
86}
87
88Channel *Thread::getChannel(std::string channelName) {
89 if (channelName.empty())
90 throw eve::Exception("Thread::getChannel: name must not be empty");
91 std::lock_guard<std::mutex> lock(mu_);
92 auto it = namedChannels_.find(channelName);
93 if (it != namedChannels_.end())
94 return it->second.get();
95 auto ch = std::make_unique<Channel>(channelName);
96 Channel *raw = ch.get();
97 namedChannels_.emplace(std::move(channelName), std::move(ch));
98 return raw;
99}
100
102
103void Thread::postMain(std::string eventName, std::string data) {
104 if (eventName.empty())
105 throw eve::Exception("Thread::postMain: name must not be empty");
106 auto *poster = cap::query<caps::IMainThreadPost>();
107 if (!poster)
108 throw eve::Exception("Thread::postMain: no main-thread queue (event module not linked)");
109 poster->prepare();
110 poster->postToMainThread(std::move(eventName), std::move(data));
111}
112
113void Thread::expose(ssq::Table &table) {
114 auto cls = table.addClass(name, Thread::create, false);
115 expose(cls);
116
117 // Avoid clashing with network::Channel ("Channel").
118 auto ch = table.addClass<Channel>(
119 "ThreadChannel", std::function<Channel *()>([]() -> Channel * { return nullptr; }), true);
120 ch.addFunc("getName", &Channel::getName);
121 ch.addFunc("push", &Channel::push);
122 ch.addFunc("pop", &Channel::pop);
123 ch.addFunc("demand", &Channel::demand);
124 ch.addFunc("supply", &Channel::supply);
125 ch.addFunc("hasData", &Channel::hasData);
126 ch.addFunc("getCount", &Channel::getCount);
127 ch.addFunc("clear", &Channel::clear);
128
129 auto task = table.addClass<Task>(
130 "Task", std::function<Task *()>([]() -> Task * { return nullptr; }), true);
131 task.addFunc("getStatus", &Task::getStatus);
132 task.addFunc("isDone", &Task::isDone);
133 task.addFunc("hasFailed", &Task::hasFailed);
134 task.addFunc("getError", &Task::getError);
135 task.addFunc("wait", &Task::wait);
136
137 auto pool = table.addClass<ThreadPool>(
138 "ThreadPool", std::function<ThreadPool *()>([]() -> ThreadPool * { return nullptr; }), true);
139 pool.addFunc("getWorkerCount", &ThreadPool::getWorkerCount);
140 pool.addFunc("getPendingCount", &ThreadPool::getPendingCount);
141 pool.addFunc("isRunning", &ThreadPool::isRunning);
142 pool.addFunc("submitSleep", &ThreadPool::submitSleep);
143 pool.addFunc("submitPush", &ThreadPool::submitPush);
144 pool.addFunc("submitPost", &ThreadPool::submitPost);
145 pool.addFunc("waitAll", &ThreadPool::waitAll);
146 pool.addFunc("stop", &ThreadPool::stop);
147}
148
149void Thread::expose(ssq::Class &cls) {
150 cls.addFunc("getName", &Thread::getName);
151 cls.addFunc("getHardwareConcurrency", &Thread::getHardwareConcurrency);
152 cls.addFunc("getPool", &Thread::getPool);
153 cls.addFunc("newThreadPool", &Thread::newThreadPool);
154 cls.addFunc("getChannel", &Thread::getChannel);
155 cls.addFunc("newChannel", &Thread::newChannel);
156 cls.addFunc("postMain", &Thread::postMain);
157}
158
159} // namespace thread
160} // namespace eve
Capability for submitting CPU work onto a worker pool.
Stable, structured diagnostics shared by engine modules.
HSQOBJECT cls
Definition ECS.cpp:21
glm::vec3 n
Definition Grass.cpp:63
std::string name
#define Module_IMPL(ModuleName, newExpr)
Definition Module.h:26
Move-only, checked operation results for the common layer.
static Diagnostic error(DiagnosticCode code, std::string message, std::string path={}, DiagnosticDetails details={}, std::string source={})
Construct an error diagnostic with the standard error severity.
Definition Diagnostic.h:125
EVENGINE_API_FOUNDATION public API.
Definition Exception.h:13
virtual const char * what() const
Returns a string containing reason for the exception.
Definition Exception.h:24
virtual std::string getName() const =0
Returns the name.
Move-only operation result carrying either a value or Status.
Definition Result.h:155
static Result success(T value)
Construct a successful result owning value.
Definition Result.h:164
static Result failure(Status status)
Construct a failed result from a structured status.
Definition Result.h:175
static Status success(StatusCode code=StatusCode::Ok)
Construct a successful status with an explicit non-error outcome.
Definition Status.h:81
Fire-and-forget CPU job submission used by the resource cache.
Definition AsyncWork.h:26
Thread-safe message queue (love2d-style Channel). Values are strings so the API stays overload-free f...
Definition Channel.h:19
void push(std::string value)
Pushes .
Definition Channel.cpp:16
bool hasData() const
True when data.
Definition Channel.cpp:54
std::string pop()
Non-blocking pop; returns "" if empty.
Definition Channel.cpp:24
std::string demand()
Block until a value is available, then pop it.
Definition Channel.cpp:33
void clear()
Clears .
Definition Channel.cpp:64
int getCount() const
Returns the count.
Definition Channel.cpp:59
std::string getName() const
Returns the name.
Definition Channel.cpp:14
std::string supply(int timeoutMs)
Block up to timeoutMs; returns "" on timeout.
Definition Channel.cpp:42
Abstract job scheduler: dependencies, parallel_for, task_group, arena.
Definition JobSystem.h:157
bool hasFailed() const
True when failed.
Definition Task.cpp:28
void wait()
Block until the task finishes (done or failed).
Definition Task.cpp:38
std::string getError() const
Returns the error.
Definition Task.cpp:33
bool isDone() const
True when done.
Definition Task.cpp:23
std::string getStatus() const
Returns the status.
Definition Task.cpp:18
Fixed-size worker pool. Owns worker std::threads; tasks run FIFO. Squirrel VM is not thread-safe — do...
Definition ThreadPool.h:26
void waitAll()
Block until idle. Throws when called by a worker belonging to this pool.
int getPendingCount() const
Returns the pending count.
Task * submitSleep(int ms)
Sleep on a worker, then mark done — useful from scripts / tests.
void stop()
Stop accepting work and join workers. Worker calls are rejected. Idempotent.
bool isRunning() const
True when running.
Task * submitPost(std::string name, std::string data="", int delayMs=0)
Sleep, then post an Event on the main queue (thread-safe). Scripts poll via event....
int getWorkerCount() const
Returns the worker count.
Task * submitPush(Channel *channel, std::string message, int delayMs=0)
Sleep, then push a message onto a channel (cross-thread signalling).
Thread module: default pool, named channels, pool factory. Script: eve.Thread() → getPool / newThread...
Definition Thread.h:24
void postMain(std::string eventName, std::string data="")
Thread-safe post onto the Event module queue (any thread). Main loop: event.pump is unrelated; just e...
Definition Thread.cpp:103
Channel * newChannel()
Anonymous channel (not registered in the name map). @ownership Owned; caller must delete when done....
Definition Thread.cpp:101
ThreadPool * getPool()
Shared default pool (created lazily with hardwareConcurrency workers). @ownership Borrowed from this ...
Definition Thread.cpp:68
JobSystem * getJobSystem()
Shared engine-wide JobSystem (created lazily with hardware concurrency workers). Use it for dependenc...
Definition Thread.cpp:81
ThreadPool * newThreadPool(int workerCount=0)
Create an independent pool. @ownership Owned; caller must delete when done. @lifetime Valid until the...
Definition Thread.cpp:75
Channel * getChannel(std::string channelName)
Named shared channel (love2d-style). Same name → same Channel instance for the lifetime of the module...
Definition Thread.cpp:88
int getHardwareConcurrency() const
Hardware concurrency hint (at least 1).
Definition Thread.cpp:63
~Thread() override
Thread.
Definition Thread.cpp:54
JobSystem * createJobSystem(int workerCount)
Create a JobSystem with the configured backend.
Definition JobSystem.cpp:14
Build metadata (engine git commit, build time, third-party version).
Definition Build.cpp:16