(
autoalloc: &mut AutoAllocState,
events: &EventStreamer,
server_directory: PathBuf,
params: QueueParameters,
queue_id: Option<QueueId>,
worker_resources: Option<ResourceDescri
| 675 | } |
| 676 | |
| 677 | fn create_queue( |
| 678 | autoalloc: &mut AutoAllocState, |
| 679 | events: &EventStreamer, |
| 680 | server_directory: PathBuf, |
| 681 | params: QueueParameters, |
| 682 | queue_id: Option<QueueId>, |
| 683 | worker_resources: Option<ResourceDescriptor>, |
| 684 | ) -> anyhow::Result<QueueId> { |
| 685 | let name = params.name.clone(); |
| 686 | let handler = create_allocation_handler(¶ms.manager, name.clone(), server_directory); |
| 687 | let queue_info = QueueInfo::new(params.clone()); |
| 688 | |
| 689 | match handler { |
| 690 | Ok(handler) => { |
| 691 | let queue = AllocationQueue::new( |
| 692 | queue_info, |
| 693 | name, |
| 694 | handler, |
| 695 | create_rate_limiter(), |
| 696 | worker_resources, |
| 697 | ); |
| 698 | let id = { |
| 699 | let id = autoalloc.add_queue(queue, queue_id); |
| 700 | if queue_id.is_none() { |
| 701 | // When queue_id is provided, then we are restoring journal, |
| 702 | // so we do not want to double log the event |
| 703 | events.on_allocation_queue_created(id, params); |
| 704 | } |
| 705 | id |
| 706 | }; |
| 707 | |
| 708 | Ok(id) |
| 709 | } |
| 710 | Err(error) => Err(anyhow::anyhow!("Could not create autoalloc queue: {error}")), |
| 711 | } |
| 712 | } |
| 713 | |
| 714 | // TODO: use proper error type |
| 715 | async fn remove_queue( |
no test coverage detected