MCPcopy Create free account
hub / github.com/It4innovations/hyperqueue / start_server

Function start_server

crates/hyperqueue/src/server/bootstrap.rs:327–400  ·  view source on GitHub ↗
(
    gsettings: &GlobalSettings,
    mut server_cfg: ServerConfig,
)

Source from the content-addressed store, hash-verified

325}
326
327async fn start_server(
328 gsettings: &GlobalSettings,
329 mut server_cfg: ServerConfig,
330) -> anyhow::Result<()> {
331 let restorer = if server_cfg.journal_path.as_ref().is_some_and(|p| p.exists()) {
332 log::info!("Loading journal ...");
333 let mut restorer = StateRestorer::default();
334 restorer.load_event_file(server_cfg.journal_path.as_ref().unwrap())?;
335 if restorer.truncate_size().is_some() {
336 log::warn!(
337 "Journal contains not fully written data; they will be removed from the log"
338 );
339 }
340 let server_uid = restorer.take_server_uid();
341 if !server_uid.is_empty() {
342 server_cfg.server_uid = Some(server_uid)
343 }
344 Some(restorer)
345 } else {
346 None
347 };
348
349 let (fut, _, state_ref, senders) = initialize_server(
350 gsettings,
351 server_cfg,
352 restorer
353 .as_ref()
354 .map(|r| r.worker_id_counter())
355 .unwrap_or(WorkerId::new(0)),
356 restorer.as_ref().map(|r| r.queue_id_counter()).unwrap_or(1),
357 restorer.as_ref().and_then(|r| r.truncate_size()),
358 )
359 .await?;
360 let new_tasks_and_queues = if let Some(restorer) = restorer {
361 let mut state = state_ref.get_mut();
362 let ra = &senders.server_control;
363 state.restore_state(&restorer);
364 Some(restorer.restore_jobs_and_queues(&mut state, ra)?)
365 } else {
366 None
367 };
368 let local_set = LocalSet::new();
369
370 if let Some((new_tasks, new_queues)) = new_tasks_and_queues {
371 let server_dir = gsettings.server_directory().to_path_buf();
372 local_set.spawn_local(async move {
373 log::debug!("Restoring old tasks into Tako");
374 for new in new_tasks {
375 senders.server_control.add_new_tasks(new).unwrap();
376 }
377 log::debug!("Restoration of old tasks is completed");
378 for queue in new_queues {
379 senders
380 .autoalloc
381 .add_queue(
382 &server_dir,
383 *queue.params,
384 Some(queue.queue_id),

Callers 1

init_hq_serverFunction · 0.70

Calls 15

initialize_serverFunction · 0.85
as_refMethod · 0.80
load_event_fileMethod · 0.80
truncate_sizeMethod · 0.80
take_server_uidMethod · 0.80
worker_id_counterMethod · 0.80
queue_id_counterMethod · 0.80
restore_stateMethod · 0.80
server_directoryMethod · 0.80
add_new_tasksMethod · 0.80
run_untilMethod · 0.80

Tested by

no test coverage detected