22//
33// SPDX-License-Identifier: Apache-2.0
44
5- use crate :: config:: { Config , Networking , ProcessAnnotation , Protocol } ;
6- use crate :: logrotate;
5+ use crate :: {
6+ config:: { Config , NetworkFilterMode , Networking , NetworkingMode , ProcessAnnotation , Protocol } ,
7+ logrotate,
8+ netd:: { self , InterfaceIdentity , PrepareRequest , Request as NetdRequest } ,
9+ } ;
710
811use anyhow:: { bail, Context , Result } ;
912use bon:: Builder ;
@@ -18,6 +21,7 @@ use dstack_vmm_rpc::{
1821use fs_err as fs;
1922use guest_api:: client:: DefaultClient as GuestClient ;
2023use id_pool:: IdPool ;
24+ use nix:: unistd:: { Uid , User } ;
2125use or_panic:: ResultOrPanic ;
2226use ra_rpc:: client:: RaClient ;
2327use serde:: { Deserialize , Serialize } ;
@@ -443,17 +447,30 @@ impl App {
443447 append_boot_separator ( & path) ;
444448 }
445449
450+ let runtime_networks = resolved_networks ( & vm_config. manifest , & self . config . cvm ) ;
446451 let devices = self . try_allocate_gpus ( & vm_config. manifest ) ?;
447452 let processes = vm_config. config_qemu ( & work_dir, & self . config . cvm , & devices) ?;
448- let runtime_networks = resolved_networks ( & vm_config. manifest , & self . config . cvm ) ;
449453 work_dir. set_runtime_networks ( & runtime_networks) ?;
454+ if let Err ( error) = self
455+ . prepare_filtered_networks ( & vm_config, & runtime_networks)
456+ . await
457+ {
458+ let _ = work_dir. clear_runtime_networks ( ) ;
459+ return Err ( error) ;
460+ }
450461 {
451462 let mut state = self . lock ( ) ;
452463 let vm_state = state. get_mut ( id) . context ( "VM not found" ) ?;
453- vm_state. state . runtime_networks = runtime_networks;
464+ vm_state. state . runtime_networks = runtime_networks. clone ( ) ;
454465 }
455466 for process in processes {
456467 if let Err ( err) = self . supervisor . deploy ( & process) . await {
468+ if let Err ( cleanup_error) = self
469+ . remove_filtered_networks ( & vm_config. manifest . id , & runtime_networks)
470+ . await
471+ {
472+ warn ! ( id, %cleanup_error, "failed to roll back filtered networking" ) ;
473+ }
457474 if let Err ( clear_err) = work_dir. clear_runtime_networks ( ) {
458475 warn ! (
459476 id,
@@ -488,6 +505,108 @@ impl App {
488505 }
489506 self . set_started ( id, false ) ?;
490507 self . stop_vm_process ( id) . await ?;
508+ let networks = self . work_dir ( id) ?. runtime_networks ( ) ;
509+ self . remove_filtered_networks ( id, & networks) . await ?;
510+ Ok ( ( ) )
511+ }
512+
513+ async fn prepare_filtered_networks (
514+ & self ,
515+ vm : & VmConfig ,
516+ networks : & [ Networking ] ,
517+ ) -> Result < ( ) > {
518+ if self . config . cvm . network_filter . mode == NetworkFilterMode :: None {
519+ return Ok ( ( ) ) ;
520+ }
521+ let qemu_uid = if self . config . cvm . user . is_empty ( ) {
522+ Uid :: effective ( ) . as_raw ( )
523+ } else {
524+ User :: from_name ( & self . config . cvm . user )
525+ . context ( "failed to resolve QEMU user" ) ?
526+ . with_context ( || format ! ( "QEMU user {} does not exist" , self . config. cvm. user) ) ?
527+ . uid
528+ . as_raw ( )
529+ } ;
530+ let mut prepared = Vec :: new ( ) ;
531+ for ( nic_index, network) in networks. iter ( ) . enumerate ( ) {
532+ if network. mode != NetworkingMode :: Bridge {
533+ continue ;
534+ }
535+ let identity = InterfaceIdentity {
536+ instance_id : self . config . cvm . instance_id . clone ( ) ,
537+ vm_id : vm. manifest . id . clone ( ) ,
538+ nic_index,
539+ } ;
540+ let request = PrepareRequest {
541+ identity : identity. clone ( ) ,
542+ bridge : network. bridge . clone ( ) ,
543+ mac : network:: mac_address_for_vm_index (
544+ & vm. manifest . id ,
545+ & network. mac_prefix_bytes ( ) ,
546+ nic_index,
547+ ) ,
548+ qemu_uid,
549+ filter : self . config . cvm . network_filter . filter . clone ( ) ,
550+ parameters : self . config . cvm . network_filter . parameters . clone ( ) ,
551+ } ;
552+ if let Err ( error) =
553+ netd:: request ( & self . config . netd . socket , & NetdRequest :: Prepare ( request) ) . await
554+ {
555+ // The client may have timed out while netd was still finishing
556+ // this Prepare. Remove the in-flight identity first; netd's
557+ // serialized accept loop processes it after Prepare completes.
558+ if let Err ( cleanup_error) = netd:: request (
559+ & self . config . netd . socket ,
560+ & NetdRequest :: Remove {
561+ identity : identity. clone ( ) ,
562+ } ,
563+ )
564+ . await
565+ {
566+ warn ! ( %cleanup_error, "failed to roll back in-flight filtered network" ) ;
567+ }
568+ for identity in prepared. into_iter ( ) . rev ( ) {
569+ if let Err ( cleanup_error) =
570+ netd:: request ( & self . config . netd . socket , & NetdRequest :: Remove { identity } )
571+ . await
572+ {
573+ warn ! ( %cleanup_error, "failed to roll back prepared filtered network" ) ;
574+ }
575+ }
576+ return Err ( error) . context ( "failed to prepare libvirt-filtered networking" ) ;
577+ }
578+ prepared. push ( identity) ;
579+ }
580+ Ok ( ( ) )
581+ }
582+
583+ pub ( crate ) async fn remove_filtered_networks (
584+ & self ,
585+ vm_id : & str ,
586+ networks : & [ Networking ] ,
587+ ) -> Result < ( ) > {
588+ if self . config . cvm . network_filter . mode == NetworkFilterMode :: None {
589+ return Ok ( ( ) ) ;
590+ }
591+ let mut first_error = None ;
592+ for ( nic_index, network) in networks. iter ( ) . enumerate ( ) . rev ( ) {
593+ if network. mode != NetworkingMode :: Bridge {
594+ continue ;
595+ }
596+ let identity = InterfaceIdentity {
597+ instance_id : self . config . cvm . instance_id . clone ( ) ,
598+ vm_id : vm_id. to_string ( ) ,
599+ nic_index,
600+ } ;
601+ if let Err ( error) =
602+ netd:: request ( & self . config . netd . socket , & NetdRequest :: Remove { identity } ) . await
603+ {
604+ first_error. get_or_insert ( error) ;
605+ }
606+ }
607+ if let Some ( error) = first_error {
608+ return Err ( error) . context ( "failed to remove libvirt-filtered networking" ) ;
609+ }
491610 Ok ( ( ) )
492611 }
493612
@@ -597,6 +716,11 @@ impl App {
597716 }
598717 }
599718
719+ let runtime_networks = self . work_dir ( id) ?. runtime_networks ( ) ;
720+ if let Err ( error) = self . remove_filtered_networks ( id, & runtime_networks) . await {
721+ warn ! ( id, %error, "failed to remove filtered networking during VM removal" ) ;
722+ }
723+
600724 // Only delete the workdir for user-initiated removal or if .removing marker exists.
601725 // Orphaned supervisor processes without the marker keep their data intact.
602726 let vm_path = self . work_dir ( id) ?;
0 commit comments