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); jobOptions = dsp_options; jobOptions.parameters = optionalVars; jobs(jobCounter).args = [{run_ids(r)}, structToNameValue(jobOptions)]; 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 nameValue = structToNameValue(options) names = fieldnames(options); nameValue = cell(1, 2*numel(names)); for i = 1:numel(names) nameValue{2*i - 1} = string(names{i}); nameValue{2*i} = options.(names{i}); 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