Rename runners to workers
This commit is contained in:
parent
78f945647c
commit
6f4793bcf2
20 changed files with 233 additions and 237 deletions
|
|
@ -3,8 +3,8 @@ mod commit;
|
|||
mod index;
|
||||
mod link;
|
||||
mod queue;
|
||||
mod runner;
|
||||
mod r#static;
|
||||
mod worker;
|
||||
|
||||
use axum::{routing::get, Router};
|
||||
|
||||
|
|
@ -46,7 +46,7 @@ pub async fn run(server: Server) -> somehow::Result<()> {
|
|||
let app = Router::new()
|
||||
.route("/", get(index::get))
|
||||
.route("/commit/:hash", get(commit::get))
|
||||
.route("/runner/:name", get(runner::get))
|
||||
.route("/worker/:name", get(worker::get))
|
||||
.route("/queue/", get(queue::get))
|
||||
.route("/queue/inner", get(queue::get_inner))
|
||||
.merge(api::router(&server))
|
||||
|
|
|
|||
|
|
@ -17,10 +17,10 @@ use tracing::debug;
|
|||
use crate::{
|
||||
config::Config,
|
||||
server::{
|
||||
runners::{RunnerInfo, Runners},
|
||||
workers::{WorkerInfo, Workers},
|
||||
BenchRepo, Server,
|
||||
},
|
||||
shared::{BenchMethod, RunnerRequest, ServerResponse, Work},
|
||||
shared::{BenchMethod, ServerResponse, Work, WorkerRequest},
|
||||
somehow,
|
||||
};
|
||||
|
||||
|
|
@ -29,8 +29,8 @@ async fn post_status(
|
|||
State(config): State<&'static Config>,
|
||||
State(db): State<SqlitePool>,
|
||||
State(bench_repo): State<Option<BenchRepo>>,
|
||||
State(runners): State<Arc<Mutex<Runners>>>,
|
||||
Json(request): Json<RunnerRequest>,
|
||||
State(workers): State<Arc<Mutex<Workers>>>,
|
||||
Json(request): Json<WorkerRequest>,
|
||||
) -> somehow::Result<Response> {
|
||||
let name = match auth::authenticate(config, auth) {
|
||||
Ok(name) => name,
|
||||
|
|
@ -46,14 +46,14 @@ async fn post_status(
|
|||
.fetch_all(&db)
|
||||
.await?;
|
||||
|
||||
let mut guard = runners.lock().unwrap();
|
||||
let mut guard = workers.lock().unwrap();
|
||||
guard.clean();
|
||||
if !guard.verify(&name, &request.secret) {
|
||||
return Ok((StatusCode::UNAUTHORIZED, "invalid secret").into_response());
|
||||
}
|
||||
guard.update(
|
||||
name.clone(),
|
||||
RunnerInfo::new(request.secret, OffsetDateTime::now_utc(), request.status),
|
||||
WorkerInfo::new(request.secret, OffsetDateTime::now_utc(), request.status),
|
||||
);
|
||||
let work = match request.request_work {
|
||||
true => guard.find_free_work(&queue),
|
||||
|
|
@ -91,5 +91,5 @@ pub fn router(server: &Server) -> Router<Server> {
|
|||
|
||||
// TODO Get repo tar
|
||||
// TODO Get bench repo tar
|
||||
Router::new().route("/api/runner/status", post(post_status))
|
||||
Router::new().route("/api/worker/status", post(post_status))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -19,7 +19,7 @@ fn is_username_valid(username: &str) -> bool {
|
|||
}
|
||||
|
||||
fn is_password_valid(password: &str, config: &'static Config) -> bool {
|
||||
password == config.web_runner_token
|
||||
password == config.web_worker_token
|
||||
}
|
||||
|
||||
pub fn authenticate(
|
||||
|
|
@ -36,7 +36,7 @@ pub fn authenticate(
|
|||
StatusCode::UNAUTHORIZED,
|
||||
[(
|
||||
header::WWW_AUTHENTICATE,
|
||||
HeaderValue::from_str("Basic realm=\"runner api\"").unwrap(),
|
||||
HeaderValue::from_str("Basic realm=\"worker api\"").unwrap(),
|
||||
)],
|
||||
"invalid credentials",
|
||||
)
|
||||
|
|
|
|||
|
|
@ -62,17 +62,17 @@ impl RunLink {
|
|||
#[template(
|
||||
ext = "html",
|
||||
source = "\
|
||||
<a href=\"{{ root }}runner/{{ name }}\">
|
||||
<a href=\"{{ root }}worker/{{ name }}\">
|
||||
{{ name }}
|
||||
</a>
|
||||
"
|
||||
)]
|
||||
pub struct RunnerLink {
|
||||
pub struct WorkerLink {
|
||||
root: String,
|
||||
name: String,
|
||||
}
|
||||
|
||||
impl RunnerLink {
|
||||
impl WorkerLink {
|
||||
pub fn new(base: &Base, name: String) -> Self {
|
||||
Self {
|
||||
root: base.root.clone(),
|
||||
|
|
|
|||
|
|
@ -11,15 +11,15 @@ use sqlx::SqlitePool;
|
|||
use crate::{
|
||||
config::Config,
|
||||
server::{
|
||||
runners::{RunnerInfo, Runners},
|
||||
util,
|
||||
workers::{WorkerInfo, Workers},
|
||||
},
|
||||
shared::RunnerStatus,
|
||||
shared::WorkerStatus,
|
||||
somehow,
|
||||
};
|
||||
|
||||
use super::{
|
||||
link::{CommitLink, RunLink, RunnerLink},
|
||||
link::{CommitLink, RunLink, WorkerLink},
|
||||
Base, Tab,
|
||||
};
|
||||
|
||||
|
|
@ -29,8 +29,8 @@ enum Status {
|
|||
Working(RunLink),
|
||||
}
|
||||
|
||||
struct Runner {
|
||||
link: RunnerLink,
|
||||
struct Worker {
|
||||
link: WorkerLink,
|
||||
status: Status,
|
||||
}
|
||||
|
||||
|
|
@ -38,33 +38,33 @@ struct Task {
|
|||
commit: CommitLink,
|
||||
since: String,
|
||||
priority: i64,
|
||||
runners: Vec<RunnerLink>,
|
||||
workers: Vec<WorkerLink>,
|
||||
odd: bool,
|
||||
}
|
||||
|
||||
fn sorted_runners(runners: &Mutex<Runners>) -> Vec<(String, RunnerInfo)> {
|
||||
let mut runners = runners
|
||||
fn sorted_workers(workers: &Mutex<Workers>) -> Vec<(String, WorkerInfo)> {
|
||||
let mut workers = workers
|
||||
.lock()
|
||||
.unwrap()
|
||||
.clean()
|
||||
.get_all()
|
||||
.into_iter()
|
||||
.collect::<Vec<_>>();
|
||||
runners.sort_unstable_by(|(a, _), (b, _)| a.cmp(b));
|
||||
runners
|
||||
workers.sort_unstable_by(|(a, _), (b, _)| a.cmp(b));
|
||||
workers
|
||||
}
|
||||
|
||||
async fn get_runners(
|
||||
async fn get_workers(
|
||||
db: &SqlitePool,
|
||||
runners: &[(String, RunnerInfo)],
|
||||
workers: &[(String, WorkerInfo)],
|
||||
base: &Base,
|
||||
) -> somehow::Result<Vec<Runner>> {
|
||||
) -> somehow::Result<Vec<Worker>> {
|
||||
let mut result = vec![];
|
||||
for (name, info) in runners {
|
||||
for (name, info) in workers {
|
||||
let status = match &info.status {
|
||||
RunnerStatus::Idle => Status::Idle,
|
||||
RunnerStatus::Busy => Status::Busy,
|
||||
RunnerStatus::Working(run) => {
|
||||
WorkerStatus::Idle => Status::Idle,
|
||||
WorkerStatus::Busy => Status::Busy,
|
||||
WorkerStatus::Working(run) => {
|
||||
let message =
|
||||
sqlx::query_scalar!("SELECT message FROM commits WHERE hash = ?", run.hash)
|
||||
.fetch_one(db)
|
||||
|
|
@ -73,8 +73,8 @@ async fn get_runners(
|
|||
}
|
||||
};
|
||||
|
||||
result.push(Runner {
|
||||
link: RunnerLink::new(base, name.clone()),
|
||||
result.push(Worker {
|
||||
link: WorkerLink::new(base, name.clone()),
|
||||
status,
|
||||
})
|
||||
}
|
||||
|
|
@ -83,17 +83,17 @@ async fn get_runners(
|
|||
|
||||
async fn get_queue(
|
||||
db: &SqlitePool,
|
||||
runners: &[(String, RunnerInfo)],
|
||||
workers: &[(String, WorkerInfo)],
|
||||
base: &Base,
|
||||
) -> somehow::Result<Vec<Task>> {
|
||||
// Group runners by commit hash
|
||||
let mut runners_by_commit: HashMap<String, Vec<RunnerLink>> = HashMap::new();
|
||||
for (name, info) in runners {
|
||||
if let RunnerStatus::Working(run) = &info.status {
|
||||
runners_by_commit
|
||||
// Group workers by commit hash
|
||||
let mut workers_by_commit: HashMap<String, Vec<WorkerLink>> = HashMap::new();
|
||||
for (name, info) in workers {
|
||||
if let WorkerStatus::Working(run) = &info.status {
|
||||
workers_by_commit
|
||||
.entry(run.hash.clone())
|
||||
.or_default()
|
||||
.push(RunnerLink::new(base, name.clone()));
|
||||
.push(WorkerLink::new(base, name.clone()));
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -112,7 +112,7 @@ async fn get_queue(
|
|||
)
|
||||
.fetch(db)
|
||||
.map_ok(|r| Task {
|
||||
runners: runners_by_commit.remove(&r.hash).unwrap_or_default(),
|
||||
workers: workers_by_commit.remove(&r.hash).unwrap_or_default(),
|
||||
commit: CommitLink::new(base, r.hash, &r.message, r.reachable),
|
||||
since: util::format_delta_from_now(r.date),
|
||||
priority: r.priority,
|
||||
|
|
@ -137,20 +137,20 @@ async fn get_queue(
|
|||
#[derive(Template)]
|
||||
#[template(path = "queue_inner.html")]
|
||||
struct QueueInnerTemplate {
|
||||
runners: Vec<Runner>,
|
||||
workers: Vec<Worker>,
|
||||
tasks: Vec<Task>,
|
||||
}
|
||||
|
||||
pub async fn get_inner(
|
||||
State(config): State<&'static Config>,
|
||||
State(db): State<SqlitePool>,
|
||||
State(runners): State<Arc<Mutex<Runners>>>,
|
||||
State(workers): State<Arc<Mutex<Workers>>>,
|
||||
) -> somehow::Result<impl IntoResponse> {
|
||||
let base = Base::new(config, Tab::Queue);
|
||||
let sorted_runners = sorted_runners(&runners);
|
||||
let runners = get_runners(&db, &sorted_runners, &base).await?;
|
||||
let tasks = get_queue(&db, &sorted_runners, &base).await?;
|
||||
Ok(QueueInnerTemplate { runners, tasks })
|
||||
let sorted_workers = sorted_workers(&workers);
|
||||
let workers = get_workers(&db, &sorted_workers, &base).await?;
|
||||
let tasks = get_queue(&db, &sorted_workers, &base).await?;
|
||||
Ok(QueueInnerTemplate { workers, tasks })
|
||||
}
|
||||
#[derive(Template)]
|
||||
#[template(path = "queue.html")]
|
||||
|
|
@ -162,14 +162,14 @@ struct QueueTemplate {
|
|||
pub async fn get(
|
||||
State(config): State<&'static Config>,
|
||||
State(db): State<SqlitePool>,
|
||||
State(runners): State<Arc<Mutex<Runners>>>,
|
||||
State(workers): State<Arc<Mutex<Workers>>>,
|
||||
) -> somehow::Result<impl IntoResponse> {
|
||||
let base = Base::new(config, Tab::Queue);
|
||||
let sorted_runners = sorted_runners(&runners);
|
||||
let runners = get_runners(&db, &sorted_runners, &base).await?;
|
||||
let tasks = get_queue(&db, &sorted_runners, &base).await?;
|
||||
let sorted_workers = sorted_workers(&workers);
|
||||
let workers = get_workers(&db, &sorted_workers, &base).await?;
|
||||
let tasks = get_queue(&db, &sorted_workers, &base).await?;
|
||||
Ok(QueueTemplate {
|
||||
base,
|
||||
inner: QueueInnerTemplate { runners, tasks },
|
||||
inner: QueueInnerTemplate { workers, tasks },
|
||||
})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -9,15 +9,15 @@ use axum::{
|
|||
|
||||
use crate::{
|
||||
config::Config,
|
||||
server::{runners::Runners, util},
|
||||
server::{util, workers::Workers},
|
||||
somehow,
|
||||
};
|
||||
|
||||
use super::{Base, Tab};
|
||||
|
||||
#[derive(Template)]
|
||||
#[template(path = "runner.html")]
|
||||
struct RunnerTemplate {
|
||||
#[template(path = "worker.html")]
|
||||
struct WorkerTemplate {
|
||||
base: Base,
|
||||
name: String,
|
||||
last_seen: String,
|
||||
|
|
@ -27,14 +27,14 @@ struct RunnerTemplate {
|
|||
pub async fn get(
|
||||
Path(name): Path<String>,
|
||||
State(config): State<&'static Config>,
|
||||
State(runners): State<Arc<Mutex<Runners>>>,
|
||||
State(workers): State<Arc<Mutex<Workers>>>,
|
||||
) -> somehow::Result<Response> {
|
||||
let info = runners.lock().unwrap().clean().get(&name);
|
||||
let info = workers.lock().unwrap().clean().get(&name);
|
||||
let Some(info) = info else {
|
||||
return Ok(StatusCode::NOT_FOUND.into_response());
|
||||
};
|
||||
|
||||
Ok(RunnerTemplate {
|
||||
Ok(WorkerTemplate {
|
||||
base: Base::new(config, Tab::None),
|
||||
name,
|
||||
last_seen: util::format_time(info.last_seen),
|
||||
|
|
@ -3,17 +3,17 @@ use std::collections::HashMap;
|
|||
use gix::hashtable::HashSet;
|
||||
use time::OffsetDateTime;
|
||||
|
||||
use crate::{config::Config, shared::RunnerStatus};
|
||||
use crate::{config::Config, shared::WorkerStatus};
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct RunnerInfo {
|
||||
pub struct WorkerInfo {
|
||||
pub secret: String,
|
||||
pub last_seen: OffsetDateTime,
|
||||
pub status: RunnerStatus,
|
||||
pub status: WorkerStatus,
|
||||
}
|
||||
|
||||
impl RunnerInfo {
|
||||
pub fn new(secret: String, last_seen: OffsetDateTime, status: RunnerStatus) -> Self {
|
||||
impl WorkerInfo {
|
||||
pub fn new(secret: String, last_seen: OffsetDateTime, status: WorkerStatus) -> Self {
|
||||
Self {
|
||||
secret,
|
||||
last_seen,
|
||||
|
|
@ -22,40 +22,40 @@ impl RunnerInfo {
|
|||
}
|
||||
}
|
||||
|
||||
pub struct Runners {
|
||||
pub struct Workers {
|
||||
config: &'static Config,
|
||||
runners: HashMap<String, RunnerInfo>,
|
||||
workers: HashMap<String, WorkerInfo>,
|
||||
}
|
||||
|
||||
impl Runners {
|
||||
impl Workers {
|
||||
pub fn new(config: &'static Config) -> Self {
|
||||
Self {
|
||||
config,
|
||||
runners: HashMap::new(),
|
||||
workers: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn clean(&mut self) -> &mut Self {
|
||||
let now = OffsetDateTime::now_utc();
|
||||
self.runners
|
||||
.retain(|_, v| now <= v.last_seen + self.config.web_runner_timeout);
|
||||
self.workers
|
||||
.retain(|_, v| now <= v.last_seen + self.config.web_worker_timeout);
|
||||
self
|
||||
}
|
||||
|
||||
pub fn verify(&self, name: &str, secret: &str) -> bool {
|
||||
let Some(runner) = self.runners.get(name) else { return true; };
|
||||
runner.secret == secret
|
||||
let Some(worker) = self.workers.get(name) else { return true; };
|
||||
worker.secret == secret
|
||||
}
|
||||
|
||||
pub fn update(&mut self, name: String, info: RunnerInfo) {
|
||||
self.runners.insert(name, info);
|
||||
pub fn update(&mut self, name: String, info: WorkerInfo) {
|
||||
self.workers.insert(name, info);
|
||||
}
|
||||
|
||||
fn oldest_working_on(&self, hash: &str) -> Option<&str> {
|
||||
self.runners
|
||||
self.workers
|
||||
.iter()
|
||||
.filter_map(|(name, info)| match &info.status {
|
||||
RunnerStatus::Working(run) if run.hash == hash => Some((name, run.start)),
|
||||
WorkerStatus::Working(run) if run.hash == hash => Some((name, run.start)),
|
||||
_ => None,
|
||||
})
|
||||
.max_by_key(|(_, since)| *since)
|
||||
|
|
@ -63,18 +63,18 @@ impl Runners {
|
|||
}
|
||||
|
||||
pub fn should_abort_work(&self, name: &str) -> bool {
|
||||
let Some(info) = self.runners.get(name) else { return false; };
|
||||
let RunnerStatus::Working ( run) = &info.status else { return false; };
|
||||
let Some(info) = self.workers.get(name) else { return false; };
|
||||
let WorkerStatus::Working ( run) = &info.status else { return false; };
|
||||
let Some(oldest) = self.oldest_working_on(&run.hash) else { return false; };
|
||||
name != oldest
|
||||
}
|
||||
|
||||
pub fn find_free_work<'a>(&self, hashes: &'a [String]) -> Option<&'a str> {
|
||||
let covered = self
|
||||
.runners
|
||||
.workers
|
||||
.values()
|
||||
.filter_map(|info| match &info.status {
|
||||
RunnerStatus::Working(run) => Some(&run.hash),
|
||||
WorkerStatus::Working(run) => Some(&run.hash),
|
||||
_ => None,
|
||||
})
|
||||
.collect::<HashSet<_>>();
|
||||
|
|
@ -85,11 +85,11 @@ impl Runners {
|
|||
.map(|hash| hash as &str)
|
||||
}
|
||||
|
||||
pub fn get(&self, name: &str) -> Option<RunnerInfo> {
|
||||
self.runners.get(name).cloned()
|
||||
pub fn get(&self, name: &str) -> Option<WorkerInfo> {
|
||||
self.workers.get(name).cloned()
|
||||
}
|
||||
|
||||
pub fn get_all(&self) -> HashMap<String, RunnerInfo> {
|
||||
self.runners.clone()
|
||||
pub fn get_all(&self) -> HashMap<String, WorkerInfo> {
|
||||
self.workers.clone()
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue