Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-09-01 09:33:46

0001 """
0002 Core Data Carousel operations shared by the synchronous and the asynchronous APIs.
0003 
0004 Each operation takes an already initialized DataCarouselInterface and returns the
0005 (success, message, data) triple that pandaserver.api.v1.data_carousel_api packs into its
0006 HTTP response. The synchronous endpoints call these functions directly, while the
0007 asynchronous endpoints queue a request in async_requests and
0008 pandaserver.asyncprocess.data_carousel_handlers calls the very same functions from the
0009 async request daemon. Keeping the bodies here is what makes the two paths interchangeable.
0010 """
0011 
0012 from concurrent.futures import ThreadPoolExecutor
0013 
0014 from pandacommon.pandalogger.LogWrapper import LogWrapper
0015 from pandacommon.pandalogger.PandaLogger import PandaLogger
0016 from pandacommon.pandautils.PandaUtils import naive_utcnow
0017 
0018 from pandaserver.taskbuffer.DataCarousel import (
0019     DataCarouselInterface,
0020     DataCarouselRequestSpec,
0021     DataCarouselRequestStatus,
0022 )
0023 
0024 # deliberately the same logger as pandaserver.api.v1.data_carousel_api, so an operation logs
0025 # under one name whether it ran inline in the API or in the async request daemon
0026 _logger = PandaLogger().getLogger("api_data_carousel")
0027 
0028 # (success, message, data) returned by every operation
0029 OperationResult = tuple[bool, str, dict | None]
0030 
0031 # cap on the threads submitting iDDS requests in parallel, so a request with many related tasks
0032 # can't spawn an unbounded number of threads in the API or daemon process
0033 IDDS_SUBMISSION_MAX_WORKERS = 8
0034 
0035 
0036 def validate_target(request_id: int | str | None = None, dataset: str | None = None) -> tuple[bool, str]:
0037     """
0038     Check that the arguments identify a request, without touching the DB.
0039 
0040     Used by the asynchronous endpoints to reject bad input before queueing anything; the
0041     synchronous endpoints don't need it since they report the failure to resolve directly.
0042 
0043     Args:
0044         request_id (int|str|None): request_id of the staging request
0045         dataset (str|None): dataset name of the staging request
0046 
0047     Returns:
0048         tuple[bool, str]: (valid, error message)
0049     """
0050     if request_id is None and dataset is None:
0051         return False, "either request_id or dataset must be provided"
0052     if request_id is not None:
0053         try:
0054             int(request_id)
0055         except (TypeError, ValueError):
0056             return False, f"invalid request_id: {request_id}"
0057     return True, ""
0058 
0059 
0060 def _resolve_request(dcif: DataCarouselInterface, request_id: int | str | None, dataset: str | None) -> DataCarouselRequestSpec | None:
0061     """
0062     Get the spec of the request specified by request_id or dataset (request_id is taken if both exist).
0063 
0064     Args:
0065         dcif (DataCarouselInterface): Data Carousel interface
0066         request_id (int|str|None): request_id of the staging request; may come in as a string
0067             since parameters round-trip through JSON on the asynchronous path
0068         dataset (str|None): dataset name of the staging request
0069 
0070     Returns:
0071         DataCarouselRequestSpec|None: spec of the request, or None if not found
0072     """
0073     if request_id is not None:
0074         # specified by request_id
0075         return dcif.get_request_by_id(int(request_id))
0076     elif dataset is not None:
0077         # specified by dataset
0078         return dcif.get_request_by_dataset(dataset)
0079     return None
0080 
0081 
0082 def change_staging_destination(dcif: DataCarouselInterface, request_id: int | str | None = None, dataset: str | None = None) -> OperationResult:
0083     """
0084     Change destination of staging
0085 
0086     The current active staging request will be cancelled, and a new request will be created with the newly selected destination RSE, excluding the original destination.
0087     The requests can be specified by request_id or dataset (if both exist, request_id is taken).
0088 
0089     Args:
0090         dcif (DataCarouselInterface): Data Carousel interface
0091         request_id (int|str|None): request_id of the staging request, e.g. `123`
0092         dataset (str|None): dataset name of the staging request in the format of Rucio DID
0093 
0094     Returns:
0095         tuple[bool, str, dict|None]: (success, message, data)
0096     """
0097     tmp_logger = LogWrapper(_logger, f"change_staging_destination request_id={request_id} dataset={dataset}")
0098     tmp_logger.debug("Start")
0099     success, message, data = False, "", None
0100     dc_req_spec_resubmitted = None
0101     to_submit_idds = False
0102     time_start = naive_utcnow()
0103 
0104     dc_req_spec = _resolve_request(dcif, request_id, dataset)
0105 
0106     if dc_req_spec is not None:
0107         dc_req_spec_resubmitted, err_msg = dcif.resubmit_request(dc_req_spec, submit_idds_request=False, exclude_prev_dst=True)
0108         if not dc_req_spec_resubmitted or err_msg:
0109             err_msg = f"failed to resubmit request_id={dc_req_spec.request_id} : {err_msg}"
0110             tmp_logger.error(err_msg)
0111             success, message = False, err_msg
0112         else:
0113             to_submit_idds = True
0114     else:
0115         err_msg = f"failed to get corresponding request"
0116         tmp_logger.error(err_msg)
0117         success, message = False, err_msg
0118 
0119     if dc_req_spec_resubmitted and dc_req_spec_resubmitted.status == DataCarouselRequestStatus.staging:
0120         success = True
0121         data = {"request_id": dc_req_spec.request_id, "new_request_id": dc_req_spec_resubmitted.request_id, "dataset": dc_req_spec_resubmitted.dataset}
0122         message = "new request resubmitted, destination changed"
0123         if to_submit_idds:
0124             new_request_id = dc_req_spec_resubmitted.request_id
0125             task_id_list = dcif._get_related_tasks(new_request_id)
0126             if task_id_list:
0127                 tmp_logger.debug(f"related tasks: {task_id_list}")
0128                 with ThreadPoolExecutor(max_workers=min(IDDS_SUBMISSION_MAX_WORKERS, len(task_id_list))) as thread_pool:
0129                     future_map = {task_id: thread_pool.submit(dcif._submit_idds_stagein_request, task_id, dc_req_spec_resubmitted) for task_id in task_id_list}
0130                 # the results must be consumed, otherwise a submission raising in its thread is silently lost
0131                 failed_task_id_list = []
0132                 for task_id, future in future_map.items():
0133                     try:
0134                         future.result()
0135                     except Exception as e:
0136                         failed_task_id_list.append(task_id)
0137                         tmp_logger.error(f"failed to submit iDDS request for task_id={task_id} : {e}")
0138                 if failed_task_id_list:
0139                     err_msg = f"submitted iDDS requests for {len(task_id_list) - len(failed_task_id_list)}/{len(task_id_list)} related tasks; failed for {failed_task_id_list}"
0140                     tmp_logger.warning(err_msg)
0141                     message += f"; {err_msg}"
0142                 else:
0143                     tmp_logger.debug(f"submitted corresponding iDDS requests for related tasks")
0144                     message += "; submitted iDDS requests"
0145 
0146             else:
0147                 err_msg = f"failed to get related tasks; skipped to submit iDDS requests"
0148                 tmp_logger.warning(err_msg)
0149                 message += f"; {err_msg}"
0150 
0151     time_delta = naive_utcnow() - time_start
0152     tmp_logger.debug(f"Done. Took {time_delta.seconds}.{time_delta.microseconds // 1000:03d} sec")
0153 
0154     return success, message, data
0155 
0156 
0157 def change_staging_source(
0158     dcif: DataCarouselInterface,
0159     request_id: int | str | None = None,
0160     dataset: str | None = None,
0161     cancel_fts: bool = False,
0162     change_src_expr: bool = False,
0163     source_rse: str | None = None,
0164 ) -> OperationResult:
0165     """
0166     Change source of staging
0167 
0168     If the request is queued, its source_rse will be rechosen, excluding the original source.
0169     If the request is staging, the source_replica_expression of its DDM rule is unset so new source can be tried.
0170     Only effective on queued or staging requests.
0171     The requests can be specified by request_id or dataset (if both exist, request_id is taken).
0172 
0173     Args:
0174         dcif (DataCarouselInterface): Data Carousel interface
0175         request_id (int|str|None): request_id of the staging request, e.g. `123`
0176         dataset (str|None): dataset name of the staging request in the format of Rucio DID
0177         cancel_fts (bool): whether to cancel current FTS requests on DDM, False by default
0178         change_src_expr (bool): whether to change source_replica_expression of the DDM rule by replacing old source with new one, instead of just dropping old source
0179         source_rse (str|None): if set, use this source RSE instead of choosing one randomly, also force change_src_expr to be True; default is None
0180 
0181     Returns:
0182         tuple[bool, str, dict|None]: (success, message, data)
0183     """
0184     tmp_logger = LogWrapper(
0185         _logger,
0186         f"change_staging_source request_id={request_id} dataset={dataset} cancel_fts={cancel_fts} change_src_expr={change_src_expr} source_rse={source_rse}",
0187     )
0188     tmp_logger.debug("Start")
0189     success, message, data = False, "", None
0190     time_start = naive_utcnow()
0191 
0192     dc_req_spec = _resolve_request(dcif, request_id, dataset)
0193 
0194     if dc_req_spec is not None:
0195         status = dc_req_spec.status
0196         orig_source_rse = dc_req_spec.source_rse
0197         if status not in [DataCarouselRequestStatus.queued, DataCarouselRequestStatus.staging]:
0198             err_msg = f"request_id={dc_req_spec.request_id} status={status} not queued or staging; skipped"
0199             tmp_logger.warning(err_msg)
0200             success, message = False, err_msg
0201         else:
0202             ret, dc_req_spec, err_msg = dcif.change_request_source_rse(dc_req_spec, cancel_fts, change_src_expr, source_rse)
0203             if not ret:
0204                 err_msg = f"failed to change source request_id={dc_req_spec.request_id} : {err_msg}"
0205                 tmp_logger.error(err_msg)
0206                 success, message = False, err_msg
0207             else:
0208                 success = True
0209                 if dc_req_spec.status == DataCarouselRequestStatus.queued or change_src_expr:
0210                     message = f"status={status} changed source_rse from {orig_source_rse} to {dc_req_spec.source_rse}"
0211                 else:
0212                     message = f"status={status} source replica expression is dropped"
0213                 data = {
0214                     "request_id": dc_req_spec.request_id,
0215                     "dataset": dc_req_spec.dataset,
0216                     "source_rse": dc_req_spec.source_rse,
0217                     "ddm_rule_id": dc_req_spec.ddm_rule_id,
0218                 }
0219     else:
0220         err_msg = f"failed to get corresponding request"
0221         tmp_logger.error(err_msg)
0222         success, message = False, err_msg
0223 
0224     time_delta = naive_utcnow() - time_start
0225     tmp_logger.debug(f"Done. Took {time_delta.seconds}.{time_delta.microseconds // 1000:03d} sec")
0226 
0227     return success, message, data
0228 
0229 
0230 def force_to_staging(dcif: DataCarouselInterface, request_id: int | str | None = None, dataset: str | None = None) -> OperationResult:
0231     """
0232     Force to staging
0233 
0234     The request will skip the queue and go to staging immediately (will submit DDM rules).
0235     Only effective on queued requests.
0236     The requests can be specified by request_id or dataset (if both exist, request_id is taken).
0237 
0238     Args:
0239         dcif (DataCarouselInterface): Data Carousel interface
0240         request_id (int|str|None): request_id of the staging request, e.g. `123`
0241         dataset (str|None): dataset name of the staging request in the format of Rucio DID
0242 
0243     Returns:
0244         tuple[bool, str, dict|None]: (success, message, data)
0245     """
0246     tmp_logger = LogWrapper(_logger, f"force_to_staging request_id={request_id} dataset={dataset}")
0247     tmp_logger.debug("Start")
0248     success, message, data = False, "", None
0249     time_start = naive_utcnow()
0250 
0251     dc_req_spec = _resolve_request(dcif, request_id, dataset)
0252 
0253     if dc_req_spec is not None:
0254         is_ok, err_msg, dc_req_spec = dcif.stage_request(dc_req_spec)
0255         if not is_ok:
0256             err_msg = f"failed to stage request_id={dc_req_spec.request_id} : {err_msg}"
0257             tmp_logger.error(err_msg)
0258             success, message = False, err_msg
0259         else:
0260             success = True
0261             message = f"status has become {dc_req_spec.status}"
0262             data = {
0263                 "request_id": dc_req_spec.request_id,
0264                 "dataset": dc_req_spec.dataset,
0265                 "status": dc_req_spec.status,
0266                 "ddm_rule_id": dc_req_spec.ddm_rule_id,
0267             }
0268     else:
0269         err_msg = f"failed to get corresponding request"
0270         tmp_logger.error(err_msg)
0271         success, message = False, err_msg
0272 
0273     time_delta = naive_utcnow() - time_start
0274     tmp_logger.debug(f"Done. Took {time_delta.seconds}.{time_delta.microseconds // 1000:03d} sec")
0275 
0276     return success, message, data
0277 
0278 
0279 def retire_unused(dcif: DataCarouselInterface, request_id: int | str | None = None, dataset: str | None = None) -> OperationResult:
0280     """
0281     Retire unused staging request
0282 
0283     If the request is done and has no related tasks, it can be retired to clean up the DDM rules and replicas.
0284     The requests can be specified by request_id or dataset (if both exist, request_id is taken).
0285 
0286     Args:
0287         dcif (DataCarouselInterface): Data Carousel interface
0288         request_id (int|str|None): request_id of the staging request, e.g. `123`
0289         dataset (str|None): dataset name of the staging request in the format of Rucio DID
0290 
0291     Returns:
0292         tuple[bool, str, dict|None]: (success, message, data)
0293     """
0294     tmp_logger = LogWrapper(_logger, f"retire_unused request_id={request_id} dataset={dataset}")
0295     tmp_logger.debug("Start")
0296     success, message, data = False, "", None
0297     time_start = naive_utcnow()
0298 
0299     dc_req_spec = _resolve_request(dcif, request_id, dataset)
0300 
0301     if dc_req_spec is not None:
0302         is_ok, dc_req_spec, err_msg = dcif.retire_unused_request(dc_req_spec)
0303         if not is_ok:
0304             err_msg = f"failed to retire request_id={dc_req_spec.request_id} : {err_msg}"
0305             tmp_logger.error(err_msg)
0306             success, message = False, err_msg
0307         else:
0308             success = True
0309             message = f"retired successfully"
0310             data = {
0311                 "request_id": dc_req_spec.request_id,
0312                 "dataset": dc_req_spec.dataset,
0313                 "status": dc_req_spec.status,
0314                 "ddm_rule_id": dc_req_spec.ddm_rule_id,
0315             }
0316     else:
0317         err_msg = f"failed to get corresponding request"
0318         tmp_logger.error(err_msg)
0319         success, message = False, err_msg
0320 
0321     time_delta = naive_utcnow() - time_start
0322     tmp_logger.debug(f"Done. Took {time_delta.seconds}.{time_delta.microseconds // 1000:03d} sec")
0323 
0324     return success, message, data
0325 
0326 
0327 # operations addressable by name, used by the asynchronous handlers to dispatch on request_type
0328 OPERATIONS = {
0329     "change_staging_destination": change_staging_destination,
0330     "change_staging_source": change_staging_source,
0331     "force_to_staging": force_to_staging,
0332     "retire_unused": retire_unused,
0333 }