1mod debug;
7#[cfg(test)]
8mod tests;
9mod vp_set;
10
11pub use vp_set::Halt;
12pub use vp_set::RequestYield;
13pub use vp_set::RunCancelled;
14pub use vp_set::RunnerCanceller;
15pub use vp_set::VpRunner;
16pub use vp_set::block_on_vp;
17
18use self::vp_set::RegisterSetError;
19#[cfg(feature = "dump")]
20use anyhow::Context as _;
21use async_trait::async_trait;
22use futures::FutureExt;
23use futures::StreamExt;
24use guestmem::GuestMemory;
25use hvdef::Vtl;
26use inspect::InspectMut;
27use mesh::Receiver;
28use mesh::rpc::Rpc;
29use mesh::rpc::RpcSend;
30use pal_async::task::Spawn;
31use state_unit::NameInUse;
32use state_unit::SpawnedUnit;
33use state_unit::StateRequest;
34use state_unit::StateUnit;
35use state_unit::UnitBuilder;
36use state_unit::UnitHandle;
37use std::sync::Arc;
38use thiserror::Error;
39use virt::InitialPageImport;
40use virt::InitialRegs;
41use virt::InitialVpStateSource;
42#[cfg(feature = "dump")]
43use virt::VpIndex;
44use vm_topology::processor::ProcessorTopology;
45use vmcore::save_restore::ProtobufSaveRestore;
46use vmcore::save_restore::RestoreError;
47use vmcore::save_restore::SaveError;
48use vmcore::save_restore::SavedStateBlob;
49use vmm_core_defs::HaltReason;
50use vp_set::VpSet;
51
52pub struct PartitionUnit {
54 handle: SpawnedUnit<PartitionUnitRunner>,
55 req_send: mesh::Sender<PartitionRequest>,
56}
57
58#[async_trait]
60pub trait VmPartition: 'static + Send + Sync + InspectMut + ProtobufSaveRestore {
61 fn initial_vp_state_source(&self) -> InitialVpStateSource;
63
64 fn freeze_time(&mut self) {}
67
68 fn thaw_time(&mut self) {}
71
72 fn reset(&mut self) -> anyhow::Result<()>;
74
75 fn scrub_vtl(&mut self, vtl: Vtl) -> anyhow::Result<()>;
77
78 fn accept_initial_pages(&mut self, pages: Vec<InitialPageImport>) -> anyhow::Result<()>;
80
81 fn guest_os_id(&self) -> u64 {
86 0
87 }
88}
89
90struct PartitionUnitRunner {
92 partition: Box<dyn VmPartition>,
93 vp_set: VpSet,
94 unit_started: bool,
95 vp_stop_count: usize,
96 needs_reset: bool,
97 halt_reason: Option<HaltReason>,
98 halt_request_recv: Receiver<InternalHaltReason>,
99 client_notify_send: mesh::Sender<HaltReason>,
100 req_recv: Receiver<PartitionRequest>,
101 topology: ProcessorTopology,
102 initial_regs: Option<Arc<InitialRegs>>,
103
104 #[cfg(feature = "gdb")]
105 debugger_state: debug::DebuggerState,
106}
107
108impl InspectMut for PartitionUnitRunner {
109 fn inspect_mut(&mut self, req: inspect::Request<'_>) {
110 req.respond()
111 .field(
112 "power_state",
113 self.halt_reason.as_ref().map_or("running", |_| "halted"),
114 )
115 .merge(&self.halt_reason)
116 .merge(&self.vp_set)
117 .field_mut_with("clear_halt", |clear| {
118 if let Some(clear) = clear {
120 match clear.parse::<bool>() {
121 Ok(x) => {
122 if x {
123 self.clear_halt();
124 }
125 Ok(x)
126 }
127 Err(err) => Err(err),
128 }
129 } else {
130 Ok(false)
131 }
132 })
133 .field("topology", &self.topology)
134 .merge(&mut self.partition);
135 }
136}
137
138enum PartitionRequest {
139 ClearHalt(Rpc<(), bool>), SetInitialRegs(Rpc<(Vtl, Arc<InitialRegs>), Result<(), InitialRegError>>),
141 AcceptInitialPages(Rpc<Vec<InitialPageImport>, Result<(), AcceptInitialPagesError>>),
142 StopVps(Rpc<(), ()>),
143 StartVps,
144 #[cfg(feature = "dump")]
146 BuildDumpPartitionState(Rpc<(), anyhow::Result<Vec<u8>>>),
147}
148
149pub struct PartitionUnitParams<'a> {
150 pub vtl_guest_memory: [Option<&'a GuestMemory>; 3],
151 pub processor_topology: &'a ProcessorTopology,
152 pub halt_vps: Arc<Halt>,
154 pub halt_request_recv: HaltReasonReceiver,
156 pub client_notify_send: mesh::Sender<HaltReason>,
159 pub debugger_rpc: Option<Receiver<vmm_core_defs::debug_rpc::DebugRequest>>,
160}
161
162pub struct HaltReasonReceiver(Receiver<InternalHaltReason>);
164
165enum InternalHaltReason {
166 Halt(HaltReason),
167 ReplayMtrrs,
168}
169
170#[derive(Debug, Error)]
172pub enum Error {
173 #[error("debugging is not supported in this build")]
174 DebuggingNotSupported,
175 #[error(transparent)]
176 NameInUse(NameInUse),
177 #[error("missing guest memory required for gdb support")]
178 MissingGuestMemory,
179}
180
181#[derive(Debug, Error)]
183pub enum InitialRegError {
184 #[error("failed to set registers")]
185 RegisterSet(#[source] RegisterSetError),
186 #[error("failed to scrub VTL state")]
187 ScrubVtl(#[source] anyhow::Error),
188}
189
190#[derive(Debug, Error)]
192pub enum AcceptInitialPagesError {
193 #[error("failed to finalize initial page imports")]
194 Finalize(#[source] anyhow::Error),
195}
196
197impl PartitionUnit {
198 pub fn new(
203 spawner: impl Spawn,
204 builder: UnitBuilder<'_>,
205 partition: impl VmPartition,
206 params: PartitionUnitParams<'_>,
207 ) -> Result<(Self, Vec<VpRunner>), Error> {
208 #[cfg(not(feature = "gdb"))]
209 if params.debugger_rpc.is_some() {
210 return Err(Error::DebuggingNotSupported);
211 }
212
213 let mut vp_set = VpSet::new(params.vtl_guest_memory.map(|m| m.cloned()), params.halt_vps);
214 let vps = params
215 .processor_topology
216 .vps_arch()
217 .map(|vp| vp_set.add(vp))
218 .collect();
219
220 let (req_send, req_recv) = mesh::channel();
221
222 let mut runner = PartitionUnitRunner {
223 partition: Box::new(partition),
224 vp_set,
225 unit_started: false,
226 vp_stop_count: 0,
227 needs_reset: false,
228 halt_reason: None,
229 halt_request_recv: params.halt_request_recv.0,
230 client_notify_send: params.client_notify_send,
231 req_recv,
232 topology: params.processor_topology.clone(),
233 initial_regs: None,
234 #[cfg(feature = "gdb")]
235 debugger_state: debug::DebuggerState::new(
236 params.vtl_guest_memory[0]
237 .ok_or(Error::MissingGuestMemory)?
238 .clone(),
239 params.debugger_rpc,
240 ),
241 };
242
243 let handle = builder
244 .spawn(spawner, async |recv| {
245 runner.run(recv).await;
246 runner
247 })
248 .unwrap();
249
250 Ok((Self { handle, req_send }, vps))
251 }
252
253 pub fn unit_handle(&self) -> &UnitHandle {
255 self.handle.handle()
256 }
257
258 pub async fn teardown(self) -> mesh::Sender<HaltReason> {
261 let runner = self.handle.remove().await;
262 runner.vp_set.teardown().await;
263 runner.client_notify_send
264 }
265
266 pub async fn clear_halt(&mut self) -> bool {
269 self.req_send
270 .call(PartitionRequest::ClearHalt, ())
271 .await
272 .unwrap()
273 }
274
275 pub async fn temporarily_stop_vps(&mut self) -> StopGuard {
278 self.req_send
279 .call(PartitionRequest::StopVps, ())
280 .await
281 .unwrap();
282
283 StopGuard(self.req_send.clone())
284 }
285
286 pub async fn set_initial_regs(
292 &mut self,
293 vtl: Vtl,
294 state: Arc<InitialRegs>,
295 ) -> Result<(), InitialRegError> {
296 self.req_send
297 .call(PartitionRequest::SetInitialRegs, (vtl, state))
298 .await
299 .unwrap()
300 }
301
302 pub async fn accept_initial_pages(
303 &mut self,
304 initial_pages: Vec<InitialPageImport>,
305 ) -> Result<(), AcceptInitialPagesError> {
306 self.req_send
307 .call(PartitionRequest::AcceptInitialPages, initial_pages)
308 .await
309 .unwrap()
310 }
311
312 #[cfg(feature = "dump")]
318 pub async fn build_dump_partition_state(&mut self) -> anyhow::Result<Vec<u8>> {
319 self.req_send
320 .call(PartitionRequest::BuildDumpPartitionState, ())
321 .await
322 .unwrap()
323 }
324}
325
326impl PartitionUnitRunner {
327 async fn run(&mut self, mut recv: Receiver<StateRequest>) {
329 loop {
330 enum Event {
331 State(Option<StateRequest>),
332 Halt(InternalHaltReason),
333 Request(PartitionRequest),
334 #[cfg(feature = "gdb")]
335 Debug(vmm_core_defs::debug_rpc::DebugRequest),
336 }
337
338 #[cfg(feature = "gdb")]
339 let debug = self.debugger_state.wait_rpc();
340 #[cfg(not(feature = "gdb"))]
341 let debug = std::future::pending();
342
343 let event = futures::select! { request = recv.next() => Event::State(request),
345 request = self.halt_request_recv.select_next_some() => Event::Halt(request),
346 request = self.req_recv.select_next_some() => Event::Request(request),
347 request = debug.fuse() => {
348 #[cfg(feature = "gdb")]
349 {
350 Event::Debug(request)
351 }
352 #[cfg(not(feature = "gdb"))]
353 {
354 let _: std::convert::Infallible = request;
355 unreachable!()
356 }
357 }
358 };
359
360 match event {
361 Event::State(request) => {
362 if let Some(request) = request {
363 request.apply(self).await;
364 } else {
365 break;
366 }
367 }
368 Event::Halt(reason) => {
369 self.vp_set.stop().await;
375 self.handle_halt(reason).await;
376 }
377 Event::Request(request) => match request {
378 PartitionRequest::ClearHalt(rpc) => rpc.handle_sync(|()| self.clear_halt()),
379 PartitionRequest::SetInitialRegs(rpc) => {
380 rpc.handle(async |(vtl, state)| self.set_initial_regs(vtl, state).await)
381 .await
382 }
383 PartitionRequest::AcceptInitialPages(rpc) => {
384 rpc.handle(async |initial_pages| {
385 self.accept_initial_pages(initial_pages).await
386 })
387 .await
388 }
389 PartitionRequest::StopVps(rpc) => {
390 rpc.handle(async |()| self.stop_vps().await).await
391 }
392 PartitionRequest::StartVps => {
393 self.resume_vps();
394 }
395 #[cfg(feature = "dump")]
396 PartitionRequest::BuildDumpPartitionState(rpc) => {
397 rpc.handle(async |()| self.build_dump_partition_state().await)
398 .await
399 }
400 },
401 #[cfg(feature = "gdb")]
402 Event::Debug(request) => {
403 self.handle_gdb(request).await;
404 }
405 }
406 }
407
408 if self.unit_started {
409 StateUnit::stop(self).await;
410 }
411 }
412
413 async fn handle_halt(&mut self, reason: InternalHaltReason) {
414 match reason {
415 InternalHaltReason::Halt(reason) => {
416 if self.halt_reason.is_none() {
420 self.halt_reason = Some(reason.clone());
421
422 #[cfg(feature = "gdb")]
424 let reported = self.debugger_state.report_halt_to_debugger(&reason);
425 #[cfg(not(feature = "gdb"))]
426 let reported = false;
427
428 if !reported {
431 self.client_notify_send.send(reason);
432 }
433 } else {
434 self.vp_set.clear_halt();
436 }
437 }
438 InternalHaltReason::ReplayMtrrs => {
439 if let Some(initial_regs) = self.initial_regs.clone() {
440 if let Err(err) = self
441 .vp_set
442 .set_initial_regs(
443 Vtl::Vtl0,
444 initial_regs,
445 vp_set::RegistersToSet::MtrrsOnly,
446 )
447 .await
448 {
449 tracing::error!(
450 error = &err as &dyn std::error::Error,
451 "failed to replay mtrrs, guest may see inconsistent results"
452 );
453 }
454 } else {
455 tracing::warn!("no initial mtrrs to replay");
456 }
457 self.vp_set.clear_halt();
458 self.try_start();
459 }
460 }
461 }
462
463 fn clear_halt(&mut self) -> bool {
466 if self.halt_reason.is_some() {
467 self.halt_reason = None;
468 self.vp_set.clear_halt();
469 self.try_start();
470 true
471 } else {
472 false
473 }
474 }
475
476 async fn set_initial_regs(
477 &mut self,
478 vtl: Vtl,
479 state: Arc<InitialRegs>,
480 ) -> Result<(), InitialRegError> {
481 assert!(!self.unit_started || self.vp_stop_count > 0);
482
483 if self.needs_reset {
486 self.partition
487 .scrub_vtl(vtl)
488 .map_err(InitialRegError::ScrubVtl)?;
489 self.vp_set
490 .scrub(vtl)
491 .await
492 .map_err(InitialRegError::ScrubVtl)?;
493 self.needs_reset = false;
494 }
495
496 match self.partition.initial_vp_state_source() {
497 InitialVpStateSource::Registers => {
498 self.vp_set
499 .set_initial_regs(vtl, state.clone(), vp_set::RegistersToSet::All)
500 .await
501 .map_err(InitialRegError::RegisterSet)?;
502 }
503 InitialVpStateSource::ImportedContext => {}
504 }
505
506 self.initial_regs = Some(state);
507 Ok(())
508 }
509
510 async fn accept_initial_pages(
511 &mut self,
512 initial_pages: Vec<InitialPageImport>,
513 ) -> Result<(), AcceptInitialPagesError> {
514 assert!(!self.unit_started);
515
516 self.partition
517 .accept_initial_pages(initial_pages)
518 .map_err(AcceptInitialPagesError::Finalize)
519 }
520
521 fn try_start(&mut self) {
522 if self.unit_started && self.halt_reason.is_none() && self.vp_stop_count == 0 {
523 self.partition.thaw_time();
525 self.needs_reset = true;
526 self.vp_set.start();
527 }
528 }
529
530 async fn stop_vps(&mut self) {
531 self.vp_set.stop().await;
532 self.vp_stop_count += 1;
533 }
534
535 fn resume_vps(&mut self) {
536 assert!(
537 self.vp_stop_count > 0,
538 "resume_vps called without matching stop"
539 );
540 self.vp_stop_count -= 1;
541 self.try_start();
542 }
543}
544
545#[cfg(feature = "dump")]
546impl PartitionUnitRunner {
547 async fn build_dump_partition_state(&mut self) -> anyhow::Result<Vec<u8>> {
551 self.stop_vps().await;
554 let result = self.build_dump_partition_state_inner().await;
555 self.resume_vps();
556 result
557 }
558
559 async fn build_dump_partition_state_inner(&mut self) -> anyhow::Result<Vec<u8>> {
560 use hyperv_dump::PartitionStateBuilder;
561 use hyperv_dump::ProcessorArch;
562
563 #[cfg(guest_arch = "x86_64")]
564 let arch = ProcessorArch::X64;
565 #[cfg(guest_arch = "aarch64")]
566 let arch = ProcessorArch::Aarch64;
567
568 let mut builder = PartitionStateBuilder::new(arch);
569 builder.set_os_id(self.partition.guest_os_id());
570
571 let vp_count = self.topology.vp_count();
572 for vp_idx in 0..vp_count {
573 let vtl = Vtl::Vtl0;
574 let vp_state = self
575 .vp_set
576 .get_dump_vp_state(VpIndex::new(vp_idx), vtl)
577 .await
578 .with_context(|| format!("failed to get state for VP {vp_idx}"))?;
579
580 builder.add_vp(vp_idx, vec![(vtl, vp_state)]);
581 }
582
583 Ok(builder.finish())
584 }
585}
586
587#[must_use = "when dropped, the VPs will be resumed"]
588pub struct StopGuard(mesh::Sender<PartitionRequest>);
589
590impl Drop for StopGuard {
591 fn drop(&mut self) {
592 self.0.send(PartitionRequest::StartVps);
593 }
594}
595
596impl StateUnit for PartitionUnitRunner {
597 async fn start(&mut self) {
598 self.partition.thaw_time();
599 self.unit_started = true;
600 self.try_start();
601 }
602
603 async fn stop(&mut self) {
604 self.vp_set.stop().await;
605 self.partition.freeze_time();
606 self.unit_started = false;
607
608 while let Ok(reason) = self.halt_request_recv.try_recv() {
611 self.handle_halt(reason).await;
612 }
613 }
614
615 async fn reset(&mut self) -> anyhow::Result<()> {
616 self.partition.reset()?;
617 self.vp_set.reset().await?;
618 self.clear_halt();
619 self.needs_reset = false;
620 Ok(())
621 }
622
623 async fn save(&mut self) -> Result<Option<SavedStateBlob>, SaveError> {
624 let state = self.save().await?;
625 Ok(Some(SavedStateBlob::new(state)))
626 }
627
628 async fn restore(&mut self, buffer: SavedStateBlob) -> Result<(), RestoreError> {
629 self.needs_reset = true;
631 self.restore(buffer.parse()?).await?;
632 Ok(())
633 }
634}
635
636mod save_restore {
637 use super::PartitionUnitRunner;
638 use virt::VpIndex;
639 use vmcore::save_restore::RestoreError;
640 use vmcore::save_restore::SaveError;
641
642 mod state {
643 use mesh::payload::Protobuf;
644 use vmcore::save_restore::SavedStateBlob;
645 use vmcore::save_restore::SavedStateRoot;
646
647 #[derive(Protobuf, SavedStateRoot)]
648 #[mesh(package = "partition")]
649 pub struct Partition {
650 #[mesh(1)]
651 pub(super) partition: SavedStateBlob,
652 #[mesh(2)]
653 pub(super) vps: Vec<Vp>,
654 }
656
657 #[derive(Protobuf)]
658 #[mesh(package = "partition")]
659 pub struct Vp {
660 #[mesh(1)]
661 pub vp_index: u32,
662 #[mesh(2)]
663 pub data: SavedStateBlob,
664 }
665 }
666
667 impl PartitionUnitRunner {
668 pub async fn save(&mut self) -> Result<state::Partition, SaveError> {
669 let partition = self.partition.save()?;
670 let vps = self.vp_set.save().await?;
671 let vps = vps
672 .into_iter()
673 .map(|(vp_index, data)| state::Vp {
674 vp_index: vp_index.index(),
675 data,
676 })
677 .collect();
678
679 Ok(state::Partition { partition, vps })
680 }
681
682 pub async fn restore(&mut self, state: state::Partition) -> Result<(), RestoreError> {
683 let state::Partition { partition, vps } = state;
684 self.partition.restore(partition)?;
685 self.vp_set
686 .restore(
687 vps.into_iter()
688 .map(|state::Vp { vp_index, data }| (VpIndex::new(vp_index), data)),
689 )
690 .await?;
691 Ok(())
692 }
693 }
694}