Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -70,11 +70,11 @@ In this section all the attributes that can be used in the DIRAC JDL job descrip
+---------------------+---------------------------------------------+-------------------------------------------------------------------------------------+
| *InputDataPolicy* | Job input data policy | InputDataPolicy = ``"DIRAC.WorkloadManagementSystem.Client.DownloadInputData";`` |
+---------------------+---------------------------------------------+-------------------------------------------------------------------------------------+
| *OutputData* | Job output data files | OutputData = ``{"output1","output2"};`` |
| *OutputData* [1] | Job output data files | OutputData = ``{"output1","output2"};`` |
+---------------------+---------------------------------------------+-------------------------------------------------------------------------------------+
| *OutputPath* | The output data path in the File Catalog | OutputPath = ``{"/myjobs/output"};`` |
| *OutputPath* [2] | The output data path in the File Catalog | OutputPath = ``{"/myjobs/output"};`` |
+---------------------+---------------------------------------------+-------------------------------------------------------------------------------------+
| *OutputSE* | The output data Storage Element | OutputSE = ``{"DIRAC-USER"};`` |
| *OutputSE* [3] | The output data Storage Element | OutputSE = ``{"DIRAC-USER"};`` |
+---------------------+---------------------------------------------+-------------------------------------------------------------------------------------+
| |
| :subtitle:`Parametric Jobs` |
Expand All @@ -91,3 +91,27 @@ In this section all the attributes that can be used in the DIRAC JDL job descrip
+---------------------+---------------------------------------------+-------------------------------------------------------------------------------------+
| *ParameterFactor* | Parameter multiplier | ParameterFactor = 1.1; (default 1.) |
+---------------------+---------------------------------------------+-------------------------------------------------------------------------------------+

1. Elements of OutputData can be specified in several forms:

- filenames; in this case files with the specified names will be looked for in the job directory and uploaded
to a location specified by the OutputPath (see below);
- filenames with wild cards, e.g. ``"*.log"`` ; same after the filenames expansion;
- output data specified in a form ``"LFN:/vo/full/destination/path/filename"``; in this case the file ``"filename"``
in the job directory will be uploaded to the specified LFN path without taking into account the OutputPath.
Note that "filename" here can be also specified with wild cards, e.g. ``"LFN:/vo/full/destination/path/*.log"`` .

2. The OutputPath can be specified in several ways

- if not given, it will be taken as the user's home directory + the job directory
for example ``"/lhcb/user/a/atsareg/1234/1234567"``, where 1234567 is the job ID;
- if given as a path starting with "/", it will be appended to the user's home
directory, e.g. outputPath = ``"/my/analysis"`` will make output files to go to the
``"/lhcb/user/a/atsareg/my/analysis"`` directory
- if given as ``"LFN:/output/path"``, it will be taken as an absolute path for
output files in the logical namespace. It is the responsibility of the user to make
sure that this path is accessible for writing for the user's data.

3. If multiple output SEs are specified, they will be tried one-by-one for redundancy for each
output file until the first successful file upload. No more than one replica will be created
for each file.
23 changes: 21 additions & 2 deletions src/DIRAC/WorkloadManagementSystem/JobWrapper/JobWrapper.py
Original file line number Diff line number Diff line change
Expand Up @@ -1056,10 +1056,25 @@ def __transferOutputDataFiles(self, outputData, outputSE, outputPath):
else:
nonlfnList.append(out)

# Check whether list of outputData has a globbable pattern
# Check whether the list of LFNs has globbable patterns
globbedLfnList = []
for lfn in lfnList:
lfnPath = os.path.dirname(lfn)
lfnLocal = os.path.basename(lfn)
globbedLfnList += [os.path.join(lfnPath, gLfn) for gLfn in getGlobbedFiles(lfnLocal)]

if globbedLfnList:
globbedLfnList = List.uniqueElements(globbedLfnList)
if globbedLfnList != lfnList:
self.log.info(
"Found a pattern in the output data LFN list, LFNs to upload are:", ", ".join(globbedLfnList)
)
lfnList = globbedLfnList

# Check whether the list of outputData has a globbable pattern
nonlfnList = [str(self.jobIDPath / x) for x in nonlfnList]
globbedOutputList = List.uniqueElements(getGlobbedFiles(nonlfnList))
if globbedOutputList != nonlfnList and globbedOutputList:
if globbedOutputList and globbedOutputList != nonlfnList:
self.log.info(
"Found a pattern in the output data file list, files to upload are:", ", ".join(globbedOutputList)
)
Expand Down Expand Up @@ -1218,6 +1233,10 @@ def __getLFNfromOutputFile(self, outputFile, outputPath=""):
# If output path is given, append it to the user path and put output files in this directory
if outputPath.startswith("/"):
outputPath = outputPath[1:]
# If output path is given with the LFN: prefix, take it as an absolute path
elif outputPath.startswith("LFN:"):
outputPath = outputPath[4:]
basePath = Path("")
else:
# By default the output path is constructed from the job id
subdir = str(int(self.jobID / 1000))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,15 @@

gLogger.setLevel("DEBUG")

uploadSandboxMock = MagicMock()
uploadSandboxMock.return_value = {"OK": True}


def uploadFileMockFunc(**kwargs):
destinationSEList = kwargs["destinationSEList"]
return {"OK": True, "Value": {"uploadedSE": destinationSEList[0]}}


# -------------------------------------------------------------------------------------------------


Expand Down Expand Up @@ -661,11 +670,13 @@ def jobIDPath():
(p / "std.err").touch()
(p / "summary_123.xml").touch()
(p / "result_dir").mkdir()
(p / "result_dir" / "file1").touch()
# Output data files
(p / "00232454_00000244.xml").touch()
(p / "00232454_00000244_1.sim").touch()
(p / "1720442808testFileUpload.txt").touch()
(p / "testFileUploadFullLFN.txt").touch()
(p / "result_dir" / "output.xml").touch()
(p / "result_dir" / "output.txt").touch()

with open(p / "pool_xml_catalog.xml", "w") as f:
f.write(
Expand Down Expand Up @@ -996,3 +1007,116 @@ def test_finalize(mocker, failedFlag, expectedRes, finalStates):
assert res == expectedRes
assert jw.jobReport.jobStatusInfo[0][0] == finalStates[0]
assert jw.jobReport.jobStatusInfo[0][1] == finalStates[1]


@pytest.mark.parametrize(
"outputData, outputPath, expectedResult",
[
(
"00232454_00000244.xml",
None,
"/dirac/user/u/unknown/0/123/00232454_00000244.xml",
),
(
"00232454_00000244*",
None,
[
"/dirac/user/u/unknown/0/123/00232454_00000244.xml, "
"/dirac/user/u/unknown/0/123/00232454_00000244_1.sim",
"/dirac/user/u/unknown/0/123/00232454_00000244_1.sim, "
"/dirac/user/u/unknown/0/123/00232454_00000244.xml",
],
),
(
"*.txt",
None,
[
"/dirac/user/u/unknown/0/123/1720442808testFileUpload.txt, "
"/dirac/user/u/unknown/0/123/testFileUploadFullLFN.txt",
"/dirac/user/u/unknown/0/123/testFileUploadFullLFN.txt, "
"/dirac/user/u/unknown/0/123/1720442808testFileUpload.txt",
],
),
(
"00232454_00000244.xml",
"/my_output_dir/00232454",
"/dirac/user/u/unknown/my_output_dir/00232454/00232454_00000244.xml",
),
(
"00232454_00000244.xml",
"LFN:/dirac/prod/00232454",
"/dirac/prod/00232454/00232454_00000244.xml",
),
(
"LFN:/dirac/prod/00232454/00232454_00000244.xml",
None,
"/dirac/prod/00232454/00232454_00000244.xml",
),
(
"LFN:/dirac/prod/00232454/00232454_00000244.xml",
"/my_output_dir/00232454",
"/dirac/prod/00232454/00232454_00000244.xml",
),
(
"result_dir",
None,
[
"/dirac/user/u/unknown/0/123/output.xml, /dirac/user/u/unknown/0/123/output.txt",
"/dirac/user/u/unknown/0/123/output.txt, /dirac/user/u/unknown/0/123/output.xml",
],
),
(
"result_dir/*.xml",
None,
"/dirac/user/u/unknown/0/123/output.xml",
),
(
"result_dir/*.xml",
"/my_output_dir/00232454",
"/dirac/user/u/unknown/my_output_dir/00232454/output.xml",
),
],
)
def test_OutputData(mocker, jobIDPath, outputData, outputPath, expectedResult):
mocker.patch(
"DIRAC.WorkloadManagementSystem.JobWrapper.JobWrapper.gConfig.getValue",
side_effect=lambda path, default=None: default,
)

mocker.patch("DIRAC.WorkloadManagementSystem.JobWrapper.JobWrapper.getVOForGroup", return_value="dirac")

ops_mock = MagicMock()

ops_mock.getValue.return_value = "user"

mocker.patch("DIRAC.WorkloadManagementSystem.JobWrapper.JobWrapper.Operations", return_value=ops_mock)
mocker.patch(
"DIRAC.DataManagementSystem.Client.FailoverTransfer.FailoverTransfer.transferAndRegisterFile",
side_effect=uploadFileMockFunc,
)
mocker.patch(
"DIRAC.WorkloadManagementSystem.Client.SandboxStoreClient.SandboxStoreClient.uploadFilesAsSandbox",
side_effect=uploadSandboxMock,
)

jw = JobWrapper(jobIDPath)
jw.jobIDPath = Path.cwd() / str(jobIDPath)
os.chdir(str(jw.jobID))
jw.jobArgs = {
"OutputData": outputData,
"OutputPath": outputPath,
"Owner": "duser",
"OutputSE": "DIRAC-disk",
"OutputSandbox": ["std.out", "std.err"],
}

jw.failedFlag = False
jw.dm = dm_mock
jw.fc = fc_mock

result = jw.processJobOutputs()
os.chdir(jw.root)
assert result["OK"]
uploaded = [v for k, v in jw.jobReport.jobParameters if k == "UploadedOutputData"]
assert uploaded, f"UploadedOutputData parameter not set: {jw.jobReport.jobParameters}"
assert uploaded[0] in expectedResult
Loading