7#include <packager/app/job_manager.h>
16#include <absl/log/check.h>
17#include <absl/synchronization/mutex.h>
19#include <packager/media/chunking/sync_point_queue.h>
20#include <packager/media/origin/origin_handler.h>
21#include <packager/status.h>
26Job::Job(
const std::string& name,
27 std::shared_ptr<OriginHandler> work,
28 OnCompleteFunction on_complete)
30 work_(std::move(work)),
31 on_complete_(on_complete),
32 status_(error::Code::UNKNOWN,
"Job uninitialized") {
36const Status& Job::Initialize() {
37 status_ = work_->Initialize();
42 thread_.reset(
new std::thread(&Job::Run,
this));
49const Status& Job::Run() {
51 status_ = work_->Run();
65JobManager::JobManager(std::unique_ptr<SyncPointQueue> sync_points)
66 : sync_points_(std::move(sync_points)) {}
68void JobManager::Add(
const std::string& name,
69 std::shared_ptr<OriginHandler> handler) {
70 jobs_.emplace_back(
new Job(
71 name, std::move(handler),
72 std::bind(&JobManager::OnJobComplete,
this, std::placeholders::_1)));
75Status JobManager::InitializeJobs() {
77 for (
auto& job : jobs_)
78 status.Update(job->Initialize());
82Status JobManager::RunJobs() {
83 std::set<Job*> active_jobs;
87 for (
auto& job : jobs_) {
90 active_jobs.insert(job.get());
96 absl::MutexLock lock(mutex_);
97 while (status.ok() && active_jobs.size()) {
99 any_job_complete_.Wait(&mutex_);
102 for (
const auto& entry : complete_) {
103 Job* job = entry.first;
104 bool complete = entry.second;
107 status.Update(job->status());
108 active_jobs.erase(job);
117 sync_points_->Cancel();
119 for (
auto& job : active_jobs)
122 for (
auto& job : active_jobs)
128void JobManager::OnJobComplete(Job* job) {
129 absl::MutexLock lock(mutex_);
131 complete_[job] =
true;
132 any_job_complete_.Signal();
135void JobManager::CancelJobs() {
137 sync_points_->Cancel();
139 for (
auto& job : jobs_)
All the methods that are virtual are virtual for mocking.