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
0025
0026 _logger = PandaLogger().getLogger("api_data_carousel")
0027
0028
0029 OperationResult = tuple[bool, str, dict | None]
0030
0031
0032
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
0075 return dcif.get_request_by_id(int(request_id))
0076 elif dataset is not None:
0077
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
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
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 }