(
gsettings: &GlobalSettings,
mut server_cfg: ServerConfig,
)
| 325 | } |
| 326 | |
| 327 | async 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), |
no test coverage detected