5#ifndef LINE_IO_WFCOMMONS_LOADER_H
6#define LINE_IO_WFCOMMONS_LOADER_H
87 for (
char& c : l) c =
static_cast<char>(std::tolower(
static_cast<unsigned char>(c)));
111 return v ==
"1.3" || v ==
"1.4" || v ==
"1.5";
114inline bool has(
const json& o,
const char* key) {
return o.is_object() && o.contains(key); }
118 if (v.is_null())
return false;
119 if (v.is_array() || v.is_object() || v.is_string())
return !v.empty();
124inline std::string
as_text(
const json& v) {
return v.is_string() ? v.get<std::string>() : v.dump(); }
128 if (!
has(data,
"schemaVersion"))
129 throw InputError(
"WfCommonsLoader: Missing schemaVersion field in WfCommons JSON.");
130 const std::string version =
as_text(data.at(
"schemaVersion"));
132 std::cerr <<
"[LINE] Warning: WfCommonsLoader: Schema version " << version
133 <<
" may not be fully supported." << std::endl;
134 if (!
has(data,
"workflow"))
135 throw InputError(
"WfCommonsLoader: Missing workflow field in WfCommons JSON.");
136 const json&
wf = data.at(
"workflow");
137 bool has_tasks =
false;
138 if (
has(
wf,
"specification")) {
139 const json& spec =
wf.at(
"specification");
140 has_tasks =
has(spec,
"tasks") &&
nonempty(spec.at(
"tasks"));
144 if (!has_tasks)
throw InputError(
"WfCommonsLoader: Workflow must have at least one task.");
149 const std::size_t slash = s.find_last_of(
"/\\");
150 std::string f = slash == std::string::npos ? s : s.substr(slash + 1);
151 const std::size_t dot = f.find_last_of(
'.');
152 if (dot != std::string::npos) f = f.substr(0, dot);
159 if (
has(data,
"name") && data.at(
"name").is_string() && !data.at(
"name").get<std::string>().empty())
160 name = data.at(
"name").get<std::string>();
164 if (!(std::isalnum(
static_cast<unsigned char>(c)) || c ==
'_')) c =
'_';
165 if (name.empty()) name =
"Workflow";
171 if (
has(task,
"id"))
return as_text(task.at(
"id"));
172 if (
has(task,
"name"))
return as_text(task.at(
"name"));
173 return "task_" + std::to_string(idx1);
182 switch (
opt.distribution_type) {
188 if (
opt.default_scv > 1.0)
190 return D::exp_mean(m);
193 return D::exp_mean(m);
200 std::map<std::string, std::string> md;
201 md[
"taskId"] =
json(
id).dump();
202 static const char* kTask[] = {
"name",
"inputFiles",
"outputFiles"};
203 for (
const char* f : kTask)
204 if (
has(task, f)) md[f] = task.at(f).dump();
206 static const char* kExec[] = {
"executedAt",
"command",
"coreCount",
"avgCPU",
207 "readBytes",
"writtenBytes",
"memoryInBytes",
208 "energyInKWh",
"avgPowerInW",
"priority",
"machines"};
209 for (
const char* f : kExec)
210 if (
has(*exec, f)) md[f] = exec->at(f).dump();
217 const std::vector<std::vector<std::size_t>>& adj) {
218 std::vector<bool> visited(adj.size(),
false);
219 std::deque<std::size_t> q(1, start);
220 visited[start] =
true;
221 std::vector<std::size_t> out;
223 const std::size_t cur = q.front();
225 for (std::size_t nx : adj[cur])
242 const std::vector<std::vector<std::size_t>>& adj,
243 const std::vector<std::size_t>& indeg) {
244 const std::size_t none = adj.size();
245 if (children.empty())
return none;
247 std::sort(common.begin(), common.end());
248 common.erase(std::unique(common.begin(), common.end()), common.end());
249 for (std::size_t c = 1; c < children.size(); ++c) {
251 std::sort(r.begin(), r.end());
252 std::vector<std::size_t> both;
253 std::set_intersection(common.begin(), common.end(), r.begin(), r.end(),
254 std::back_inserter(both));
257 for (std::size_t node : common) {
258 if (indeg[node] < children.size())
continue;
259 bool all_direct =
true;
260 for (std::size_t ch : children)
261 if (std::find(adj[ch].begin(), adj[ch].end(), node) == adj[ch].end()) {
265 if (all_direct)
return node;
276 const json& jwf = data.at(
"workflow");
277 const bool legacy = !
has(jwf,
"specification");
278 const json& tasks = legacy ? jwf.at(
"tasks") : jwf.at(
"specification").at(
"tasks");
279 if (!tasks.is_array())
280 throw InputError(
"WfCommonsLoader: the workflow's 'tasks' must be an array of task objects");
281 const std::size_t n = tasks.size();
283 std::vector<std::string> ids(n);
284 for (std::size_t i = 0; i < n; ++i) ids[i] =
task_id(tasks[i], i + 1);
287 std::map<std::string, const json*> exec_map;
288 if (
opt.use_execution_data) {
290 for (std::size_t i = 0; i < n; ++i) exec_map[ids[i]] = &tasks[i];
291 }
else if (
has(jwf,
"execution") &&
has(jwf.at(
"execution"),
"tasks")) {
292 for (
const json& et : jwf.at(
"execution").at(
"tasks"))
293 if (
has(et,
"id")) exec_map[
as_text(et.at(
"id"))] = &et;
298 std::map<std::string, std::size_t> idx_of;
299 for (std::size_t i = 0; i < n; ++i) {
300 const std::string&
id = ids[i];
301 double runtime =
opt.default_runtime;
302 const auto it = exec_map.find(
id);
303 const json* exec = it == exec_map.end() ? nullptr : it->second;
304 if (exec &&
has(*exec,
"runtimeInSeconds") && exec->at(
"runtimeInSeconds").is_number())
305 runtime = exec->at(
"runtimeInSeconds").get<
double>();
312 std::map<std::string, std::size_t> name_idx;
313 std::map<std::string, int> name_count;
314 if (
opt.resolve_children_by_name)
315 for (std::size_t i = 0; i < n; ++i)
316 if (
has(tasks[i],
"name") && tasks[i].at(
"name").is_string()) {
317 const std::string nm = tasks[i].at(
"name").get<std::string>();
324 std::vector<std::vector<std::size_t>> adj(n);
325 std::vector<std::size_t> indeg(n, 0), outdeg(n, 0);
326 std::vector<std::vector<bool>> seen(n, std::vector<bool>(n,
false));
327 std::vector<std::string> unresolved;
328 for (
int pass = 0; pass < 2; ++pass) {
329 const char* field = pass == 0 ?
"children" :
"parents";
330 for (std::size_t i = 0; i < n; ++i) {
331 if (!
has(tasks[i], field))
continue;
332 const json& rj = tasks[i].at(field);
333 std::vector<std::string> refs;
334 if (rj.is_string()) refs.push_back(rj.get<std::string>());
335 else if (rj.is_array())
336 for (
const json& c : rj) refs.push_back(
as_text(c));
337 for (
const std::string& k : refs) {
339 const auto a = idx_of.find(k);
340 if (a != idx_of.end()) j = a->second;
341 else if (
opt.resolve_children_by_name) {
342 const auto b = name_idx.find(k);
343 if (b != name_idx.end() && name_count[k] == 1) j = b->second;
346 unresolved.push_back(k);
349 const std::size_t from = pass == 0 ? i : j, to = pass == 0 ? j : i;
350 if (seen[from][to])
continue;
351 seen[from][to] =
true;
352 adj[from].push_back(to);
358 if (!unresolved.empty()) {
359 std::vector<std::string> shown;
360 for (
const std::string& u : unresolved)
361 if (std::find(shown.begin(), shown.end(), u) == shown.end()) shown.push_back(u);
362 std::cerr <<
"[LINE] Warning: WfCommonsLoader: " << unresolved.size()
363 <<
" child/parent reference(s) match no task id"
364 << (
opt.resolve_children_by_name ?
" or unique task name" :
"") <<
" and were dropped: ";
365 for (std::size_t u = 0; u < shown.size() && u < 5; ++u) std::cerr << (u ?
", " :
"") << shown[u];
366 std::cerr << (shown.size() > 5 ?
", ..." :
"") <<
"\n";
370 std::vector<std::vector<bool>> processed(n, std::vector<bool>(n,
false));
371 for (std::size_t i = 0; i < n; ++i) {
372 if (outdeg[i] <= 1)
continue;
373 const std::vector<std::size_t>& children = adj[i];
375 if (join == n)
continue;
376 std::vector<std::string> posts;
377 for (std::size_t c : children) posts.push_back(ids[c]);
378 wf.add_precedence(W::AndFork(ids[i], posts));
379 wf.add_precedence(W::AndJoin(posts, ids[join]));
380 for (std::size_t c : children) {
381 processed[i][c] =
true;
382 processed[c][join] =
true;
385 for (std::size_t i = 0; i < n; ++i)
386 for (std::size_t j : adj[i])
387 if (!processed[i][j])
wf.add_precedence(W::Serial(ids[i], ids[j]));
393 std::ifstream in(path.c_str());
394 if (!in)
throw InputError(
"WfCommonsLoader: cannot open '" + path +
"'");
399 }
catch (
const json::parse_error& e) {
400 throw InputError(
"WfCommonsLoader: malformed JSON in '" + path +
"': " + e.what());
428namespace wfcommons_detail {
432 const char*
env = std::getenv(
"PATH");
433 const std::string path =
env ?
env :
"";
435 while (b <= path.size()) {
436 const std::size_t e = path.find(
':', b);
437 const std::string dir = path.substr(b, e == std::string::npos ? std::string::npos : e - b);
438 if (!dir.empty() && ::access((dir +
"/" + name).c_str(), X_OK) == 0)
return dir +
"/" + name;
439 if (e == std::string::npos)
break;
442 return std::string();
449inline std::string
fetch_https(
const std::string& url,
int timeout_ms) {
450 const int secs = timeout_ms > 0 ? (timeout_ms + 999) / 1000 : 0;
452 const std::string dest = dir.
file(
"workflow.json");
453 std::vector<std::string> argv;
456 argv = {curl,
"-sS",
"-fL",
"--proto",
"=https",
"--proto-redir",
"=https",
"--tlsv1.2",
"-o", dest};
458 argv.push_back(
"--max-time");
459 argv.push_back(std::to_string(secs));
464 throw UnsupportedError(
"WfCommonsLoader: " + url +
" is https, which this port fetches with curl or "
465 "wget, and neither is on PATH. Install one, or download the file and read it "
466 "with wfcommons_load");
467 argv = {wget,
"-q",
"--https-only",
"--tries=1",
"-O", dest};
468 if (secs > 0) argv.push_back(
"--timeout=" + std::to_string(secs));
472 if (r.
timedOut)
throw InputError(
"WfCommonsLoader: GET " + url +
" did not finish within " +
473 std::to_string(secs) +
" s");
475 throw InputError(
"WfCommonsLoader: GET " + url +
" failed (" + argv[0] +
" exit " +
477 std::ifstream in(dest.c_str(), std::ios::binary);
478 std::stringstream ss;
497 int timeout_ms = 60000) {
499 if (url.compare(0, 8,
"https://") == 0) {
504 throw InputError(
"WfCommonsLoader: GET " + url +
" answered HTTP " +
505 std::to_string(resp.
status) +
" (redirects are not followed)");
510 data = nlohmann::json::parse(body);
511 }
catch (
const nlohmann::json::parse_error& e) {
512 throw InputError(
"WfCommonsLoader: malformed JSON from " + url +
": " + e.what());
UnsupportedError(const std::string &what)
std::string file(const std::string &name) const
A file inside it.
The moment fitters the reference distributions carry as STATIC FACTORIES: Erlang.fitMeanAndOrder,...
The exception types the port throws.
Minimal HTTP/1.1 client, enough to talk to a line-*-rest service.
Enumerations and the minimal distribution descriptor shared by the model layer of the C++ port.
Response get(const std::string &url, int timeoutMillis)
GET a URL.
std::string task_id(const json &task, std::size_t idx1)
Port of getTaskId.
std::string fetch_https(const std::string &url, int timeout_ms)
GET an https:// URL with curl, else wget, into memory.
std::string file_stem(const std::string &s)
[~, name, ~] = fileparts(s).
std::map< std::string, std::string > extract_metadata(const json &task, const std::string &id, const json *exec)
Port of extractMetadata; values are the JSON text of each field.
std::string extract_name(const json &data, const std::string &default_name)
Port of extractName: the document's name, else the default's stem, sanitized.
bool nonempty(const json &v)
~isempty(x) on a decoded JSON value.
bool supported_version(const std::string &v)
WfCommonsLoader.SUPPORTED_SCHEMA_VERSIONS.
std::size_t find_common_join(const std::vector< std::size_t > &children, const std::vector< std::vector< std::size_t > > &adj, const std::vector< std::size_t > &indeg)
Port of findCommonJoin.
std::string which_on_path(const std::string &name)
The first executable name on PATH, empty when there is none.
lang::Distrib< T > fit_distribution(double runtime, const WfCommonsOptions &opt)
Port of fitDistribution.
void validate_schema(const json &data)
Port of validateSchema.
workflow::Workflow< T > build_workflow(const json &data, const std::string &wf_name, const WfCommonsOptions &opt)
Port of buildWorkflow (with buildAdjacency and addPrecedences).
std::vector< std::size_t > reachable_nodes(std::size_t start, const std::vector< std::vector< std::size_t > > &adj)
Port of getReachableNodes: BFS order, the start excluded.
bool has(const json &o, const char *key)
std::string as_text(const json &v)
A JSON scalar as the text strcmp compares: strings verbatim, numbers as written.
json parse_file(const std::string &path)
Parse a whole file as JSON, naming the file on failure.
WfDistributionType
options.distributionType, the JAR's WfCommonsOptions.DistributionType.
workflow::Workflow< T > wfcommons_load(const std::string &path, const WfCommonsOptions &options=WfCommonsOptions())
Port of WfCommonsLoader.load(jsonFile, options).
bool wfcommons_validate(const std::string &path)
Port of WfCommonsLoader.validateFile(jsonFile): true when the file parses and passes the schema check...
workflow::Workflow< T > wfcommons_load_from_url(const std::string &url, const WfCommonsOptions &options=WfCommonsOptions(), int timeout_ms=60000)
Port of WfCommonsLoader.loadFromUrl(urlString, options).
WfDistributionType wf_distribution_type_from_string(const std::string &s)
lower(options.distributionType) as the reference switches on it: exp, det, aph, hyperexp.
workflow::Workflow< T > wfcommons_load_from_json(const nlohmann::json &data, const WfCommonsOptions &options=WfCommonsOptions())
Port of WfCommonsLoader.loadFromStruct(data, options), on a parsed document.
Distrib< T > aph_fit_mean_scv(const T &mean, const T &scv)
APH.fitMeanAndSCV(MEAN, SCV), through mam::aph_fit_mean_scv.
Distrib< T > hyperexp_fit_mean_scv(const T &mean, const T &scv)
HyperExp.fitMeanAndSCV(MEAN, SCV), which is map_hyperexp at p = 0.99 read back as (p,...
ProcResult capture(const std::vector< std::string > &argv, int timeoutSeconds, bool mergeStderr=false)
Runs a command, capturing stdout and discarding stderr.
std::string trim(const std::string &s)
Trims ASCII whitespace from both ends, as Java's String.trim() does.
Conservation laws of a layered queueing network, enumerated from its structure.
Number-type abstraction for the templated API port.
An HTTP response, with the body already de-chunked.
std::string body
Response body, decoded.
int status
HTTP status code.
The loader options; the defaults are the reference's parseOptions.
bool resolve_children_by_name
Resolve a child matching no task id by task name (see the file comment).
double default_scv
defaultSCV, for APH and HyperExp
double default_runtime
defaultRuntime, when a task has no runtime
bool use_execution_data
useExecutionData
bool store_metadata
storeMetadata, onto WorkflowActivity::metadata()
WfDistributionType distribution_type
distributionType
static constexpr double FineTol
Outcome of a captured command.
int exitCode
Exit status, or -1 when the command could not run.
bool timedOut
True when the deadline expired and the child was killed.
std::string out
Everything the command wrote to stdout.
Running an external command and capturing its output, with a deadline.
A scratch directory for the subprocess wrappers, the port's lineTempName.
An activity workflow reduced to one phase-type law.