1818
1919from keepersdk .helpers .pam_user_record_facade import PamUserRecordFacade
2020from keepersdk .helpers .keeper_dag .jobs import Jobs
21- from keepersdk .helpers .keeper_dag .dag_types import (CredentialBase , DiscoveryDelta , DiscoveryObject , JobItem , UserAcl , DirectoryInfo ,
21+ from keepersdk .helpers .keeper_dag .dag_types import (CredentialBase , DiscoveryDelta , DiscoveryObject , JobItem , Settings , UserAcl , DirectoryInfo ,
2222 BulkRecordConvert , BulkRecordAdd , BulkRecordSuccess , BulkProcessResults , NormalizedRecord , BulkRecordFail , PromptResult ,
2323 PromptActionEnum , RecordField )
2424from keepersdk .helpers .keeper_dag .dag_vertex import DAGVertex
@@ -152,7 +152,7 @@ def print_job_detail(vault: vault_online.VaultOnline,
152152 job_id : str ):
153153
154154 def _find_job (configuration_record ) -> Optional [Dict ]:
155- jobs_obj = Jobs (record = configuration_record )
155+ jobs_obj = Jobs (record = configuration_record , vault = vault )
156156 job_item = jobs_obj .get_job (job_id )
157157 if job_item is not None :
158158 return {
@@ -167,7 +167,7 @@ def _find_job(configuration_record) -> Optional[Dict]:
167167 if gateway_context is not None :
168168 jobs = payload ["jobs" ]
169169 job = jobs .get_job (job_id )
170- infra = Infrastructure (record = gateway_context .configuration )
170+ infra = Infrastructure (record = gateway_context .configuration , vault = vault )
171171
172172 status = "RUNNING"
173173 if job .end_ts is not None and not job .error :
@@ -296,7 +296,7 @@ def execute(self, context: KeeperParams, **kwargs):
296296 if len (gateway_context .gateway_name ) > max_gateway_name :
297297 max_gateway_name = len (gateway_context .gateway_name )
298298
299- jobs = Jobs (record = configuration_record )
299+ jobs = Jobs (record = configuration_record , vault = vault )
300300 if show_history is True :
301301 job_list = reversed (jobs .history )
302302 else :
@@ -391,7 +391,7 @@ def execute(self, context: KeeperParams, **kwargs):
391391 multi_conf_msg (gateway , err )
392392 return
393393
394- jobs = Jobs (record = gateway_context .configuration )
394+ jobs = Jobs (record = gateway_context .configuration , vault = vault )
395395 current_job_item = jobs .current_job
396396 removed_prior_job = None
397397 if current_job_item is not None :
@@ -467,15 +467,20 @@ def execute(self, context: KeeperParams, **kwargs):
467467 setattr (c , key , obj [key ])
468468 credentials .append (c .model_dump ())
469469
470+ user_map_entries = self .make_protobuf_user_map (
471+ context = context ,
472+ gateway_context = gateway_context
473+ )
474+ if len (user_map_entries ) == 0 :
475+ logger .info (
476+ "No pamUser records are linked to this configuration; "
477+ "discovery will run without an existing user map."
478+ )
479+
470480 action_inputs = GatewayActionDiscoverJobStartInputs (
471481 configuration_uid = gateway_context .configuration_uid ,
472482 resource_uid = kwargs .get ('resource_uid' ),
473- user_map = gateway_context .encrypt (
474- self .make_protobuf_user_map (
475- context = context ,
476- gateway_context = gateway_context
477- )[0 ]
478- ),
483+ user_map = gateway_context .encrypt ({"users" : user_map_entries }),
479484
480485 shared_folder_uid = gateway_context .default_shared_folder_uid ,
481486 languages = [kwargs .get ('language' )],
@@ -507,16 +512,39 @@ def execute(self, context: KeeperParams, **kwargs):
507512 logger .error (f"The router returned a failure." )
508513 return
509514
515+ discovery_settings = Settings (
516+ credentials = [CredentialBase (** c ) for c in credentials ],
517+ default_shared_folder_uid = gateway_context .default_shared_folder_uid ,
518+ include_azure_aadds = kwargs .get ('include_azure_aadds' , False ),
519+ skip_rules = kwargs .get ('skip_rules' , False ),
520+ skip_machines = kwargs .get ('skip_machines' , False ),
521+ skip_databases = kwargs .get ('skip_databases' , False ),
522+ skip_directories = kwargs .get ('skip_directories' , False ),
523+ skip_cloud_users = kwargs .get ('skip_cloud_users' , False ),
524+ user_map = user_map_entries or None ,
525+ )
526+ job_id = jobs .start (
527+ settings = discovery_settings ,
528+ resource_uid = kwargs .get ('resource_uid' ),
529+ conversation_id = conversation_id ,
530+ )
531+ jobs .close ()
532+
510533 if "has been queued" in data .get ("Response" , "" ):
511534
512535 if removed_prior_job is None :
513- logger .info ("The discovery job is currently running." )
536+ logger .info (f"Discovery job { job_id } is running." )
514537 else :
515- logger .info (f"Active discovery job { removed_prior_job } has been removed and new discovery job is running." )
538+ logger .info (
539+ f"Active discovery job { removed_prior_job } has been removed; "
540+ f"discovery job { job_id } is running."
541+ )
516542 logger .info (f"To check the status, use the command 'pam action discover status'." )
517- logger .info (f"To stop and remove the current job, use the command 'pam action discover remove -j <Job ID> '." )
543+ logger .info (f"To stop and remove the current job, use the command 'pam action discover remove -j { job_id } '." )
518544 else :
519545 router_utils .print_router_response (router_response , "job_info" , conversation_id , gateway_uid = gateway_context .gateway_uid )
546+ logger .info (f"Discovery job { job_id } was recorded on the configuration." )
547+ logger .info (f"To check the status, use the command 'pam action discover status -j { job_id } '." )
520548
521549 @staticmethod
522550 def make_protobuf_user_map (context : KeeperParams , gateway_context : GatewayContext ) -> List [dict ]:
@@ -580,7 +608,7 @@ def execute(self, context: KeeperParams, **kwargs):
580608 all_gateways = GatewayContext .all_gateways (vault )
581609
582610 def _find_job (configuration_record ) -> Optional [Dict ]:
583- jobs_obj = Jobs (record = configuration_record )
611+ jobs_obj = Jobs (record = configuration_record , vault = vault )
584612 job_item = jobs_obj .get_job (job_id )
585613 if job_item is not None :
586614 return {
@@ -1775,7 +1803,7 @@ def _get_directory_info(domain: str,
17751803 def remove_job (context : KeeperParams , configuration_record : vault_record .KeeperRecord , job_id : str ):
17761804
17771805 try :
1778- jobs = Jobs (record = configuration_record , context = context )
1806+ jobs = Jobs (record = configuration_record , vault = context . vault )
17791807 jobs .cancel (job_id )
17801808 logger .info (f"No items left to process. Removing completed discovery job." )
17811809 except Exception as err :
@@ -1786,7 +1814,7 @@ def preview(self, job_item: JobItem, context: KeeperParams, gateway_context: Gat
17861814
17871815 sync_point = job_item .sync_point
17881816 infra = Infrastructure (record = gateway_context .configuration ,
1789- context = context ,
1817+ vault = context . vault ,
17901818 logger = logger ,
17911819 debug_level = debug_level )
17921820 infra .load (sync_point )
@@ -1941,7 +1969,7 @@ def execute(self, context: KeeperParams, **kwargs):
19411969
19421970 # Get the current job.
19431971 # There can only be one active job.
1944- jobs = Jobs (record = configuration_record , context = context , logger = logger , debug_level = debug_level )
1972+ jobs = Jobs (record = configuration_record , vault = vault , logger = logger , debug_level = debug_level )
19451973 job_item = jobs .current_job
19461974 if job_item is None :
19471975 continue
0 commit comments