107 lines
3.9 KiB
Matlab
107 lines
3.9 KiB
Matlab
function [results, wh] = submitJobs(run_ids, dsp_options, submit_mode, submit_options)
|
|
% SUBMITJOBS Submits dsp_runid jobs for processing (parallel or serial)
|
|
%
|
|
% [results, wh] = submitJobs(run_ids, dsp_options, submit_mode, submit_options)
|
|
%
|
|
% run_ids : scalar or vector of run IDs
|
|
% dsp_options : struct of dsp_runid name/value options
|
|
% submit_mode : processingMode.parallel or .serial
|
|
% submit_options : struct with fields
|
|
% .waitbar (logical)
|
|
% .wh (DataStorage object)
|
|
%
|
|
% results : cell(nJobsPerRunId, nRunIds)
|
|
% wh : updated DataStorage
|
|
|
|
arguments
|
|
run_ids int32 = 0
|
|
dsp_options struct = struct()
|
|
submit_mode processingMode = processingMode.serial
|
|
submit_options.waitbar (1,1) logical = true
|
|
submit_options.wh = DataStorage(struct())
|
|
end
|
|
|
|
% Normalize
|
|
run_ids = run_ids(:)';
|
|
nRunIds = numel(run_ids);
|
|
nJobsPerRunId = submit_options.wh.getLastLinIndice();
|
|
totalJobs = nRunIds * nJobsPerRunId;
|
|
|
|
wh = submit_options.wh;
|
|
results = cell(nJobsPerRunId, nRunIds);
|
|
jobs = repmat(struct('args', {{}}, 'label', "", 'meta', struct()), 1, totalJobs);
|
|
|
|
jobCounter = 0;
|
|
for r = 1:nRunIds
|
|
for k = 1:nJobsPerRunId
|
|
jobCounter = jobCounter + 1;
|
|
optionalVars = buildOptionalVars(k, submit_options.wh);
|
|
|
|
jobs(jobCounter).args = { ...
|
|
run_ids(r), ...
|
|
"database_type", dsp_options.database_type, ...
|
|
"dataBase", dsp_options.dataBase, ...
|
|
"append_to_db", dsp_options.append_to_db, ...
|
|
"load_file_path", dsp_options.load_file_path, ...
|
|
"max_occurences", dsp_options.max_occurences, ...
|
|
"storage_path", dsp_options.storage_path, ...
|
|
"mode", dsp_options.mode, ...
|
|
"parameters", optionalVars};
|
|
jobs(jobCounter).label = sprintf('RunID %d, Job %d', run_ids(r), k);
|
|
jobs(jobCounter).meta.runIndex = r;
|
|
jobs(jobCounter).meta.jobIndex = k;
|
|
jobs(jobCounter).meta.run_id = run_ids(r);
|
|
end
|
|
end
|
|
|
|
linearResults = runBatch(@dsp_runid, jobs, ...
|
|
"mode", submit_mode, ...
|
|
"waitbar", submit_options.waitbar, ...
|
|
"waitbarMessage", "Processing Jobs...", ...
|
|
"numWorkers", 11, ...
|
|
"idleTimeout", 300, ...
|
|
"cancelExistingQueue", true, ...
|
|
"resultHandler", @storeResult, ...
|
|
"errorHandler", @handleError);
|
|
|
|
for idx = 1:totalJobs
|
|
r = jobs(idx).meta.runIndex;
|
|
k = jobs(idx).meta.jobIndex;
|
|
results{k, r} = linearResults{idx};
|
|
end
|
|
|
|
%% Local helpers %%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
|
|
|
|
function handleError(ME, job, ~)
|
|
fprintf('[RunID %d, Job %d] ERROR [%s]: %s\n', job.meta.run_id, job.meta.jobIndex, ...
|
|
ME.identifier, ME.message);
|
|
for st = ME.stack'
|
|
fprintf(' %s:%d (%s)\n', st.file, st.line, st.name);
|
|
end
|
|
fprintf('Full report:\n%s\n', getReport(ME,'extended'));
|
|
end
|
|
|
|
function optionalVars = buildOptionalVars(jobIndex, wh)
|
|
optionalVars = struct();
|
|
if ~isempty(wh.getDimension())
|
|
[vals, names] = wh.getPhysIndicesByLinIndex(jobIndex);
|
|
for pi = 1:numel(names)
|
|
optionalVars.(names{pi}) = vals{pi};
|
|
end
|
|
end
|
|
end
|
|
|
|
function storeResult(val, job, ~)
|
|
if ~isempty(wh)
|
|
jobIndex = job.meta.jobIndex;
|
|
wh.addValueToStorageByLinIdx(val.ffe_package, 'ffe_package', jobIndex);
|
|
wh.addValueToStorageByLinIdx(val.mlse_package, 'mlse_package', jobIndex);
|
|
wh.addValueToStorageByLinIdx(val.vnle_package, 'vnle_package', jobIndex);
|
|
wh.addValueToStorageByLinIdx(val.dbtgt_package,'dbtgt_package',jobIndex);
|
|
wh.addValueToStorageByLinIdx(val.dbenc_package,'dbenc_package',jobIndex);
|
|
wh.addValueToStorageByLinIdx(val.mlmlse_package,'mlmlse_package',jobIndex);
|
|
end
|
|
end
|
|
|
|
end
|