Skip to content

Commit 703f89e

Browse files
committed
Unified resource management
1 parent c907dc1 commit 703f89e

59 files changed

Lines changed: 3124 additions & 2146 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

crates/hyperqueue/src/bin/hq.rs

Lines changed: 15 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -40,9 +40,11 @@ use hyperqueue::server::bootstrap::get_client_session;
4040
use hyperqueue::transfer::messages::{
4141
FromClientMessage, IdSelector, JobInfoRequest, ToClientMessage,
4242
};
43-
use hyperqueue::worker::hwdetect::{detect_cpus, detect_cpus_no_ht, detect_generic_resources};
43+
use hyperqueue::worker::hwdetect::{
44+
detect_additional_resources, detect_cpus, prune_hyper_threading,
45+
};
4446
use hyperqueue::WorkerId;
45-
use tako::resources::ResourceDescriptor;
47+
use tako::resources::{ResourceDescriptor, ResourceDescriptorItem, CPU_RESOURCE_NAME};
4648

4749
#[cfg(feature = "jemalloc")]
4850
#[global_allocator]
@@ -188,7 +190,7 @@ struct WorkerWaitOpts {
188190
struct HwDetectOpts {
189191
/// Detect only physical cores
190192
#[clap(long)]
191-
no_hyperthreading: bool,
193+
no_hyper_threading: bool,
192194
}
193195

194196
// Job CLI options
@@ -406,15 +408,18 @@ async fn command_worker_wait(
406408
}
407409

408410
fn command_worker_hwdetect(gsettings: &GlobalSettings, opts: HwDetectOpts) -> anyhow::Result<()> {
409-
let cpus = if opts.no_hyperthreading {
410-
detect_cpus_no_ht()?
411-
} else {
412-
detect_cpus()?
413-
};
414-
let generic = detect_generic_resources()?;
411+
let mut cpus = detect_cpus()?;
412+
if opts.no_hyper_threading {
413+
cpus = prune_hyper_threading(&cpus)?;
414+
}
415+
let mut resources = vec![ResourceDescriptorItem {
416+
name: CPU_RESOURCE_NAME.to_string(),
417+
kind: cpus,
418+
}];
419+
detect_additional_resources(&mut resources)?;
415420
gsettings
416421
.printer()
417-
.print_hw(&ResourceDescriptor::new(cpus, generic));
422+
.print_hw(&ResourceDescriptor::new(resources));
418423
Ok(())
419424
}
420425

crates/hyperqueue/src/client/commands/autoalloc.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ use crate::transfer::connection::ClientSession;
1414
use crate::transfer::messages::{
1515
AllocationQueueParams, AutoAllocRequest, AutoAllocResponse, FromClientMessage, ToClientMessage,
1616
};
17-
use crate::worker::parser::{ArgCpuDefinition, ArgGenericResourceDef};
17+
use crate::worker::parser::{ArgCpuDefinition, ArgResourceItemDef};
1818

1919
#[derive(Parser)]
2020
pub struct AutoAllocOpts {
@@ -105,7 +105,7 @@ struct SharedQueueOpts {
105105

106106
/// What resources should the workers spawned inside allocations contain
107107
#[clap(long, multiple_occurrences(true))]
108-
resource: Vec<PassThroughArgument<ArgGenericResourceDef>>,
108+
resource: Vec<PassThroughArgument<ArgResourceItemDef>>,
109109

110110
/// Behavior when a connection to a server is lost
111111
#[clap(long, default_value = "finish-running", arg_enum)]

crates/hyperqueue/src/client/commands/submit/command.rs

Lines changed: 40 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -6,16 +6,16 @@ use std::{fs, io};
66
use anyhow::{anyhow, bail};
77
use bstr::BString;
88
use clap::Parser;
9-
use tako::gateway::{GenericResourceRequest, ResourceRequest};
9+
use tako::gateway::{ResourceRequest, ResourceRequestEntries, ResourceRequestEntry};
1010
use tako::program::{ProgramDefinition, StdioDef};
11-
use tako::resources::{CpuRequest, GenericResourceAmount, NumOfNodes};
11+
use tako::resources::{AllocationRequest, NumOfNodes, CPU_RESOURCE_NAME};
1212

1313
use super::directives::parse_hq_directives;
1414
use crate::client::commands::submit::directives::parse_hq_directives_from_file;
1515
use crate::client::commands::wait::{wait_for_jobs, wait_for_jobs_with_progress};
1616
use crate::client::globalsettings::GlobalSettings;
1717
use crate::client::job::get_worker_map;
18-
use crate::client::resources::{parse_cpu_request, parse_resource_request};
18+
use crate::client::resources::{parse_allocation_request, parse_resource_request};
1919
use crate::client::status::Status;
2020
use crate::common::arraydef::IntArray;
2121
use crate::common::placeholders::{
@@ -66,10 +66,10 @@ const DEFAULT_STDERR_PATH: &str = const_format::concatcp!(
6666
".stderr"
6767
);
6868

69-
crate::arg_wrapper!(ArgCpuRequest, CpuRequest, parse_cpu_request);
69+
crate::arg_wrapper!(ArgCpuRequest, AllocationRequest, parse_allocation_request);
7070
crate::arg_wrapper!(
7171
ArgNamedResourceRequest,
72-
(String, GenericResourceAmount),
72+
(String, AllocationRequest),
7373
parse_resource_request
7474
);
7575

@@ -275,7 +275,7 @@ impl SubmitJobConfOpts {
275275
}
276276
}
277277

278-
#[derive(clap::ArgEnum, Clone, PartialEq)]
278+
#[derive(clap::ArgEnum, Clone, Eq, PartialEq)]
279279
pub enum DirectivesMode {
280280
Auto,
281281
File,
@@ -328,36 +328,53 @@ pub struct JobSubmitOpts {
328328
}
329329

330330
impl JobSubmitOpts {
331-
fn resource_request(&self) -> ResourceRequest {
332-
let generic_resources = self
331+
fn resource_request(&self) -> anyhow::Result<ResourceRequest> {
332+
let mut resources: ResourceRequestEntries = self
333333
.conf
334334
.resource
335335
.iter()
336-
.map(|gr| {
337-
let rq = gr.get().clone();
338-
GenericResourceRequest {
339-
resource: rq.0,
340-
amount: rq.1,
336+
.map(|r| {
337+
let r = r.get().clone();
338+
ResourceRequestEntry {
339+
resource: r.0,
340+
policy: r.1,
341341
}
342342
})
343343
.collect();
344344

345-
ResourceRequest {
345+
let has_cpus = resources.iter().any(|r| r.resource == CPU_RESOURCE_NAME);
346+
347+
if let Some(cpus) = &self.conf.cpus {
348+
if has_cpus {
349+
anyhow::bail!("--cpus and --resource cpus=... cannot be combined");
350+
}
351+
resources.insert(
352+
0,
353+
ResourceRequestEntry {
354+
resource: CPU_RESOURCE_NAME.to_string(),
355+
policy: cpus.get().clone(),
356+
},
357+
)
358+
} else if !has_cpus {
359+
resources.insert(
360+
0,
361+
ResourceRequestEntry {
362+
resource: CPU_RESOURCE_NAME.to_string(),
363+
policy: AllocationRequest::Compact(1),
364+
},
365+
)
366+
}
367+
368+
Ok(ResourceRequest {
346369
n_nodes: self.conf.nodes.unwrap_or(0),
347-
cpus: self
348-
.conf
349-
.cpus
350-
.as_ref()
351-
.map(|c| c.get().clone())
352-
.unwrap_or_default(),
353370
min_time: self
354371
.conf
355372
.time_request
356373
.as_ref()
357374
.map(|t| *t.get())
358375
.unwrap_or_else(|| std::time::Duration::from_millis(0)),
359-
generic: generic_resources,
360-
}
376+
resources,
377+
})
361378
}
362379
}
363380

@@ -422,7 +439,7 @@ pub async fn submit_computation(
422439

423440
let opts = handle_directives(opts, stdin.as_deref())?;
424441

425-
let resources = opts.resource_request();
442+
let resources = opts.resource_request()?;
426443
let (ids, entries) = get_ids_and_entries(&opts)?;
427444
let task_count = ids.id_count();
428445

crates/hyperqueue/src/client/commands/submit/directives.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -82,7 +82,7 @@ fn extract_metadata(data: &BStr) -> crate::Result<FileMetadata> {
8282
));
8383
}
8484

85-
match line.get(0) {
85+
match line.first() {
8686
Some(b'#') | None => continue,
8787
_ => break,
8888
}

crates/hyperqueue/src/client/commands/worker.rs

Lines changed: 35 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,8 @@
1+
use anyhow::bail;
12
use std::path::PathBuf;
23
use std::time::Duration;
34

4-
use tako::resources::ResourceDescriptor;
5+
use tako::resources::{ResourceDescriptor, ResourceDescriptorItem, CPU_RESOURCE_NAME};
56
use tako::worker::{ServerLostPolicy, WorkerConfiguration};
67
use tako::Map;
78

@@ -22,8 +23,8 @@ use crate::transfer::messages::{
2223
use crate::worker::bootstrap::{
2324
finalize_configuration, initialize_worker, try_get_pbs_info, try_get_slurm_info,
2425
};
25-
use crate::worker::hwdetect::{detect_cpus, detect_cpus_no_ht, detect_generic_resources};
26-
use crate::worker::parser::{ArgCpuDefinition, ArgGenericResourceDef, CpuDefinition};
26+
use crate::worker::hwdetect::{detect_additional_resources, detect_cpus, prune_hyper_threading};
27+
use crate::worker::parser::{ArgCpuDefinition, ArgResourceItemDef};
2728
use crate::WorkerId;
2829

2930
#[derive(clap::ArgEnum, Clone)]
@@ -58,17 +59,21 @@ impl From<ArgServerLostPolicy> for ServerLostPolicy {
5859
#[derive(Parser)]
5960
pub struct WorkerStartOpts {
6061
/// How many cores should be allocated for the worker
61-
#[clap(long, default_value = "auto")]
62-
pub cpus: ArgCpuDefinition,
62+
#[clap(long)]
63+
pub cpus: Option<ArgCpuDefinition>,
6364

6465
/// Resources
6566
#[clap(long, multiple_occurrences(true))]
66-
pub resource: Vec<ArgGenericResourceDef>,
67+
pub resource: Vec<ArgResourceItemDef>,
6768

6869
#[clap(long = "no-detect-resources")]
6970
/// Disable auto-detection of resources
7071
pub no_detect_resources: bool,
7172

73+
#[clap(long = "no-hyper-threading")]
74+
/// Ignore hyper-threading while detecting CPU cores
75+
pub no_hyper_threading: bool,
76+
7277
/// How often should the worker announce its existence to the server. (default: "8s")
7378
#[clap(long, default_value = "8s")]
7479
pub heartbeat: ArgDuration,
@@ -121,24 +126,33 @@ fn gather_configuration(opts: WorkerStartOpts) -> anyhow::Result<WorkerConfigura
121126

122127
let hostname = get_hostname(opts.hostname);
123128

124-
let cpus = match opts.cpus.unpack() {
125-
CpuDefinition::Detect => detect_cpus()?,
126-
CpuDefinition::DetectNoHyperThreading => detect_cpus_no_ht()?,
127-
CpuDefinition::Custom(cpus) => cpus,
128-
};
129+
let mut resources: Vec<_> = opts.resource.into_iter().map(|x| x.unpack()).collect();
130+
if !resources.iter().any(|x| x.name == CPU_RESOURCE_NAME) {
131+
resources.push(ResourceDescriptorItem {
132+
name: CPU_RESOURCE_NAME.to_string(),
133+
kind: if let Some(cpus) = opts.cpus {
134+
cpus.unpack()
135+
} else {
136+
detect_cpus()?
137+
},
138+
})
139+
} else if opts.cpus.is_some() {
140+
bail!("Parameters --cpus and --resource cpus=... cannot be combined");
141+
}
129142

130-
let mut generic = if opts.no_detect_resources {
131-
Vec::new()
132-
} else {
133-
detect_generic_resources()?
134-
};
135-
for def in opts.resource {
136-
let descriptor = def.unpack();
137-
generic.retain(|desc| desc.name != descriptor.name);
138-
generic.push(descriptor)
143+
if opts.no_hyper_threading {
144+
let cpus = resources
145+
.iter_mut()
146+
.find(|x| x.name == CPU_RESOURCE_NAME)
147+
.unwrap();
148+
cpus.kind = prune_hyper_threading(&cpus.kind)?;
149+
}
150+
151+
if !opts.no_detect_resources {
152+
detect_additional_resources(&mut resources)?;
139153
}
140154

141-
let resources = ResourceDescriptor::new(cpus, generic);
155+
let resources = ResourceDescriptor::new(resources);
142156
resources.validate()?;
143157

144158
let (work_dir, log_dir) = {

0 commit comments

Comments
 (0)