@@ -28,7 +28,7 @@ use curvine_common::state::{
2828 CreateFileOpts , FileBlocks , FileStatus , HeartbeatStatus , OpenFlags , RenameFlags ,
2929} ;
3030use curvine_common:: utils:: ProtoUtils ;
31- use curvine_common:: version:: Version ;
31+ use curvine_common:: version:: { CompatibilityPolicy , CompatibilityResult , Version , VersionChecker } ;
3232use curvine_common:: FsResult ;
3333use orpc:: err_box;
3434use orpc:: handler:: MessageHandler ;
@@ -46,6 +46,8 @@ pub struct MasterHandler {
4646 pub ( crate ) job_handler : JobHandler ,
4747 pub ( crate ) mount_manager : Arc < MountManager > ,
4848 pub ( crate ) replication_handler : Option < MasterReplicationHandler > ,
49+ pub ( crate ) worker_version_checker : VersionChecker ,
50+ pub ( crate ) client_version_checker : VersionChecker ,
4951}
5052
5153impl MasterHandler {
@@ -59,6 +61,34 @@ impl MasterHandler {
5961 job_handler : JobHandler ,
6062 replication_manager : Arc < MasterReplicationManager > ,
6163 ) -> Self {
64+ let master_version = Version :: current ( ) ;
65+
66+ // Initialize Worker version checker
67+ let worker_min_version =
68+ Version :: from_str ( & conf. master . min_worker_version ) . unwrap_or_else ( |e| {
69+ warn ! (
70+ "Failed to parse min_worker_version '{}': {}, using default 0.1.0" ,
71+ conf. master. min_worker_version, e
72+ ) ;
73+ Version :: new ( 0 , 1 , 0 )
74+ } ) ;
75+ let worker_policy = CompatibilityPolicy :: new ( worker_min_version) ;
76+ let worker_version_checker = VersionChecker :: new ( master_version, worker_policy) ;
77+
78+ // Initialize Client version checker (looser policy - no upper bound)
79+ let client_min_version =
80+ Version :: from_str ( & conf. master . min_client_version ) . unwrap_or_else ( |e| {
81+ warn ! (
82+ "Failed to parse min_client_version '{}': {}, using default 0.1.0" ,
83+ conf. master. min_client_version, e
84+ ) ;
85+ Version :: new ( 0 , 1 , 0 )
86+ } ) ;
87+ let client_policy = CompatibilityPolicy :: new ( client_min_version) ;
88+ // Client uses a very large upper bound for compatibility
89+ let client_version_checker =
90+ VersionChecker :: new ( Version :: new ( u32:: MAX , u32:: MAX , u32:: MAX ) , client_policy) ;
91+
6292 Self {
6393 fs,
6494 retry_cache,
@@ -68,6 +98,8 @@ impl MasterHandler {
6898 mount_manager,
6999 job_handler,
70100 replication_handler : Some ( MasterReplicationHandler :: new ( replication_manager) ) ,
101+ worker_version_checker,
102+ client_version_checker,
71103 }
72104 }
73105
@@ -384,7 +416,36 @@ impl MasterHandler {
384416 }
385417
386418 pub fn get_master_info ( & self , ctx : & mut RpcContext < ' _ > ) -> FsResult < Message > {
387- let _: GetMasterInfoRequest = ctx. parse_header ( ) ?;
419+ let header: GetMasterInfoRequest = ctx. parse_header ( ) ?;
420+
421+ // Check client version if provided
422+ if let Some ( client_version_str) = & header. client_version {
423+ if !client_version_str. is_empty ( ) {
424+ let client_version = Version :: from_str ( client_version_str) . map_err ( |e| {
425+ FsError :: common ( format ! (
426+ "Failed to parse client version '{}': {}" ,
427+ client_version_str, e
428+ ) )
429+ } ) ?;
430+
431+ let compat_result = self
432+ . client_version_checker
433+ . check_compatibility ( & client_version) ;
434+ if let CompatibilityResult :: Incompatible ( reason) = compat_result {
435+ warn ! ( "Client version {} rejected: {}" , client_version, reason) ;
436+ return Err ( FsError :: version_incompatible ( reason) ) ;
437+ }
438+
439+ // Version check passed
440+ log:: info!( "Client version {} accepted" , client_version) ;
441+ } else {
442+ // Empty version string - backward compatibility
443+ warn ! ( "Client connected without version information (backward compatibility mode)" ) ;
444+ }
445+ } else {
446+ // No version field - old client (backward compatibility)
447+ warn ! ( "Client connected without version information (backward compatibility mode)" ) ;
448+ }
388449
389450 let info = self . fs . master_info ( ) ?;
390451 let rep_header = ProtoUtils :: master_info_to_pb ( info) ;
@@ -393,7 +454,6 @@ impl MasterHandler {
393454
394455 pub fn worker_heartbeat ( & self , ctx : & mut RpcContext < ' _ > ) -> FsResult < Message > {
395456 let header: WorkerHeartbeatRequest = ctx. parse_header ( ) ?;
396- let mut wm = self . fs . worker_manager . write ( ) ;
397457
398458 // Parse worker version from request
399459 let worker_version_str = & header. version ;
@@ -410,10 +470,25 @@ impl MasterHandler {
410470 Version :: new ( 0 , 1 , 0 )
411471 } ;
412472
473+ // Check version compatibility at Handler layer
474+ let worker_addr = ProtoUtils :: worker_address_from_pb ( & header. address ) ;
475+ let compat_result = self
476+ . worker_version_checker
477+ . check_compatibility ( & worker_version) ;
478+ if let CompatibilityResult :: Incompatible ( reason) = compat_result {
479+ warn ! (
480+ "Worker {} registration rejected due to version incompatibility: {}" ,
481+ worker_addr, reason
482+ ) ;
483+ return Err ( FsError :: version_incompatible ( reason) ) ;
484+ }
485+
486+ // Version check passed, proceed with heartbeat
487+ let mut wm = self . fs . worker_manager . write ( ) ;
413488 let cmds = wm. heartbeat (
414489 & header. cluster_id ,
415490 HeartbeatStatus :: from ( header. status ) ,
416- ProtoUtils :: worker_address_from_pb ( & header . address ) ,
491+ worker_addr ,
417492 ProtoUtils :: storage_info_list_from_pb ( header. storages ) ,
418493 worker_version,
419494 ) ?;
0 commit comments