Skip to main content

hermit/scheduler/
mod.rs

1#![allow(clippy::type_complexity)]
2
3use alloc::boxed::Box;
4use alloc::collections::{BTreeMap, VecDeque};
5use alloc::rc::Rc;
6use alloc::sync::Arc;
7#[cfg(feature = "smp")]
8use alloc::vec::Vec;
9use core::cell::RefCell;
10use core::ptr;
11use core::sync::atomic::{AtomicI32, AtomicU32, Ordering};
12
13use ahash::RandomState;
14use crossbeam_utils::Backoff;
15use hashbrown::{HashMap, hash_map};
16use hermit_sync::*;
17#[cfg(target_arch = "riscv64")]
18use riscv::register::sstatus;
19use timer_interrupts::TimerList;
20
21use crate::arch::kernel;
22use crate::arch::kernel::core_local::*;
23#[cfg(feature = "preemptive")]
24use crate::arch::kernel::processor;
25use crate::arch::kernel::scheduler::TaskStacks;
26#[cfg(target_arch = "riscv64")]
27use crate::arch::kernel::switch::switch_to_task;
28#[cfg(target_arch = "x86_64")]
29use crate::arch::kernel::switch::{switch_to_fpu_owner, switch_to_task};
30use crate::arch::kernel::{get_processor_count, interrupts};
31use crate::errno::Errno;
32use crate::fd::{Fd, RawFd};
33use crate::io;
34use crate::scheduler::task::*;
35
36#[cfg(all(
37	any(target_arch = "x86_64", target_arch = "riscv64"),
38	feature = "smp",
39	not(feature = "idle-poll")
40))]
41pub mod sleep_state;
42pub mod task;
43pub mod timer_interrupts;
44
45static NO_TASKS: AtomicU32 = AtomicU32::new(0);
46/// Length of a preemptive scheduling time slice, in microseconds (10 ms). With
47/// rhyve's external-interrupt exiting this also bounds the guest VM-exit rate
48/// during contention (one exit per slice).
49#[cfg(feature = "preemptive")]
50const PREEMPTION_SLICE_US: u64 = 10_000;
51/// Map between Core ID and per-core scheduler
52#[cfg(feature = "smp")]
53static SCHEDULER_INPUTS: SpinMutex<Vec<&InterruptTicketMutex<SchedulerInput>>> =
54	SpinMutex::new(Vec::new());
55/// Map between Task ID and Queue of waiting tasks
56static WAITING_TASKS: InterruptTicketMutex<BTreeMap<TaskId, VecDeque<TaskHandle>>> =
57	InterruptTicketMutex::new(BTreeMap::new());
58/// Map between Task ID and TaskHandle
59static TASKS: InterruptTicketMutex<BTreeMap<TaskId, TaskHandle>> =
60	InterruptTicketMutex::new(BTreeMap::new());
61
62/// Unique identifier for a core.
63pub type CoreId = u32;
64
65#[cfg(feature = "smp")]
66pub(crate) struct SchedulerInput {
67	/// Queue of new tasks
68	new_tasks: VecDeque<NewTask>,
69	/// Queue of task, which are wakeup by another core
70	wakeup_tasks: VecDeque<TaskHandle>,
71}
72
73#[cfg(feature = "smp")]
74impl SchedulerInput {
75	pub fn new() -> Self {
76		Self {
77			new_tasks: VecDeque::new(),
78			wakeup_tasks: VecDeque::new(),
79		}
80	}
81}
82
83#[cfg_attr(any(target_arch = "x86_64", target_arch = "aarch64"), repr(align(128)))]
84#[cfg_attr(
85	not(any(target_arch = "x86_64", target_arch = "aarch64")),
86	repr(align(64))
87)]
88pub(crate) struct PerCoreScheduler {
89	/// Core ID of this per-core scheduler
90	#[cfg(feature = "smp")]
91	core_id: CoreId,
92	/// Task which is currently running
93	current_task: Rc<RefCell<Task>>,
94	/// Idle Task
95	idle_task: Rc<RefCell<Task>>,
96	/// Task that currently owns the FPU
97	#[cfg(any(target_arch = "x86_64", target_arch = "aarch64"))]
98	fpu_owner: Rc<RefCell<Task>>,
99	/// Queue of tasks, which are ready
100	ready_queue: PriorityTaskQueue,
101	/// Queue of tasks, which are finished and can be released
102	finished_tasks: VecDeque<Rc<RefCell<Task>>>,
103	/// Queue of blocked tasks, sorted by wakeup time.
104	blocked_tasks: BlockedTaskQueue,
105	/// Queue of timer interrupts.
106	pub timers: TimerList,
107}
108
109pub(crate) trait PerCoreSchedulerExt {
110	/// Triggers the scheduler to reschedule the tasks.
111	/// Interrupt flag will be cleared during the reschedule
112	fn reschedule(self);
113
114	/// Terminate the current task on the current core.
115	fn exit(self, exit_code: i32) -> !;
116}
117
118impl PerCoreSchedulerExt for &mut PerCoreScheduler {
119	#[cfg(target_arch = "x86_64")]
120	fn reschedule(self) {
121		without_interrupts(|| {
122			let Some(last_stack_pointer) = self.scheduler() else {
123				return;
124			};
125
126			let (new_stack_pointer, is_idle) = {
127				let borrowed = self.current_task.borrow();
128				(
129					borrowed.last_stack_pointer,
130					borrowed.status == TaskStatus::Idle,
131				)
132			};
133
134			if is_idle || Rc::ptr_eq(&self.current_task, &self.fpu_owner) {
135				unsafe {
136					switch_to_fpu_owner(last_stack_pointer, new_stack_pointer.as_u64() as usize);
137				}
138			} else {
139				unsafe {
140					switch_to_task(last_stack_pointer, new_stack_pointer.as_u64() as usize);
141				}
142			}
143		});
144	}
145
146	/// Trigger an interrupt to reschedule the system
147	#[cfg(target_arch = "aarch64")]
148	fn reschedule(self) {
149		use aarch64_cpu::asm::barrier::{NSH, SY, dsb, isb};
150		use arm_gic::IntId;
151		use arm_gic::gicv3::{GicCpuInterface, SgiTarget, SgiTargetGroup};
152
153		use crate::arch::kernel::interrupts::SGI_RESCHED;
154
155		dsb(NSH);
156		isb(SY);
157
158		let reschedid = IntId::sgi(SGI_RESCHED.into());
159		#[cfg(feature = "smp")]
160		let core_id = self.core_id;
161		#[cfg(not(feature = "smp"))]
162		let core_id = 0;
163
164		GicCpuInterface::send_sgi(
165			reschedid,
166			SgiTarget::List {
167				affinity3: 0,
168				affinity2: 0,
169				affinity1: 0,
170				target_list: 1 << core_id,
171			},
172			SgiTargetGroup::CurrentGroup1,
173		)
174		.unwrap();
175
176		interrupts::enable();
177	}
178
179	#[cfg(target_arch = "riscv64")]
180	fn reschedule(self) {
181		without_interrupts(|| self.scheduler());
182	}
183
184	fn exit(self, exit_code: i32) -> ! {
185		without_interrupts(|| {
186			// Get the current task.
187			let mut current_task_borrowed = self.current_task.borrow_mut();
188			assert_ne!(
189				current_task_borrowed.status,
190				TaskStatus::Idle,
191				"Trying to terminate the idle task"
192			);
193
194			// Finish the task and reschedule.
195			debug!(
196				"Finishing task {} with exit code {}",
197				current_task_borrowed.id, exit_code
198			);
199			current_task_borrowed.status = TaskStatus::Finished;
200			NO_TASKS.fetch_sub(1, Ordering::SeqCst);
201
202			let current_id = current_task_borrowed.id;
203			drop(current_task_borrowed);
204
205			// wakeup tasks, which are waiting for task with the identifier id
206			if let Some(mut queue) = WAITING_TASKS.lock().remove(&current_id) {
207				while let Some(task) = queue.pop_front() {
208					self.custom_wakeup(task);
209				}
210			}
211
212			TASKS.lock().remove(&current_id);
213		});
214
215		self.reschedule();
216		unreachable!()
217	}
218}
219
220struct NewTask {
221	tid: TaskId,
222	func: unsafe extern "C" fn(usize),
223	arg: usize,
224	prio: Priority,
225	core_id: CoreId,
226	stacks: TaskStacks,
227	object_map: Arc<RwSpinLock<HashMap<RawFd, Arc<async_lock::RwLock<Fd>>, RandomState>>>,
228}
229
230impl From<NewTask> for Task {
231	fn from(value: NewTask) -> Self {
232		let NewTask {
233			tid,
234			func,
235			arg,
236			prio,
237			core_id,
238			stacks,
239			object_map,
240		} = value;
241		let mut task = Self::new(tid, core_id, TaskStatus::Ready, prio, stacks, object_map);
242		task.create_stack_frame(func, arg);
243		task
244	}
245}
246
247impl PerCoreScheduler {
248	/// Spawn a new task.
249	pub unsafe fn spawn(
250		func: unsafe extern "C" fn(usize),
251		arg: usize,
252		prio: Priority,
253		core_id: CoreId,
254		stack_size: usize,
255	) -> TaskId {
256		// Create the new task.
257		let tid = get_tid();
258		let stacks = TaskStacks::new(stack_size);
259		let new_task = NewTask {
260			tid,
261			func,
262			arg,
263			prio,
264			core_id,
265			stacks,
266			object_map: core_scheduler().get_current_task_object_map(),
267		};
268
269		// Add it to the task lists.
270		let wakeup = {
271			#[cfg(feature = "smp")]
272			let mut input_locked = get_scheduler_input(core_id).lock();
273			WAITING_TASKS.lock().insert(tid, VecDeque::with_capacity(1));
274			TASKS.lock().insert(
275				tid,
276				TaskHandle::new(
277					tid,
278					prio,
279					#[cfg(feature = "smp")]
280					core_id,
281				),
282			);
283			NO_TASKS.fetch_add(1, Ordering::SeqCst);
284
285			#[cfg(feature = "smp")]
286			if core_id == core_scheduler().core_id {
287				let task = Rc::new(RefCell::new(Task::from(new_task)));
288				core_scheduler().ready_queue.push(task);
289				false
290			} else {
291				input_locked.new_tasks.push_back(new_task);
292				true
293			}
294			#[cfg(not(feature = "smp"))]
295			if core_id == 0 {
296				let task = Rc::new(RefCell::new(Task::from(new_task)));
297				core_scheduler().ready_queue.push(task);
298				false
299			} else {
300				panic!("Invalid core_id {core_id}!")
301			}
302		};
303
304		debug!("Creating task {tid} with priority {prio} on core {core_id}");
305
306		if wakeup {
307			kernel::wakeup_core(core_id);
308		}
309
310		tid
311	}
312
313	#[cfg(feature = "newlib")]
314	fn clone_impl(&self, func: extern "C" fn(usize), arg: usize) -> TaskId {
315		static NEXT_CORE_ID: AtomicU32 = AtomicU32::new(1);
316
317		// Get the Core ID of the next CPU.
318		let core_id: CoreId = {
319			// Increase the CPU number by 1.
320			let id = NEXT_CORE_ID.fetch_add(1, Ordering::SeqCst);
321
322			// Check for overflow.
323			if id == get_processor_count() {
324				NEXT_CORE_ID.store(0, Ordering::SeqCst);
325				0
326			} else {
327				id
328			}
329		};
330
331		// Get the current task.
332		let current_task_borrowed = self.current_task.borrow();
333
334		// Clone the current task.
335		let tid = get_tid();
336		let clone_task = NewTask {
337			tid,
338			func,
339			arg,
340			prio: current_task_borrowed.prio,
341			core_id,
342			stacks: TaskStacks::new(current_task_borrowed.stacks.get_user_stack_size()),
343			object_map: current_task_borrowed.object_map.clone(),
344		};
345
346		// Add it to the task lists.
347		let wakeup = {
348			#[cfg(feature = "smp")]
349			let mut input_locked = get_scheduler_input(core_id).lock();
350			WAITING_TASKS.lock().insert(tid, VecDeque::with_capacity(1));
351			TASKS.lock().insert(
352				tid,
353				TaskHandle::new(
354					tid,
355					current_task_borrowed.prio,
356					#[cfg(feature = "smp")]
357					core_id,
358				),
359			);
360			NO_TASKS.fetch_add(1, Ordering::SeqCst);
361			#[cfg(feature = "smp")]
362			if core_id == core_scheduler().core_id {
363				let clone_task = Rc::new(RefCell::new(Task::from(clone_task)));
364				core_scheduler().ready_queue.push(clone_task);
365				false
366			} else {
367				input_locked.new_tasks.push_back(clone_task);
368				true
369			}
370			#[cfg(not(feature = "smp"))]
371			if core_id == 0 {
372				let clone_task = Rc::new(RefCell::new(Task::from(clone_task)));
373				core_scheduler().ready_queue.push(clone_task);
374				false
375			} else {
376				panic!("Invalid core_id {core_id}!");
377			}
378		};
379
380		// Wake up the CPU
381		if wakeup {
382			kernel::wakeup_core(core_id);
383		}
384
385		tid
386	}
387
388	#[cfg(feature = "newlib")]
389	pub fn clone(&self, func: extern "C" fn(usize), arg: usize) -> TaskId {
390		without_interrupts(|| self.clone_impl(func, arg))
391	}
392
393	/// Returns `true` if a reschedule is required
394	#[inline]
395	#[cfg(all(any(target_arch = "x86_64", target_arch = "riscv64"), feature = "smp"))]
396	pub fn is_scheduling(&self) -> bool {
397		self.current_task.borrow().prio < self.ready_queue.get_highest_priority()
398	}
399
400	#[inline]
401	pub fn handle_waiting_tasks(&mut self) {
402		without_interrupts(|| {
403			crate::executor::run();
404			self.blocked_tasks
405				.handle_waiting_tasks(&mut self.ready_queue);
406		});
407	}
408
409	#[cfg(not(feature = "smp"))]
410	pub fn custom_wakeup(&mut self, task: TaskHandle) {
411		without_interrupts(|| {
412			let task = self.blocked_tasks.custom_wakeup(task);
413			self.ready_queue.push(task);
414			// The woken task is now ready. Arm the preemption timer so the running
415			// task — which may be CPU-bound and never yield — is preempted soon and
416			// the newcomer gets to run.
417			#[cfg(feature = "preemptive")]
418			{
419				let deadline = processor::get_timer_ticks() + PREEMPTION_SLICE_US;
420				timer_interrupts::create_timer_abs(timer_interrupts::Source::Preemption, deadline);
421			}
422		});
423	}
424
425	#[cfg(feature = "smp")]
426	pub fn custom_wakeup(&mut self, task: TaskHandle) {
427		if task.get_core_id() == self.core_id {
428			without_interrupts(|| {
429				let task = self.blocked_tasks.custom_wakeup(task);
430				self.ready_queue.push(task);
431				// The woken task is now ready on this core. Arm the preemption timer
432				// so the running (possibly non-yielding) task is preempted soon.
433				#[cfg(feature = "preemptive")]
434				{
435					let deadline = processor::get_timer_ticks() + PREEMPTION_SLICE_US;
436					timer_interrupts::create_timer_abs(
437						timer_interrupts::Source::Preemption,
438						deadline,
439					);
440				}
441			});
442		} else {
443			get_scheduler_input(task.get_core_id())
444				.lock()
445				.wakeup_tasks
446				.push_back(task);
447			// Wake up the CPU
448			kernel::wakeup_core(task.get_core_id());
449		}
450	}
451
452	#[inline]
453	pub fn block_current_task(&mut self, wakeup_time: Option<u64>) {
454		without_interrupts(|| {
455			self.blocked_tasks
456				.add(self.current_task.clone(), wakeup_time);
457		});
458	}
459
460	#[inline]
461	pub fn get_current_task_handle(&self) -> TaskHandle {
462		without_interrupts(|| {
463			let current_task_borrowed = self.current_task.borrow();
464
465			TaskHandle::new(
466				current_task_borrowed.id,
467				current_task_borrowed.prio,
468				#[cfg(feature = "smp")]
469				current_task_borrowed.core_id,
470			)
471		})
472	}
473
474	#[inline]
475	pub fn get_current_task_id(&self) -> TaskId {
476		without_interrupts(|| self.current_task.borrow().id)
477	}
478
479	#[inline]
480	pub fn get_current_task_object_map(
481		&self,
482	) -> Arc<RwSpinLock<HashMap<RawFd, Arc<async_lock::RwLock<Fd>>, RandomState>>> {
483		without_interrupts(|| self.current_task.borrow().object_map.clone())
484	}
485
486	/// Map a file descriptor to their IO interface and returns
487	/// the shared reference
488	#[inline]
489	pub fn get_object(&self, fd: RawFd) -> io::Result<Arc<async_lock::RwLock<Fd>>> {
490		without_interrupts(|| {
491			let current_task = self.current_task.borrow();
492			let object_map = current_task.object_map.read();
493			object_map.get(&fd).cloned().ok_or(Errno::Badf)
494		})
495	}
496
497	/// Creates a new map between file descriptor and their IO interface and
498	/// clone the standard descriptors.
499	#[cfg(feature = "common-os")]
500	#[cfg_attr(not(target_arch = "x86_64"), expect(dead_code))]
501	pub fn recreate_objmap(&self) -> io::Result<()> {
502		let mut map = HashMap::<RawFd, Arc<async_lock::RwLock<Fd>>, RandomState>::with_hasher(
503			RandomState::with_seeds(0, 0, 0, 0),
504		);
505
506		without_interrupts(|| {
507			let mut current_task = self.current_task.borrow_mut();
508			let object_map = current_task.object_map.read();
509
510			// clone standard file descriptors
511			for i in 0..3 {
512				if let Some(obj) = object_map.get(&i) {
513					map.insert(i, obj.clone());
514				}
515			}
516
517			drop(object_map);
518			current_task.object_map = Arc::new(RwSpinLock::new(map));
519		});
520
521		Ok(())
522	}
523
524	/// Insert a new IO interface and returns a file descriptor as
525	/// identifier to this object
526	pub fn insert_object(&self, obj: Arc<async_lock::RwLock<Fd>>) -> io::Result<RawFd> {
527		without_interrupts(|| {
528			let current_task = self.current_task.borrow();
529			let mut object_map = current_task.object_map.write();
530
531			let new_fd = || -> io::Result<RawFd> {
532				let mut fd: RawFd = 0;
533				loop {
534					if !object_map.contains_key(&fd) {
535						break Ok(fd);
536					} else if fd == RawFd::MAX {
537						break Err(Errno::Overflow);
538					}
539
540					fd = fd.saturating_add(1);
541				}
542			};
543
544			let fd = new_fd()?;
545			object_map.insert(fd, obj.clone());
546			Ok(fd)
547		})
548	}
549
550	/// Duplicate a IO interface and returns a new file descriptor as
551	/// identifier to the new copy
552	pub fn dup_object(&self, fd: RawFd) -> io::Result<RawFd> {
553		without_interrupts(|| {
554			let current_task = self.current_task.borrow();
555			let mut object_map = current_task.object_map.write();
556
557			let obj = (*(object_map.get(&fd).ok_or(Errno::Inval)?)).clone();
558
559			let new_fd = || -> io::Result<RawFd> {
560				let mut fd: RawFd = 0;
561				loop {
562					if !object_map.contains_key(&fd) {
563						break Ok(fd);
564					} else if fd == RawFd::MAX {
565						break Err(Errno::Overflow);
566					}
567
568					fd = fd.saturating_add(1);
569				}
570			};
571
572			let fd = new_fd()?;
573			match object_map.entry(fd) {
574				hash_map::Entry::Occupied(_occupied_entry) => Err(Errno::Mfile),
575				hash_map::Entry::Vacant(vacant_entry) => {
576					vacant_entry.insert(obj);
577					Ok(fd)
578				}
579			}
580		})
581	}
582
583	pub fn dup_object2(&self, fd1: RawFd, fd2: RawFd) -> io::Result<RawFd> {
584		without_interrupts(|| {
585			let current_task = self.current_task.borrow();
586			let mut object_map = current_task.object_map.write();
587
588			let obj = object_map.get(&fd1).cloned().ok_or(Errno::Badf)?;
589
590			match object_map.entry(fd2) {
591				hash_map::Entry::Occupied(_occupied_entry) => Err(Errno::Mfile),
592				hash_map::Entry::Vacant(vacant_entry) => {
593					vacant_entry.insert(obj);
594					Ok(fd2)
595				}
596			}
597		})
598	}
599
600	/// Remove a IO interface, which is named by the file descriptor
601	pub fn remove_object(&self, fd: RawFd) -> io::Result<Arc<async_lock::RwLock<Fd>>> {
602		without_interrupts(|| {
603			let current_task = self.current_task.borrow();
604			let mut object_map = current_task.object_map.write();
605
606			object_map.remove(&fd).ok_or(Errno::Badf)
607		})
608	}
609
610	#[inline]
611	pub fn get_current_task_prio(&self) -> Priority {
612		without_interrupts(|| self.current_task.borrow().prio)
613	}
614
615	/// Returns reference to prio_bitmap
616	#[allow(dead_code)]
617	#[inline]
618	pub fn get_priority_bitmap(&self) -> &u64 {
619		self.ready_queue.get_priority_bitmap()
620	}
621
622	#[cfg(target_arch = "x86_64")]
623	pub fn set_current_kernel_stack(&self) {
624		let current_task_borrowed = self.current_task.borrow();
625		let tss = unsafe { &mut *CoreLocal::get().tss.get() };
626
627		let rsp = current_task_borrowed.stacks.get_kernel_stack()
628			+ current_task_borrowed.stacks.get_kernel_stack_size() as u64
629			- TaskStacks::MARKER_SIZE as u64;
630		tss.privilege_stack_table[0] = rsp.into();
631		CoreLocal::get().kernel_stack.set(rsp.as_mut_ptr());
632		let ist_start = current_task_borrowed.stacks.get_interrupt_stack()
633			+ current_task_borrowed.stacks.get_interrupt_stack_size() as u64
634			- TaskStacks::MARKER_SIZE as u64;
635		tss.interrupt_stack_table[0] = ist_start.into();
636	}
637
638	pub fn set_current_task_priority(&mut self, prio: Priority) {
639		without_interrupts(|| {
640			trace!("Change priority of the current task");
641			self.current_task.borrow_mut().prio = prio;
642		});
643	}
644
645	pub fn set_priority(&mut self, id: TaskId, prio: Priority) -> Result<(), ()> {
646		trace!("Change priority of task {id} to priority {prio}");
647
648		without_interrupts(|| {
649			let task = get_task_handle(id).ok_or(())?;
650			#[cfg(feature = "smp")]
651			let other_core = task.get_core_id() != self.core_id;
652			#[cfg(not(feature = "smp"))]
653			let other_core = false;
654
655			if other_core {
656				warn!("Have to change the priority on another core");
657			} else if self.current_task.borrow().id == task.get_id() {
658				self.current_task.borrow_mut().prio = prio;
659			} else {
660				self.ready_queue
661					.set_priority(task, prio)
662					.expect("Do not find valid task in ready queue");
663			}
664
665			Ok(())
666		})
667	}
668
669	#[cfg(target_arch = "riscv64")]
670	pub fn set_current_kernel_stack(&self) {
671		let current_task_borrowed = self.current_task.borrow();
672
673		let stack = (current_task_borrowed.stacks.get_kernel_stack()
674			+ current_task_borrowed.stacks.get_kernel_stack_size() as u64
675			- TaskStacks::MARKER_SIZE as u64)
676			.as_u64();
677		CoreLocal::get().kernel_stack.set(stack);
678	}
679
680	/// Save the FPU context for the current FPU owner and restore it for the current task,
681	/// which wants to use the FPU now.
682	#[cfg(any(target_arch = "x86_64", target_arch = "aarch64"))]
683	pub fn fpu_switch(&mut self) {
684		if !Rc::ptr_eq(&self.current_task, &self.fpu_owner) {
685			debug!(
686				"Switching FPU owner from task {} to {}",
687				self.fpu_owner.borrow().id,
688				self.current_task.borrow().id
689			);
690
691			self.fpu_owner.borrow_mut().last_fpu_state.save();
692			self.current_task.borrow().last_fpu_state.restore();
693			self.fpu_owner = self.current_task.clone();
694		}
695	}
696
697	/// Check if a finished task could be deleted.
698	fn cleanup_tasks(&mut self) {
699		// Pop the first finished task and remove it from the TASKS list, which implicitly deallocates all associated memory.
700		while let Some(finished_task) = self.finished_tasks.pop_front() {
701			debug!("Cleaning up task {}", finished_task.borrow().id);
702		}
703	}
704
705	#[cfg(feature = "smp")]
706	pub fn check_input(&mut self) {
707		let mut input_locked = CoreLocal::get().scheduler_input.lock();
708
709		while let Some(task) = input_locked.wakeup_tasks.pop_front() {
710			let task = self.blocked_tasks.custom_wakeup(task);
711			self.ready_queue.push(task);
712		}
713
714		while let Some(new_task) = input_locked.new_tasks.pop_front() {
715			let task = Rc::new(RefCell::new(Task::from(new_task)));
716			self.ready_queue.push(task.clone());
717		}
718	}
719
720	/// Only the idle task should call this function.
721	/// Set the idle task to halt state if not another
722	/// available.
723	pub fn run() -> ! {
724		let backoff = Backoff::new();
725
726		loop {
727			let core_scheduler = core_scheduler();
728			interrupts::disable();
729
730			// run async tasks
731			crate::executor::run();
732
733			// do housekeeping
734			#[cfg(feature = "smp")]
735			core_scheduler.check_input();
736			core_scheduler.cleanup_tasks();
737
738			if core_scheduler.ready_queue.is_empty() {
739				if backoff.is_completed() {
740					interrupts::enable_and_wait();
741					backoff.reset();
742				} else {
743					interrupts::enable();
744					backoff.snooze();
745				}
746			} else {
747				interrupts::enable();
748				core_scheduler.reschedule();
749				backoff.reset();
750			}
751		}
752	}
753
754	#[inline]
755	#[cfg(target_arch = "aarch64")]
756	pub fn get_last_stack_pointer(&self) -> memory_addresses::VirtAddr {
757		self.current_task.borrow().last_stack_pointer
758	}
759
760	/// Triggers the scheduler to reschedule the tasks.
761	/// Interrupt flag must be cleared before calling this function.
762	pub fn scheduler(&mut self) -> Option<*mut usize> {
763		// run background tasks
764		crate::executor::run();
765
766		// Someone wants to give up the CPU
767		// => we have time to cleanup the system
768		self.cleanup_tasks();
769
770		// Get information about the current task.
771		let (id, last_stack_pointer, prio, status) = {
772			let mut borrowed = self.current_task.borrow_mut();
773			(
774				borrowed.id,
775				ptr::from_mut(&mut borrowed.last_stack_pointer).cast::<usize>(),
776				borrowed.prio,
777				borrowed.status,
778			)
779		};
780
781		let mut new_task = None;
782
783		if status == TaskStatus::Running {
784			// A task is currently running.
785			// Check if a task with a equal or higher priority is available.
786			if let Some(task) = self.ready_queue.pop_with_prio(prio) {
787				new_task = Some(task);
788			}
789		} else {
790			if status == TaskStatus::Finished {
791				// Mark the finished task as invalid and add it to the finished tasks for a later cleanup.
792				self.current_task.borrow_mut().status = TaskStatus::Invalid;
793				self.finished_tasks.push_back(self.current_task.clone());
794			}
795
796			// No task is currently running.
797			// Check if there is any available task and get the one with the highest priority.
798			if let Some(task) = self.ready_queue.pop() {
799				// This available task becomes the new task.
800				debug!("Task is available.");
801				new_task = Some(task);
802			} else if status != TaskStatus::Idle {
803				// The Idle task becomes the new task.
804				debug!("Only Idle Task is available.");
805				new_task = Some(self.idle_task.clone());
806			}
807		}
808
809		let task = new_task?;
810		// There is a new task we want to switch to.
811
812		// Handle the current task.
813		if status == TaskStatus::Running {
814			// Mark the running task as ready again and add it back to the queue.
815			self.current_task.borrow_mut().status = TaskStatus::Ready;
816			self.ready_queue.push(self.current_task.clone());
817		}
818
819		// Handle the new task and get information about it.
820		let (new_id, new_stack_pointer) = {
821			let mut borrowed = task.borrow_mut();
822			if borrowed.status != TaskStatus::Idle {
823				// Mark the new task as running.
824				borrowed.status = TaskStatus::Running;
825			}
826
827			(borrowed.id, borrowed.last_stack_pointer)
828		};
829
830		if id == new_id {
831			return None;
832		}
833
834		// Preemptive round-robin: arm a one-shot time-slice timer for the task we
835		// are about to run, but only while other tasks are ready — with no
836		// contention the system stays tickless. When the slice expires, the timer
837		// interrupt's `reschedule()` switches to the next ready task.
838		#[cfg(feature = "preemptive")]
839		if !self.ready_queue.is_empty() {
840			let deadline = processor::get_timer_ticks() + PREEMPTION_SLICE_US;
841			timer_interrupts::create_timer_abs(timer_interrupts::Source::Preemption, deadline);
842		}
843
844		// Tell the scheduler about the new task.
845		debug!(
846			"Switching task from {} to {} (stack {:#X} => {:p})",
847			id,
848			new_id,
849			unsafe { *last_stack_pointer },
850			new_stack_pointer
851		);
852		#[cfg(not(target_arch = "riscv64"))]
853		{
854			self.current_task = task;
855		}
856
857		// Finally return the context of the new task.
858		#[cfg(not(target_arch = "riscv64"))]
859		return Some(last_stack_pointer);
860
861		#[cfg(target_arch = "riscv64")]
862		{
863			if sstatus::read().fs() == sstatus::FS::Dirty {
864				self.current_task.borrow_mut().last_fpu_state.save();
865			}
866			task.borrow().last_fpu_state.restore();
867			self.current_task = task;
868			unsafe {
869				switch_to_task(last_stack_pointer, new_stack_pointer.as_usize());
870			}
871			None
872		}
873	}
874}
875
876fn get_tid() -> TaskId {
877	static TID_COUNTER: AtomicI32 = AtomicI32::new(0);
878	let guard = TASKS.lock();
879
880	loop {
881		let id = TaskId::from(TID_COUNTER.fetch_add(1, Ordering::SeqCst));
882		if !guard.contains_key(&id) {
883			return id;
884		}
885	}
886}
887
888#[inline]
889pub(crate) fn abort() -> ! {
890	core_scheduler().exit(-1)
891}
892
893/// Add a per-core scheduler for the current core.
894pub(crate) fn add_current_core() {
895	// Create an idle task for this core.
896	let core_id = core_id();
897	let tid = get_tid();
898	let idle_task = Rc::new(RefCell::new(Task::new_idle(tid, core_id)));
899
900	// Add the ID -> Task mapping.
901	WAITING_TASKS.lock().insert(tid, VecDeque::with_capacity(1));
902	TASKS.lock().insert(
903		tid,
904		TaskHandle::new(
905			tid,
906			IDLE_PRIO,
907			#[cfg(feature = "smp")]
908			core_id,
909		),
910	);
911	// Initialize a scheduler for this core.
912	debug!("Initializing scheduler for core {core_id} with idle task {tid}");
913	let boxed_scheduler = Box::new(PerCoreScheduler {
914		#[cfg(feature = "smp")]
915		core_id,
916		current_task: idle_task.clone(),
917		#[cfg(any(target_arch = "x86_64", target_arch = "aarch64"))]
918		fpu_owner: idle_task.clone(),
919		idle_task,
920		ready_queue: PriorityTaskQueue::new(),
921		finished_tasks: VecDeque::new(),
922		blocked_tasks: BlockedTaskQueue::new(),
923		timers: TimerList::new(),
924	});
925
926	let scheduler = Box::into_raw(boxed_scheduler);
927	set_core_scheduler(scheduler);
928	#[cfg(feature = "smp")]
929	{
930		SCHEDULER_INPUTS.lock().insert(
931			core_id.try_into().unwrap(),
932			&CoreLocal::get().scheduler_input,
933		);
934		#[cfg(all(
935			any(target_arch = "x86_64", target_arch = "riscv64"),
936			not(feature = "idle-poll")
937		))]
938		sleep_state::install_for_core(core_id);
939	}
940}
941
942#[inline]
943#[cfg(feature = "smp")]
944fn get_scheduler_input(core_id: CoreId) -> &'static InterruptTicketMutex<SchedulerInput> {
945	SCHEDULER_INPUTS.lock()[usize::try_from(core_id).unwrap()]
946}
947
948pub unsafe fn spawn(
949	func: unsafe extern "C" fn(usize),
950	arg: usize,
951	prio: Priority,
952	stack_size: usize,
953	selector: isize,
954) -> TaskId {
955	static CORE_COUNTER: AtomicU32 = AtomicU32::new(1);
956
957	let core_id = if selector < 0 {
958		// use Round Robin to schedule the cores
959		CORE_COUNTER.fetch_add(1, Ordering::SeqCst) % get_processor_count()
960	} else {
961		selector as u32
962	};
963
964	unsafe { PerCoreScheduler::spawn(func, arg, prio, core_id, stack_size) }
965}
966
967#[allow(clippy::result_unit_err)]
968pub fn join(id: TaskId) -> Result<(), ()> {
969	let core_scheduler = core_scheduler();
970
971	debug!(
972		"Task {} is waiting for task {}",
973		core_scheduler.get_current_task_id(),
974		id
975	);
976
977	loop {
978		let mut waiting_tasks_guard = WAITING_TASKS.lock();
979
980		let Some(queue) = waiting_tasks_guard.get_mut(&id) else {
981			return Ok(());
982		};
983
984		queue.push_back(core_scheduler.get_current_task_handle());
985		core_scheduler.block_current_task(None);
986
987		// Switch to the next task.
988		drop(waiting_tasks_guard);
989		core_scheduler.reschedule();
990	}
991}
992
993pub fn shutdown(arg: i32) -> ! {
994	crate::syscalls::shutdown(arg)
995}
996
997fn get_task_handle(id: TaskId) -> Option<TaskHandle> {
998	TASKS.lock().get(&id).copied()
999}
1000
1001#[cfg(all(target_arch = "x86_64", feature = "common-os"))]
1002pub(crate) static BOOT_ROOT_PAGE_TABLE: OnceCell<usize> = OnceCell::new();
1003
1004#[cfg(all(target_arch = "x86_64", feature = "common-os"))]
1005pub(crate) fn get_root_page_table() -> usize {
1006	let current_task_borrowed = core_scheduler().current_task.borrow_mut();
1007	current_task_borrowed.root_page_table
1008}