LINE Solver (C++)
Templated C++ port of the LINE queueing solver
Loading...
Searching...
No Matches
lqn2qn.h
Go to the documentation of this file.
1/*
2 * Copyright (c) 2012-2026, QORE Lab, Imperial College London
3 * All rights reserved.
4 */
5#ifndef LINE_IO_LQN2QN_H
6#define LINE_IO_LQN2QN_H
7
8/**
9 * @file
10 * @ingroup line_io
11 * Convert a layered queueing network into a single queueing network in which
12 * synchronous call blocking is carried by REPLY signals.
13 *
14 * Port of matlab/src/io/LQN2QN.m (the reference), mirrored by
15 * jar/src/main/java/jline/io/LQN2QN.java and python/line_solver/io/__init__.py.
16 * The construction is the reference's, step for step, and so is the order in
17 * which nodes and classes are created, so that a model converted here and one
18 * converted by MATLAB or Python have the same node and class tables:
19 *
20 * - One station per host processor replica (a Delay when the processor is an
21 * infinite server), one Delay per reference task replica for its think time.
22 * - One class per STEP of the expanded activity graph: an activity, one call
23 * stage of an activity, or a merge/trigger step. Reference-task steps are
24 * closed classes of population 0, except the think class, which carries the
25 * task multiplicity; open-arrival steps are open classes.
26 * - A synchronous call site that holds its caller's server gets a REPLY signal
27 * bound to its class (sn.syncreply); the callee's replying step switches into
28 * that signal, which returns to the caller station and releases the server.
29 * - A call mean m is unrolled into floor(m) mandatory stages plus one stage
30 * taken with probability m - floor(m), capped at 20 stages.
31 * - OR branches and loops follow the graph weights; AND forks become a Fork plus
32 * one Router per branch and AND joins a Join, with a PARTIAL strategy when the
33 * declared quorum is below the branch count.
34 * - A CacheTask becomes a Cache node whose read step switches into the hit and
35 * the miss class; delayed-hit retrieval adds a PS fetch station per replica.
36 * - Asynchronous calls are non-blocking visits, forwarding splits the reply
37 * exits of the forwarding entry, and the multiplicity of a non-reference task
38 * is a thread pool enforced by one finite capacity region with one linear
39 * admission row per task replica.
40 * - Phase-2 activities run after the reply, spawned by the replying step's
41 * completions (sn.classspawn) and destroyed at the chain end by a NEGATIVE
42 * signal (closed chains) or at the Sink (open chains).
43 * - Activity think times are steps on a shared ActivityThink delay; setup tasks
44 * carry their setup/delay-off pair onto their host station.
45 * - Replication is materialised (one station and one step-graph copy per
46 * replica, calls reaching the fan-out block {(i*f+k) mod r}) or pooled (one
47 * station of r times the servers), selected by 'auto', 'materialize' or
48 * 'pool' exactly as in the reference.
49 *
50 * REPLIES. The reference reads which activity replies to which entry from
51 * `lsn.replygraph`, which holds the EXPLICIT replies plus the implicit ones
52 * getStruct.m infers for leaf activities. The C++ LqnStruct keeps no replygraph,
53 * so the LqnModel overload recovers it with lqn::detail::reply_activities, the
54 * port of that same pass. The LqnStruct overload can only infer the implicit
55 * replies: an entry whose reply is declared on a non-leaf activity, which is
56 * what a phase-2 activity is, needs the LqnModel overload.
57 *
58 * Construction choices that are the reference's own and not yet represented
59 * there either (retrieval on a cache read with phase-2 successors, the thread
60 * pool of a task with an internal AND-fork, and the others listed in LQN2QN.m)
61 * are reported as warnings, exactly where the reference calls line_warning.
62 */
63
64#include <algorithm>
65#include <cmath>
66#include <cstddef>
67#include <iostream>
68#include <limits>
69#include <map>
70#include <set>
71#include <string>
72#include <utility>
73#include <vector>
74
79#include "line/num/number.h"
80#include "line/util/error.h"
81#include "line/util/matrix.h"
82
83namespace line {
84namespace io {
85
86namespace detail {
87
88/** The LQN2QN walk: the MATLAB nested functions as members, their shared state as fields. */
89template <class T>
90class Lqn2Qn {
91 public:
92 typedef std::pair<std::size_t, std::size_t> Key; ///< (element index, 0-based replica)
93 typedef lang::Distrib<T> Dist;
94
95 Lqn2Qn(const lqn::LqnStruct<T>& lsn, const std::set<Key>& replies,
96 const std::string& replication, const std::string& name,
97 std::vector<std::string>* warnings)
98 : l_(lsn), replies_(replies), replication_(replication), warnings_(warnings),
99 net_(name + "-QN") {}
100
101 qn::Network<T> run() {
102 build();
103 return net_;
104 }
105
106 private:
107 static const std::size_t NONE = static_cast<std::size_t>(-1);
108 static const std::size_t MAXCALLSTAGES = 20; // guard against unrolling a huge call mean
109 static const std::size_t MAXREPLINSTANCES = 128; // above this 'auto' pools instead
110
111 /** The port a step is left through; `sig` means through its REPLY signal. */
112 struct Port {
113 std::size_t step;
114 bool sig;
115 Port(std::size_t s = NONE, bool g = false) : step(s), sig(g) {}
116 bool operator==(const Port& o) const { return step == o.step && sig == o.sig; }
117 };
118 struct Exit {
119 std::size_t step;
120 bool sig;
121 T prob;
122 };
123 struct Flow {
124 std::size_t from, to;
125 T p;
126 bool fromSig, inTarget;
127 };
128 struct Reply {
129 std::size_t exit, owner;
130 bool viaSig;
131 T p;
132 };
133 struct Ph2Exit {
134 std::size_t step;
135 bool sig;
136 T p;
137 Key ref;
138 };
139 struct CacheWire {
140 std::size_t node, readStep, hitStep, missStep, eidx, fetch;
141 Dist fetchSvc;
142 };
143 struct Stage {
144 std::size_t target;
145 T prob;
146 bool isasync;
147 };
148 struct Step {
149 std::size_t aidx;
150 Key host;
151 Dist svc; ///< disabled = no service declared by the step
152 std::string name;
153 bool blocks, isthink;
154 Key ref; ///< (reference task, replica); (0,0) on an open chain
155 std::size_t node; ///< Fork/Join/Router/Cache/ActivityThink node, 0 for stations
156 std::size_t owner;
157 std::vector<Key> tasks; ///< thread-pool task replicas holding a thread here
158 };
159 struct Expansion {
160 std::size_t first = NONE;
161 std::vector<Exit> replies, terms;
162 };
163 /** State of one expandActivities call, the closure of the MATLAB walk. */
164 struct WalkCtx {
165 std::size_t eidx, tidx, trep;
166 Key ref;
167 std::vector<std::size_t> localActs;
168 std::map<std::size_t, std::pair<std::size_t, Port> > visited;
169 std::map<std::size_t, std::size_t> joinOf;
170 std::map<std::size_t, std::size_t> joinFork; ///< join activity -> fork step it closes
171 std::vector<std::size_t> forkOwnerStack;
172 bool sawReply = false;
173 std::vector<Exit> replyExits, terminals;
174 };
175
176 const lqn::LqnStruct<T>& l_;
177 const std::set<Key>& replies_;
178 std::string replication_;
179 std::vector<std::string>* warnings_;
180 qn::Network<T> net_;
181
182 std::vector<double> replRaw_;
183 bool materialize_ = false;
184 std::vector<bool> fcrTask_;
185 std::map<std::size_t, std::size_t> hostNTasks_; // host element -> number of tasks it runs
186 std::vector<std::size_t> refTasks_, openEntries_;
187 std::map<Key, std::size_t> hostStation_, thinkNode_, cacheNodeOf_, fetchNodeOf_;
188 std::map<Key, bool> hostIsDelay_;
189 std::vector<Step> steps_;
190 std::vector<Flow> flow_;
191 std::vector<Reply> reply_;
192 std::vector<std::pair<std::size_t, std::size_t> > spawnPairs_, joinQuorum_;
193 std::vector<Ph2Exit> ph2Exits_;
194 std::vector<Key> entryStack_, threadStack_;
195 std::vector<CacheWire> cacheWiring_;
196 std::map<std::string, int> usedNames_;
197 std::size_t actThinkNode_ = 0;
198 std::size_t srcNode_ = 0, snkNode_ = 0;
199
200 // ------------------------------------------------------------ small helpers
201
202 static T tnum(double v) { return num_traits<T>::from_double(v); }
203 static double dbl(const T& v) { return num_traits<T>::to_double(v); }
204
205 void warn(const std::string& msg) {
206 if (warnings_ != nullptr)
207 warnings_->push_back("LQN2QN: " + msg);
208 else
209 std::cerr << "[LINE] Warning: LQN2QN: " << msg << std::endl;
210 }
211
212 /** A distribution that is declared, not Immediate and above tolerance. */
213 static bool timed(const Dist& d) {
214 return !d.disabled && !d.is_immediate() &&
215 dbl(d.mean) > lang::GlobalConstants::FineTol;
216 }
217
218 std::size_t parent(std::size_t idx) const { return l_.parent[idx]; }
219 const std::string& nm(std::size_t idx) const { return l_.names[idx]; }
220
221 bool isAndJoinPre(std::size_t aidx) const {
222 return aidx < l_.actpretype.size() && l_.actpretype[aidx] == lang::PrecedenceType::PRE_AND;
223 }
224 bool isPostAnd(std::size_t aidx) const {
225 return aidx < l_.actposttype.size() &&
226 l_.actposttype[aidx] == lang::PrecedenceType::POST_AND;
227 }
228 bool isAndFork(const std::vector<std::size_t>& succ) const {
229 if (succ.size() < 2) return false;
230 for (std::size_t s : succ)
231 if (!isPostAnd(s)) return false;
232 return true;
233 }
234 bool repliesTo(std::size_t aidx, std::size_t eidx) const {
235 return replies_.count(Key(aidx, eidx)) > 0;
236 }
237 bool isRouter(std::size_t node) const {
238 return node != 0 && net_raw().nodes[node - 1].nodetype == lang::NodeType::Router;
239 }
240 const qn::NetworkStruct<T>& net_raw() const {
241 return const_cast<qn::Network<T>&>(net_).raw_struct();
242 }
243
244 std::size_t nrep(std::size_t idx) const {
245 return materialize_ ? static_cast<std::size_t>(replRaw_[idx]) : 1;
246 }
247 double poolFactor(std::size_t idx) const { return materialize_ ? 1.0 : replRaw_[idx]; }
248
249 /** Callee replicas one caller replica reaches; an unset fan-out solves repl(a)*f = repl(b). */
250 std::size_t fanOutOf(std::size_t a, std::size_t b, std::size_t rb) const {
251 if (rb <= 1) return 1;
252 double f = l_.fanout_at(a, b);
253 if (f <= 0.0) {
254 const double ra = replRaw_[a];
255 f = double(rb) > ra ? std::max(1.0, std::floor(double(rb) / ra)) : 1.0;
256 }
257 f = std::min(std::max(1.0, std::round(f)), double(rb));
258 return static_cast<std::size_t>(f);
259 }
260
261 std::vector<std::size_t> targetReplicas(std::size_t a, std::size_t arep, std::size_t b) const {
262 const std::size_t rb = nrep(b);
263 if (rb <= 1) return std::vector<std::size_t>(1, 0);
264 const std::size_t f = fanOutOf(a, b, rb);
265 std::vector<std::size_t> out;
266 for (std::size_t k = 0; k < f; ++k) out.push_back((arep * f + k) % rb);
267 return out;
268 }
269
270 Key hostKey(std::size_t tidx, std::size_t trep) const {
271 const std::size_t h = parent(tidx);
272 return Key(h, trep % nrep(h));
273 }
274
275 /** Replica 1 keeps the plain name; the suffix stays in [A-Za-z0-9_] (JSON object keys). */
276 static std::string suffixed(const std::string& name, std::size_t rep) {
277 return rep == 0 ? name : name + "_r" + std::to_string(rep + 1);
278 }
279
280 /** A subgraph copied per call site or replica repeats names; a repeat gets _d<n>. */
281 std::string uniqueName(const std::string& base) {
282 std::map<std::string, int>::iterator it = usedNames_.find(base);
283 if (it == usedNames_.end()) {
284 usedNames_[base] = 1;
285 return base;
286 }
287 it->second += 1;
288 return base + "_d" + std::to_string(it->second);
289 }
290
291 /** Local activity successors of aidx within an entry, in ascending index order. */
292 std::vector<std::size_t> localSucc(std::size_t aidx, const std::vector<std::size_t>& local) const {
293 std::vector<std::size_t> out;
294 for (std::size_t s : l_.graph.succ(aidx))
295 if (std::find(local.begin(), local.end(), s) != local.end()) out.push_back(s);
296 return out;
297 }
298
299 /** Rows (target entry, probability) of the forwarding calls of an entry. */
300 std::vector<std::pair<std::size_t, T> > forwardingOf(std::size_t eidx) const {
301 std::vector<std::pair<std::size_t, T> > fwd;
302 const T zero = num_traits<T>::from_int(0), one = num_traits<T>::from_int(1);
303 for (std::size_t c = 1; c <= l_.ncalls; ++c) {
304 if (l_.calltype[c] != lang::CallType::FWD || l_.callpair_src[c] != eidx) continue;
305 T p = l_.callproc_mean[c];
306 if (p < zero) p = zero;
307 if (p > one) p = one;
308 if (dbl(p) > lang::GlobalConstants::FineTol) fwd.push_back(std::make_pair(l_.callpair_dst[c], p));
309 }
310 return fwd;
311 }
312
313 bool taskHasAndFork(std::size_t tidx) const {
314 for (std::size_t a = l_.ashift + 1; a <= l_.ashift + l_.nacts; ++a)
315 if (parent(a) == tidx && isPostAnd(a)) return true;
316 return false;
317 }
318
319 // -------------------------------------------------------------- replication
320
321 /** Copies of the task step graphs a materialised expansion would create. */
322 double replInstantiations() const {
323 std::map<std::size_t, double> memo;
324 std::set<std::size_t> seeds(refTasks_.begin(), refTasks_.end());
325 for (std::size_t e : openEntries_) seeds.insert(parent(e));
326 double n = 0.0;
327 for (std::size_t t : seeds) n += replRaw_[t] * taskCost(t, std::vector<std::size_t>(), memo);
328 return n;
329 }
330
331 double taskCost(std::size_t tidx, const std::vector<std::size_t>& stack,
332 std::map<std::size_t, double>& memo) const {
333 if (std::find(stack.begin(), stack.end(), tidx) != stack.end()) return 1.0;
334 std::map<std::size_t, double>::const_iterator it = memo.find(tidx);
335 if (it != memo.end()) return it->second;
336 std::vector<std::size_t> deeper = stack;
337 deeper.push_back(tidx);
338 double c = 1.0;
339 for (std::size_t eidx : l_.entriesof[tidx]) {
340 std::vector<std::size_t> targets;
341 for (std::size_t a : l_.actsof[eidx])
342 for (std::size_t cidx : l_.callsof[a])
343 if (l_.calltype[cidx] == lang::CallType::SYNC ||
344 l_.calltype[cidx] == lang::CallType::ASYNC)
345 targets.push_back(l_.callpair_dst[cidx]);
346 // A forwarded entry is expanded per replica just as a call is
347 const std::vector<std::pair<std::size_t, T> > fwd = forwardingOf(eidx);
348 for (std::size_t f = 0; f < fwd.size(); ++f) targets.push_back(fwd[f].first);
349 for (std::size_t te : targets) {
350 const std::size_t b = parent(te);
351 c += double(fanOutOf(tidx, b, static_cast<std::size_t>(replRaw_[b]))) *
352 taskCost(b, deeper, memo);
353 }
354 }
355 memo[tidx] = c;
356 return c;
357 }
358
359 // -------------------------------------------------------------------- steps
360
361 std::size_t addStep(std::size_t aidx, Key host, const Dist& svc, const std::string& name,
362 bool blocks, bool isthink, Key ref) {
363 Step s;
364 s.aidx = aidx;
365 s.host = host;
366 s.svc = svc;
367 s.name = name;
368 s.blocks = blocks;
369 s.isthink = isthink;
370 s.ref = ref;
371 s.node = 0;
372 s.owner = steps_.size();
373 // Every thread-pool task replica on the stack holds a thread here; a sync caller releases it only on reply
374 std::set<Key> held;
375 for (const Key& k : threadStack_)
376 if (fcrTask_[k.first]) held.insert(k);
377 s.tasks.assign(held.begin(), held.end());
378 steps_.push_back(s);
379 return steps_.size() - 1;
380 }
381
382 /** A step on a Fork, Join or Router node: no station, no service, no class of its own. */
383 std::size_t addAuxStep(std::size_t node, std::size_t ownerStep, const std::string& name, Key ref) {
384 const std::size_t id = addStep(0, Key(0, 0), Dist::disabled_dist(), name, false, false, ref);
385 steps_[id].node = node;
386 steps_[id].owner = steps_[ownerStep].owner;
387 return id;
388 }
389
390 void addRoute(const Port& from, std::size_t to, const T& p) {
391 Flow f;
392 f.from = from.step;
393 f.to = to;
394 f.p = p;
395 f.fromSig = from.sig;
396 f.inTarget = false;
397 flow_.push_back(f);
398 }
399
400 /** Leaving a Cache node: the node switches into hit or miss, so the route is in the target class. */
401 void addCacheRoute(std::size_t from, std::size_t to) {
402 Flow f;
403 f.from = from;
404 f.to = to;
405 f.p = num_traits<T>::from_int(1);
406 f.fromSig = false;
407 f.inTarget = true;
408 flow_.push_back(f);
409 }
410
411 std::size_t cacheNodeOf(std::size_t tidx, std::size_t trep) {
412 const Key k(tidx, trep);
413 std::map<Key, std::size_t>::const_iterator it = cacheNodeOf_.find(k);
414 if (it != cacheNodeOf_.end()) return it->second;
415 qn::CacheParam<T> cp;
416 cp.nitems = l_.nitems[tidx];
417 cp.itemcap = l_.itemcap[tidx];
418 cp.replacestrat = l_.replacestrat[tidx];
419 const std::size_t nd = net_.add_cache(suffixed(nm(tidx) + "_Cache", trep), cp);
420 cacheNodeOf_[k] = nd;
421 return nd;
422 }
423
424 bool hasRetrieval(std::size_t tidx) const {
425 return tidx < l_.hasretrieval.size() && l_.hasretrieval[tidx];
426 }
427
428 /** One PS fetch station per cache replica, as in the LN cache sublayer. */
429 std::size_t fetchNodeOf(std::size_t tidx, std::size_t trep) {
430 const Key k(tidx, trep);
431 std::map<Key, std::size_t>::const_iterator it = fetchNodeOf_.find(k);
432 if (it != fetchNodeOf_.end()) return it->second;
433 const std::size_t nd =
434 net_.add_queue(suffixed(nm(tidx) + "_Cache_Fetch", trep), lang::SchedStrategy::PS);
435 fetchNodeOf_[k] = nd;
436 return nd;
437 }
438
439 std::size_t actThinkStation() {
440 if (actThinkNode_ == 0) actThinkNode_ = net_.add_delay("ActivityThink");
441 return actThinkNode_;
442 }
443
444 std::vector<Stage> synchCallStages(std::size_t aidx) {
445 std::vector<Stage> stages;
446 const T one = num_traits<T>::from_int(1);
447 for (std::size_t cidx : l_.callsof[aidx]) {
448 const lang::CallType ct = l_.calltype[cidx];
449 if (ct != lang::CallType::SYNC && ct != lang::CallType::ASYNC) continue;
450 const bool isasync = ct == lang::CallType::ASYNC;
451 const T m = l_.callproc_mean[cidx];
452 std::size_t nfull =
453 static_cast<std::size_t>(std::floor(dbl(m) + lang::GlobalConstants::FineTol));
454 T frac = T(m - num_traits<T>::from_int(static_cast<long>(nfull)));
455 if (nfull > MAXCALLSTAGES) {
456 warn("Call multiplicity " + lqn::detail::fmt_num(dbl(m)) + " on " + l_.callnames[cidx] +
457 " truncated to " + std::to_string(MAXCALLSTAGES) + " stages.");
458 nfull = MAXCALLSTAGES;
459 frac = num_traits<T>::from_int(0);
460 }
461 for (std::size_t k = 0; k < nfull; ++k) stages.push_back(Stage{l_.callpair_dst[cidx], one, isasync});
462 if (dbl(frac) > lang::GlobalConstants::FineTol)
463 stages.push_back(Stage{l_.callpair_dst[cidx], frac, isasync});
464 }
465 return stages;
466 }
467
468 // ---------------------------------------------------------------- expansion
469
470 Expansion expandEntry(std::size_t eidx, Key ref, std::size_t trep) {
471 for (const Key& k : entryStack_)
472 if (k.first == eidx && k.second == trep) {
473 warn("Recursive call cycle at entry " + nm(eidx) + " truncated.");
474 return Expansion();
475 }
476 entryStack_.push_back(Key(eidx, trep));
477 threadStack_.push_back(Key(parent(eidx), trep));
478 Expansion out = expandEntryBody(eidx, ref, trep);
479 entryStack_.pop_back();
480 threadStack_.pop_back();
481 return out;
482 }
483
484 Expansion expandEntryBody(std::size_t eidx, Key ref, std::size_t trep) {
485 Expansion out;
486 if (eidx >= l_.actsof.size() || l_.actsof[eidx].empty()) return out;
487 // The bound activity is the activity successor of the entry
488 std::size_t bound = NONE;
489 for (std::size_t s : l_.graph.succ(eidx)) {
490 if (l_.type[s] != lang::LqnElement::ACTIVITY) continue;
491 if (std::find(l_.actsof[eidx].begin(), l_.actsof[eidx].end(), s) == l_.actsof[eidx].end())
492 continue;
493 bound = s;
494 break;
495 }
496 if (bound == NONE) return out;
497 out = expandActivities(bound, eidx, ref, trep);
498
499 // Forwarding: with prob p the entry hands off to a target that replies to the original caller
500 const std::vector<std::pair<std::size_t, T> > fwd = forwardingOf(eidx);
501 if (!fwd.empty() && !out.replies.empty()) {
502 const std::vector<Exit> ownPorts = out.replies;
503 std::vector<Exit> fwdExits;
504 T pforw = num_traits<T>::from_int(0);
505 // The forwarder's thread is released at the handoff
506 const Key fwdThread = threadStack_.back();
507 threadStack_.pop_back();
508 const std::size_t fwdTidx = parent(eidx);
509 for (std::size_t f = 0; f < fwd.size(); ++f) {
510 const std::vector<std::size_t> reps = targetReplicas(fwdTidx, trep, parent(fwd[f].first));
511 const T p = fwd[f].second;
512 const T nreps = num_traits<T>::from_int(static_cast<long>(reps.size()));
513 std::size_t reached = 0;
514 for (std::size_t mrep : reps) {
515 Expansion fe = expandEntry(fwd[f].first, ref, mrep);
516 if (fe.first == NONE) continue;
517 ++reached;
518 for (const Exit& op : ownPorts) addRoute(Port(op.step, op.sig), fe.first, T(op.prob * p / nreps));
519 fwdExits.insert(fwdExits.end(), fe.replies.begin(), fe.replies.end());
520 fwdExits.insert(fwdExits.end(), fe.terms.begin(), fe.terms.end());
521 }
522 pforw = T(pforw + p * num_traits<T>::from_int(static_cast<long>(reached)) / nreps);
523 }
524 threadStack_.push_back(fwdThread);
525 // What is left of each of this entry's own ports still replies
526 T resid = T(num_traits<T>::from_int(1) - pforw);
527 if (resid < num_traits<T>::from_int(0)) resid = num_traits<T>::from_int(0);
528 for (Exit& ex : out.replies) ex.prob = T(ex.prob * resid);
529 out.replies.insert(out.replies.end(), fwdExits.begin(), fwdExits.end());
530 }
531 return out;
532 }
533
534 Expansion expandActivities(std::size_t a0, std::size_t eidx, Key ref, std::size_t trep) {
535 WalkCtx c;
536 c.eidx = eidx;
537 c.tidx = parent(a0);
538 c.trep = trep;
539 c.ref = ref;
540 c.localActs = l_.actsof[eidx];
541 Expansion out;
542 out.first = walk(c, a0).first;
543 // Every chain end returns the token to the caller: as a reply if the entry replies anywhere
544 if (c.sawReply) {
545 c.replyExits.insert(c.replyExits.end(), c.terminals.begin(), c.terminals.end());
546 c.terminals.clear();
547 }
548 out.replies = c.replyExits;
549 out.terms = c.terminals;
550 return out;
551 }
552
553 std::string sname(const WalkCtx& c, const std::string& name) {
554 return uniqueName(suffixed(name, c.trep));
555 }
556
557 /** Move the terminals appended since n0 into the phase-2 chain ends. */
558 void takePh2Terminals(WalkCtx& c, std::size_t n0) {
559 for (std::size_t r = n0; r < c.terminals.size(); ++r) {
560 Ph2Exit pe;
561 pe.step = c.terminals[r].step;
562 pe.sig = c.terminals[r].sig;
563 pe.p = c.terminals[r].prob;
564 pe.ref = c.ref;
565 ph2Exits_.push_back(pe);
566 }
567 c.terminals.resize(n0);
568 }
569
570 std::pair<std::size_t, Port> walk(WalkCtx& c, std::size_t aidx) {
571 {
572 typename std::map<std::size_t, std::pair<std::size_t, Port> >::const_iterator it =
573 c.visited.find(aidx);
574 if (it != c.visited.end()) return it->second;
575 }
576 const T one = num_traits<T>::from_int(1);
577 const std::size_t tidx = c.tidx;
578 // The activity bound to an ItemEntry of a CacheTask is the read step, on the Cache node
579 std::size_t cacheNode = 0;
580 if (tidx < l_.iscache.size() && l_.iscache[tidx] && dbl(l_.graph.get(c.eidx, aidx)) > 0.0)
581 cacheNode = cacheNodeOf(tidx, c.trep);
582 std::pair<std::size_t, Port> se = makeActivitySteps(c, aidx, cacheNode);
583 const std::size_t entryStep = se.first;
584 Port exitPort = se.second;
585 c.visited[aidx] = se;
586
587 // The reply is deferred to the chain ends, not emitted here
588 const bool repliesHere = repliesTo(aidx, c.eidx);
589 if (repliesHere) c.sawReply = true;
590
591 const std::vector<std::size_t> succ = localSucc(aidx, c.localActs);
592 if (succ.empty()) {
593 c.terminals.push_back(Exit{exitPort.step, exitPort.sig, one});
594 return se;
595 }
596
597 // Phase 2 runs after the reply via sn.classspawn and needs a station departure in the step's own class
598 if (repliesHere) {
599 const Key hk = hostKey(tidx, c.trep);
600 if (cacheNode != 0 && succ.size() >= 2) {
601 // Phase 2 at a cache read: each hit/miss outcome routes through its own immediate trigger step
602 const std::size_t trigH =
603 addStep(aidx, hk, Dist::disabled_dist(), sname(c, nm(aidx) + "_ph2h"), false, false, c.ref);
604 const std::size_t trigM =
605 addStep(aidx, hk, Dist::disabled_dist(), sname(c, nm(aidx) + "_ph2m"), false, false, c.ref);
606 addCacheRoute(entryStep, trigH);
607 addCacheRoute(entryStep, trigM);
608 if (hasRetrieval(tidx))
609 warn("Delayed-hit retrieval of " + nm(tidx) +
610 " is not represented on a cache read with phase-2 successors.");
611 cacheWiring_.push_back(CacheWire{cacheNode, entryStep, trigH, trigM, c.eidx, 0,
613 c.replyExits.push_back(Exit{trigH, false, one});
614 c.replyExits.push_back(Exit{trigM, false, one});
615 const std::vector<Key> savedStack = threadStack_;
616 threadStack_.assign(1, Key(tidx, c.trep));
617 const std::size_t nT0 = c.terminals.size();
618 for (std::size_t hm = 0; hm < 2; ++hm) {
619 std::size_t sEntry = walk(c, succ[hm]).first;
620 if (steps_[sEntry].node != 0) {
621 const std::size_t head = addStep(aidx, hk, Dist::disabled_dist(),
622 sname(c, nm(aidx) + "_ph2b" + std::to_string(hm + 1)),
623 false, false, c.ref);
624 addRoute(Port(head, false), sEntry, one);
625 sEntry = head;
626 }
627 spawnPairs_.push_back(std::make_pair(hm == 0 ? trigH : trigM, sEntry));
628 }
629 takePh2Terminals(c, nT0);
630 threadStack_ = savedStack;
631 return se;
632 }
633 // Inside an AND-fork branch the lift applies only at the branch tail
634 const bool okCtx = cacheNode == 0 && (c.forkOwnerStack.empty() || isAndJoinPre(aidx));
635 // A merge step at the host station normalises a call-site exit to a station departure
636 if (okCtx && (exitPort.sig || isRouter(steps_[exitPort.step].node))) {
637 const std::size_t trig =
638 addStep(aidx, hk, Dist::disabled_dist(), sname(c, nm(aidx) + "_ph2t"), false, false, c.ref);
639 addRoute(exitPort, trig, one);
640 exitPort = Port(trig, false);
641 }
642 const bool ph2Spawn = okCtx && !exitPort.sig && steps_[exitPort.step].node == 0;
643 if (!ph2Spawn) {
644 warn("Phase-2 activities of " + nm(aidx) +
645 " run before the reply: the boundary is not a station departure, a degenerate "
646 "cache read, or mid-branch inside an AND-fork.");
647 } else {
648 c.replyExits.push_back(Exit{exitPort.step, exitPort.sig, one});
649 const std::vector<Key> savedStack = threadStack_;
650 threadStack_.assign(1, Key(tidx, c.trep));
651 const std::size_t nT0 = c.terminals.size();
652 std::vector<std::size_t> posSucc;
653 for (std::size_t s : succ)
654 if (dbl(l_.graph.get(aidx, s)) > 0.0) posSucc.push_back(s);
655 std::size_t target = NONE;
656 if (isAndFork(succ)) {
657 // Phase 2 opens with an AND-fork: spawn into an immediate head and fork from there
658 const std::size_t head =
659 addStep(aidx, hk, Dist::disabled_dist(), sname(c, nm(aidx) + "_ph2"), false, false, c.ref);
660 wireAndFork(c, Port(head, false), succ, aidx);
661 target = head;
662 } else if (isAndJoinPre(aidx)) {
663 // At an AND-join branch tail: the spawned head stands in for this branch at the Join
664 const std::size_t head =
665 addStep(aidx, hk, Dist::disabled_dist(), sname(c, nm(aidx) + "_ph2"), false, false, c.ref);
666 wireAndJoin(c, Port(head, false), succ[0]);
667 target = head;
668 } else if (posSucc.size() == 1) {
669 const std::size_t sEntry = walk(c, posSucc[0]).first;
670 if (steps_[sEntry].node == 0) target = sEntry;
671 }
672 if (target == NONE) {
673 // Branching phase 2, or a head on a non-station node: an immediate head carries the branches
674 const std::size_t head =
675 addStep(aidx, hk, Dist::disabled_dist(), sname(c, nm(aidx) + "_ph2"), false, false, c.ref);
676 for (std::size_t s2 : posSucc) {
677 const std::size_t sEntry = walk(c, s2).first;
678 addRoute(Port(head, false), sEntry, l_.graph.get(aidx, s2));
679 }
680 target = head;
681 }
682 spawnPairs_.push_back(std::make_pair(exitPort.step, target));
683 takePh2Terminals(c, nT0);
684 threadStack_ = savedStack;
685 return se;
686 }
687 }
688
689 if (cacheNode != 0) {
690 // CacheAccess: successors are the hit then the miss branch, and the Cache node decides
691 if (succ.size() < 2) {
692 warn("Cache read " + nm(aidx) + " has no hit/miss pair; treated as an ordinary activity.");
693 } else {
694 const std::size_t hEntry = walk(c, succ[0]).first;
695 const std::size_t mEntry = walk(c, succ[1]).first;
696 addCacheRoute(entryStep, hEntry);
697 addCacheRoute(entryStep, mEntry);
698 std::size_t fetch = 0;
699 Dist fetchSvc = Dist::disabled_dist();
700 if (hasRetrieval(tidx)) {
701 // The fetch is what the miss branch does, so its demand moves to the fetch station
702 fetch = fetchNodeOf(tidx, c.trep);
703 fetchSvc = steps_[mEntry].svc;
704 steps_[mEntry].svc = Dist::disabled_dist();
705 if (!l_.callsof[succ[1]].empty())
706 warn("Calls of miss activity " + nm(succ[1]) +
707 " stay outside the retrieval system, so they are not coalesced across "
708 "concurrent misses.");
709 }
710 cacheWiring_.push_back(CacheWire{cacheNode, entryStep, hEntry, mEntry, c.eidx, fetch, fetchSvc});
711 return se;
712 }
713 }
714
715 if (isAndFork(succ)) {
716 wireAndFork(c, exitPort, succ, aidx);
717 return se;
718 }
719 if (isAndJoinPre(aidx)) {
720 wireAndJoin(c, exitPort, succ[0]);
721 return se;
722 }
723 for (std::size_t s : succ) {
724 const T p = l_.graph.get(aidx, s);
725 if (!(dbl(p) > 0.0)) continue;
726 const std::size_t sEntry = walk(c, s).first;
727 addRoute(exitPort, sEntry, p);
728 }
729 return se;
730 }
731
732 /** AND-fork: a Fork replicates the job, one Router per branch since a Fork cannot switch class per link. */
733 void wireAndFork(WalkCtx& c, const Port& from, std::vector<std::size_t> fsucc, std::size_t aidx) {
734 const T one = num_traits<T>::from_int(1);
735 const std::string forkName = sname(c, "Fork_" + nm(aidx));
736 const std::size_t forkNode = net_.add_fork(forkName);
737 const std::size_t forkStep = addAuxStep(forkNode, from.step, forkName, c.ref);
738 addRoute(from, forkStep, one);
739 c.forkOwnerStack.push_back(forkStep);
740 // Walk a replying branch first so the Join and post-join subgraph are created in its phase-2 context
741 std::vector<std::size_t> ordered;
742 for (std::size_t s : fsucc)
743 if (branchReplies(c, s)) ordered.push_back(s);
744 for (std::size_t s : fsucc)
745 if (!branchReplies(c, s)) ordered.push_back(s);
746 for (std::size_t b = 0; b < ordered.size(); ++b) {
747 const std::string routerName = sname(c, "Fork_" + nm(aidx) + "_" + std::to_string(b + 1));
748 const std::size_t routerNode = net_.add_router(routerName);
749 const std::size_t routerStep = addAuxStep(routerNode, forkStep, routerName, c.ref);
750 addRoute(Port(forkStep, false), routerStep, one);
751 const std::size_t sEntry = walk(c, ordered[b]).first;
752 addRoute(Port(routerStep, false), sEntry, one);
753 }
754 c.forkOwnerStack.pop_back();
755 }
756
757 /** Route a branch tail into the AND-join, created on first arrival, in the class that entered the fork. */
758 void wireAndJoin(WalkCtx& c, const Port& from, std::size_t joinAidx) {
759 const T one = num_traits<T>::from_int(1);
760 std::map<std::size_t, std::size_t>::const_iterator it = c.joinOf.find(joinAidx);
761 if (it != c.joinOf.end()) {
762 // A Join closes one fork, so every branch tail must arrive with that fork innermost
763 if (c.forkOwnerStack.empty() || c.forkOwnerStack.back() != c.joinFork.at(joinAidx))
764 throw InputError("LQN2QN: AND-join at " + nm(joinAidx) +
765 " merges branches of different AND-forks: the fork-join structure is not"
766 " nested, so it has no Fork/Join representation.");
767 addRoute(from, it->second, one);
768 return;
769 }
770 if (c.forkOwnerStack.empty()) {
771 warn("AND-join at " + nm(joinAidx) + " has no enclosing AND-fork; branches are serialised.");
772 const std::size_t sEntry = walk(c, joinAidx).first;
773 addRoute(from, sEntry, one);
774 return;
775 }
776 const std::string joinName = sname(c, "Join_" + nm(joinAidx));
777 const std::size_t forkOwner = c.forkOwnerStack.back();
778 const std::size_t joinNode = net_.add_join(joinName, steps_[forkOwner].node);
779 const std::size_t joinStep = addAuxStep(joinNode, forkOwner, joinName, c.ref);
780 addRoute(from, joinStep, one);
781 // Registered before the walk, which runs outside the fork the Join closes
782 c.joinOf[joinAidx] = joinStep;
783 c.joinFork[joinAidx] = forkOwner;
784 const std::vector<std::size_t> savedForks = c.forkOwnerStack;
785 c.forkOwnerStack.pop_back();
786 const std::size_t sEntry = walk(c, joinAidx).first;
787 c.forkOwnerStack = savedForks;
788 addRoute(Port(joinStep, false), sEntry, one);
789 joinQuorum_.push_back(std::make_pair(joinStep, joinAidx));
790 }
791
792 /** True if the branch rooted at a0 replies to the current entry, searching up to the closing AND-join. */
793 bool branchReplies(const WalkCtx& c, std::size_t a0) const {
794 std::vector<std::size_t> stack(1, a0), seen;
795 while (!stack.empty()) {
796 const std::size_t a = stack.back();
797 stack.pop_back();
798 if (std::find(seen.begin(), seen.end(), a) != seen.end()) continue;
799 seen.push_back(a);
800 if (repliesTo(a, c.eidx)) return true;
801 if (isAndJoinPre(a)) continue;
802 const std::vector<std::size_t> nxt = localSucc(a, c.localActs);
803 stack.insert(stack.end(), nxt.begin(), nxt.end());
804 }
805 return false;
806 }
807
808 /** One step for the host demand, plus one per unrolled call stage; returns (entry step, exit port). */
809 std::pair<std::size_t, Port> makeActivitySteps(WalkCtx& c, std::size_t aidx, std::size_t cacheNode) {
810 const T one = num_traits<T>::from_int(1);
811 const std::size_t tidx = c.tidx, trep = c.trep;
812 const Key hidx = hostKey(tidx, trep);
813 if (cacheNode != 0) {
814 // A read step holds no demand and issues no call: the work is on the hit or miss branch
815 const std::size_t entryStep = addStep(aidx, hidx, Dist::disabled_dist(),
816 uniqueName(suffixed(nm(aidx), trep)), false, false, c.ref);
817 steps_[entryStep].node = cacheNode;
818 if (!l_.callsof[aidx].empty())
819 warn("Calls issued by cache read activity " + nm(aidx) + " are ignored.");
820 if (timed(l_.hostdem[aidx]))
821 warn("Host demand of cache read activity " + nm(aidx) + " is ignored.");
822 return std::make_pair(entryStep, Port(entryStep, false));
823 }
824 const Dist svc = timed(l_.hostdem[aidx]) ? l_.hostdem[aidx] : Dist::disabled_dist();
825 const std::vector<Stage> callStages = synchCallStages(aidx);
826 const std::size_t entryStep =
827 addStep(aidx, hidx, svc, uniqueName(suffixed(nm(aidx), trep)), false, false, c.ref);
828 Port cur(entryStep, false);
829
830 // Activity think time: a delay in series with the host demand on a shared INF station
831 if (aidx < l_.actthink.size() && timed(l_.actthink[aidx])) {
832 const std::size_t thinkStep = addStep(aidx, hidx, l_.actthink[aidx],
833 uniqueName(suffixed(nm(aidx) + "_think", trep)), false,
834 false, c.ref);
835 steps_[thinkStep].node = actThinkStation();
836 addRoute(cur, thinkStep, one);
837 cur = Port(thinkStep, false);
838 }
839
840 // A call holds the caller's server only at a finite, non-pool, closed-chain task alone on its host
841 const bool hostBlocks = !hostIsDelay_[hidx] && !fcrTask_[tidx] && c.ref.first != 0 &&
842 std::isfinite(l_.mult[tidx]) && l_.sched[tidx] != lang::SchedStrategy::INF &&
843 hostNTasks_.at(l_.parent[tidx]) == 1;
844
845 for (std::size_t k = 0; k < callStages.size(); ++k) {
846 const Stage& stage = callStages[k];
847 // An async call does not hold the caller's server; it stays serialised, the approximation
848 const bool blocks = hostBlocks && !stage.isasync;
849 // An async send also releases the caller's thread for the callee expansion
850 Key asyncThread(0, 0);
851 if (stage.isasync) {
852 asyncThread = threadStack_.back();
853 threadStack_.pop_back();
854 }
855 // The stage reaches the callee replicas this caller replica addresses, sharing the mean uniformly
856 const std::vector<std::size_t> reps = targetReplicas(tidx, trep, parent(stage.target));
857 std::vector<std::size_t> calleeFirsts;
858 std::vector<Exit> calleeReplies;
859 for (std::size_t mrep : reps) {
860 Expansion ce = expandEntry(stage.target, c.ref, mrep);
861 if (ce.first == NONE) continue;
862 calleeFirsts.push_back(ce.first);
863 // A callee path that neither replies nor continues still returns the token
864 calleeReplies.insert(calleeReplies.end(), ce.replies.begin(), ce.replies.end());
865 calleeReplies.insert(calleeReplies.end(), ce.terms.begin(), ce.terms.end());
866 }
867 if (stage.isasync) threadStack_.push_back(asyncThread);
868 if (calleeFirsts.empty()) continue; // callee not expandable: drop the call, never block
869 const T share = T(one / num_traits<T>::from_int(static_cast<long>(calleeFirsts.size())));
870
871 // A merge step is needed when the call may be skipped, another stage follows, the call does
872 // not block, or an AND-join branch tail must reach the Join in an ordinary class
873 const bool needsMerge =
874 stage.prob < one || k + 1 < callStages.size() || !blocks || isAndJoinPre(aidx);
875 std::size_t nxt = NONE;
876 if (needsMerge) {
877 const std::string retName =
878 uniqueName(suffixed(nm(aidx) + "_c" + std::to_string(k + 1) + "_ret", trep));
879 nxt = addStep(aidx, hidx, Dist::disabled_dist(), retName, false, false, c.ref);
880 // A non-blocking return carries no signal, so its merge point can sit on a Router
881 if (!blocks) steps_[nxt].node = net_.add_router(retName);
882 }
883
884 if (blocks) {
885 // The first mandatory call binds to the service class itself: a class switch would release the server
886 std::size_t blk;
887 if (k == 0 && !(stage.prob < one) && cur == Port(entryStep, false)) {
888 blk = entryStep;
889 steps_[blk].blocks = true;
890 } else {
891 blk = addStep(aidx, hidx, Dist::disabled_dist(),
892 uniqueName(suffixed(nm(aidx) + "_c" + std::to_string(k + 1), trep)), true,
893 false, c.ref);
894 addRoute(cur, blk, stage.prob);
895 if (stage.prob < one) addRoute(cur, nxt, T(one - stage.prob));
896 }
897 // Every reached replica replies into the same signal, so the call site blocks once
898 for (std::size_t cf : calleeFirsts) addRoute(Port(blk, false), cf, share);
899 for (const Exit& r : calleeReplies) reply_.push_back(Reply{r.step, blk, r.sig, r.prob});
900 if (needsMerge) {
901 addRoute(Port(blk, true), nxt, one);
902 cur = Port(nxt, false);
903 } else {
904 cur = Port(blk, true);
905 }
906 } else {
907 for (std::size_t cf : calleeFirsts) addRoute(cur, cf, T(stage.prob * share));
908 if (stage.prob < one) addRoute(cur, nxt, T(one - stage.prob));
909 for (const Exit& r : calleeReplies) addRoute(Port(r.step, r.sig), nxt, r.prob);
910 cur = Port(nxt, false);
911 }
912 }
913 return std::make_pair(entryStep, cur);
914 }
915
916 std::size_t stationOf(std::size_t i) {
917 const Step& s = steps_[i];
918 if (s.node != 0) return s.node;
919 if (s.isthink) return thinkNode_.at(s.ref);
920 return hostStation_.at(s.host);
921 }
922
923 // -------------------------------------------------------------------- build
924
925 void build() {
926 const lqn::LqnStruct<T>& l = l_;
927 const T zero = num_traits<T>::from_int(0), one = num_traits<T>::from_int(1);
928 const double FineTol = lang::GlobalConstants::FineTol;
929 if (replication_ != "auto" && replication_ != "materialize" && replication_ != "pool")
930 throw InputError("LQN2QN: replication must be 'auto', 'materialize' or 'pool'.");
931 const std::size_t NT = l.nhosts + l.ntasks;
932
933 for (std::size_t t = l.tshift + 1; t <= l.tshift + l.ntasks; ++t)
934 if (l.isref[t]) refTasks_.push_back(t);
935
936 // Entries with an open arrival process: Source -> open classes -> Sink
937 for (std::size_t e = l.eshift + 1; e <= l.eshift + l.nentries; ++e) {
938 if (e >= l.arrival.size() || !l.has_arrival[e] || l.arrival[e].disabled) continue;
939 const double m = dbl(l.arrival[e].mean);
940 if (std::isfinite(m) && m > FineTol) openEntries_.push_back(e);
941 }
942 if (refTasks_.empty() && openEntries_.empty())
943 throw InputError("LQN2QN: LQN must have at least one reference task or open arrival.");
944
945 for (std::size_t c = 1; c <= l.ncalls; ++c)
946 if (l.calltype[c] == lang::CallType::ASYNC) {
947 warn("Asynchronous calls are represented by LQN2QN as non-blocking visits: the caller "
948 "releases its server but remains serialised behind the callee.");
949 break;
950 }
951
952 // ---- replication
953 replRaw_.assign(NT + 1, 1.0);
954 std::vector<std::size_t> replicated;
955 for (std::size_t i = 1; i <= NT && i < l.repl.size(); ++i) {
956 replRaw_[i] = std::max(1.0, std::round(l.repl[i]));
957 if (replRaw_[i] > 1.0) replicated.push_back(i);
958 }
959 materialize_ = !replicated.empty() &&
960 (replication_ == "materialize" ||
961 (replication_ == "auto" && replInstantiations() <= double(MAXREPLINSTANCES)));
962 if (!replicated.empty() && !materialize_)
963 warn("Replication of " + nm(replicated[0]) +
964 " is pooled: its replicas become one station of r times the servers, one admission row "
965 "of r times the bound and one reference class of r times the population, which is exact "
966 "at an infinite-server host and optimistic elsewhere. Pass 'materialize' for one station "
967 "and one step-graph copy per replica.");
968
969 // A processor hosting any task besides the caller is not held across a synchronous call: LQN releases the
970 // processor while a thread waits, and holding it blocks the other tasks and deadlocks a call back onto it
971 hostNTasks_.clear();
972 for (std::size_t t = l.tshift + 1; t <= l.tshift + l.ntasks; ++t) ++hostNTasks_[l.parent[t]];
973
974 // ---- tasks whose multiplicity is a thread pool
975 fcrTask_.assign(NT + 1, false);
976 for (std::size_t t = l.tshift + 1; t <= l.tshift + l.ntasks; ++t) {
977 if (l.isref[t] || (t < l.iscache.size() && l.iscache[t])) continue;
978 if (!std::isfinite(l.mult[t]) || l.sched[t] == lang::SchedStrategy::INF) continue;
979 if (taskHasAndFork(t)) {
980 warn("Multiplicity of task " + nm(t) +
981 " is not enforced: an AND-fork inside a task cannot be capped by a finite capacity "
982 "region, whose job count would double-count the forked siblings.");
983 continue;
984 }
985 fcrTask_[t] = true;
986 }
987 // A task called transitively from inside an AND-fork branch is excluded too
988 std::vector<std::size_t> branchActs;
989 for (std::size_t a = l.ashift + 1; a <= l.ashift + l.nacts; ++a) {
990 if (!isPostAnd(a)) continue;
991 std::vector<std::size_t> frontier(1, a);
992 while (!frontier.empty()) {
993 const std::size_t cur = frontier.front();
994 frontier.erase(frontier.begin());
995 if (std::find(branchActs.begin(), branchActs.end(), cur) != branchActs.end()) continue;
996 branchActs.push_back(cur);
997 if (isAndJoinPre(cur)) continue; // branch tail: do not traverse past the join
998 for (std::size_t s : l.graph.succ(cur))
999 if (s > l.ashift && parent(s) == parent(cur)) frontier.push_back(s);
1000 }
1001 }
1002 if (!branchActs.empty() && std::find(fcrTask_.begin(), fcrTask_.end(), true) != fcrTask_.end()) {
1003 std::vector<std::size_t> front;
1004 for (std::size_t a : branchActs)
1005 for (std::size_t cidx : l.callsof[a]) front.push_back(parent(l.callpair_dst[cidx]));
1006 std::vector<bool> shadow(NT + 1, false);
1007 while (!front.empty()) {
1008 const std::size_t t = front.front();
1009 front.erase(front.begin());
1010 if (shadow[t]) continue;
1011 shadow[t] = true;
1012 for (std::size_t a = l.ashift + 1; a <= l.ashift + l.nacts; ++a)
1013 if (parent(a) == t)
1014 for (std::size_t cidx : l.callsof[a]) front.push_back(parent(l.callpair_dst[cidx]));
1015 }
1016 for (std::size_t t = 1; t <= NT; ++t)
1017 if (shadow[t] && fcrTask_[t]) {
1018 fcrTask_[t] = false;
1019 warn("Multiplicity of task " + nm(t) +
1020 " is not enforced: it is called from inside an AND-fork branch, whose flows the "
1021 "fork-join transformation retags outside the admission constraint.");
1022 }
1023 }
1024
1025 // ---- stations: one per host processor replica
1026 for (std::size_t h = 1; h <= l.nhosts; ++h) {
1027 const double nservers = l.mult[h];
1028 for (std::size_t m = 0; m < nrep(h); ++m) {
1029 const Key k(h, m);
1030 if (std::isinf(nservers) || l.sched[h] == lang::SchedStrategy::INF) {
1031 hostStation_[k] = net_.add_delay(suffixed(nm(h), m));
1032 hostIsDelay_[k] = true;
1033 } else {
1034 const std::size_t q = net_.add_queue(suffixed(nm(h), m), l.sched[h]);
1035 // A pooled processor carries the servers of all its replicas
1036 net_.set_number_of_servers(q, nservers * poolFactor(h));
1037 hostStation_[k] = q;
1038 hostIsDelay_[k] = false;
1039 }
1040 }
1041 }
1042 // ---- think delays, one per reference task replica
1043 for (std::size_t rt : refTasks_)
1044 for (std::size_t m = 0; m < nrep(rt); ++m)
1045 thinkNode_[Key(rt, m)] = net_.add_delay(suffixed(nm(rt) + "_Think", m));
1046
1047 // ---- pass 1: one closed chain per reference task replica
1048 for (std::size_t rt : refTasks_)
1049 for (std::size_t rep = 0; rep < nrep(rt); ++rep) {
1050 const Key refKey(rt, rep);
1051 const std::size_t thinkStep =
1052 addStep(0, Key(0, 0), Dist::disabled_dist(),
1053 uniqueName(suffixed(nm(rt) + "_Think", rep)), false, true, refKey);
1054 for (std::size_t eidx : l.entriesof[rt]) {
1055 const Expansion ex = expandEntry(eidx, refKey, rep);
1056 if (ex.first == NONE) continue;
1057 addRoute(Port(thinkStep, false), ex.first, one);
1058 // A reference task has no caller: replies and dead ends both close at the think delay
1059 for (const Exit& s : ex.replies) addRoute(Port(s.step, s.sig), thinkStep, s.prob);
1060 for (const Exit& s : ex.terms) addRoute(Port(s.step, s.sig), thinkStep, s.prob);
1061 }
1062 }
1063
1064 // ---- pass 1b: open arrival chains
1065 struct OpenWire {
1066 std::size_t eidx, first;
1067 std::vector<Exit> exits;
1068 };
1069 std::vector<OpenWire> openWiring;
1070 if (!openEntries_.empty()) {
1071 srcNode_ = net_.add_source("Source");
1072 snkNode_ = net_.add_sink("Sink");
1073 }
1074 for (std::size_t eidx : openEntries_) {
1075 // Each replica of the entry's task receives its own arrival stream
1076 for (std::size_t rep = 0; rep < nrep(parent(eidx)); ++rep) {
1077 Expansion ex = expandEntry(eidx, Key(0, 0), rep);
1078 if (ex.first == NONE) {
1079 warn("Open arrival entry " + nm(eidx) + " has no bound activity; ignored.");
1080 continue;
1081 }
1082 OpenWire ow;
1083 ow.eidx = eidx;
1084 ow.first = ex.first;
1085 ow.exits = ex.replies;
1086 ow.exits.insert(ow.exits.end(), ex.terms.begin(), ex.terms.end());
1087 openWiring.push_back(ow);
1088 }
1089 }
1090
1091 // ---- pass 2: classes and reply signals
1092 const std::size_t nsteps = steps_.size();
1093 std::vector<std::size_t> stepClass(nsteps, 0), stepSignal(nsteps, 0);
1094 for (std::size_t i = 0; i < nsteps; ++i) {
1095 const Step& s = steps_[i];
1096 if (s.owner != i) continue; // Fork/Join/Router steps travel in the class that entered the fork
1097 if (s.ref.first == 0) {
1098 stepClass[i] = net_.add_open_class(s.name);
1099 } else if (s.isthink) {
1100 // A pooled reference task holds the population of all its replicas
1101 stepClass[i] = net_.add_closed_class(s.name, l.mult[s.ref.first] * poolFactor(s.ref.first),
1102 thinkNode_.at(s.ref));
1103 } else {
1104 stepClass[i] = net_.add_closed_class(s.name, 0.0, thinkNode_.at(s.ref));
1105 }
1106 }
1107 for (std::size_t i = 0; i < nsteps; ++i) stepClass[i] = stepClass[steps_[i].owner];
1108
1109 // AND-join quorum, in the class the siblings are matched in
1110 for (const std::pair<std::size_t, std::size_t>& jq : joinQuorum_) {
1111 const std::size_t joinAidx = jq.second;
1112 if (joinAidx >= l.actquorum.size()) continue;
1113 const std::size_t quorum = l.actquorum[joinAidx];
1114 std::size_t nb = 0;
1115 for (std::size_t p : l.graph.pred(joinAidx))
1116 if (p != joinAidx && isAndJoinPre(p)) ++nb;
1117 // A quorum equal to the branch count is the default wait-for-all
1118 if (quorum < 1 || nb < 1 || quorum >= nb) continue;
1119 net_.set_join_strategy(steps_[jq.first].node, lang::JoinStrategy::PARTIAL, double(quorum));
1120 }
1121
1122 for (std::size_t i = 0; i < nsteps; ++i) {
1123 if (!steps_[i].blocks) continue;
1124 const std::size_t sig = net_.add_closed_class(steps_[i].name + "_Reply", 0.0, thinkNode_.at(steps_[i].ref));
1125 net_.set_reply_signal_class(stepClass[i], sig);
1126 stepSignal[i] = sig;
1127 }
1128
1129 // Spawn bindings for phase-2 continuations
1130 for (const std::pair<std::size_t, std::size_t>& sp : spawnPairs_)
1131 net_.set_class_spawn(stepClass[sp.first], stepClass[sp.second]);
1132
1133 // Phase-2 token destructor on closed chains: a NEGATIVE signal at a station nothing visits
1134 std::size_t ph2Dump = 0;
1135 std::map<Key, std::size_t> ph2DestructorOf;
1136 {
1137 std::set<Key> ph2Ref;
1138 for (const Ph2Exit& pe : ph2Exits_)
1139 if (pe.ref.first > 0) ph2Ref.insert(pe.ref);
1140 if (!ph2Ref.empty()) {
1141 ph2Dump = net_.add_queue("Ph2Sink", lang::SchedStrategy::FCFS);
1142 for (const Key& k : ph2Ref) {
1143 const std::size_t sig =
1144 net_.add_closed_class(suffixed("Ph2End_" + nm(k.first), k.second), 0.0, thinkNode_.at(k));
1145 net_.set_signal(sig, lang::SignalType::NEGATIVE);
1146 net_.set_service(ph2Dump, sig, Dist::immediate());
1147 ph2DestructorOf[k] = sig;
1148 }
1149 }
1150 }
1151
1152 // ---- pass 3: service times
1153 for (std::size_t i = 0; i < nsteps; ++i) {
1154 const Step& s = steps_[i];
1155 if (s.node != 0) {
1156 if (s.owner == i && actThinkNode_ != 0 && s.node == actThinkNode_) {
1157 net_.set_service(actThinkNode_, stepClass[i], s.svc);
1158 continue;
1159 }
1160 // A Router-hosted merge step owns a class, declared Immediate where the class is referenced
1161 if (s.owner == i && isRouter(s.node)) {
1162 if (s.ref.first == 0)
1163 net_.set_service(hostStation_.at(s.host), stepClass[i], Dist::immediate());
1164 else
1165 net_.set_service(thinkNode_.at(s.ref), stepClass[i], Dist::immediate());
1166 }
1167 continue;
1168 }
1169 if (s.isthink) {
1170 const Dist& th = l.think[s.ref.first];
1171 net_.set_service(thinkNode_.at(s.ref), stepClass[i], timed(th) ? th : Dist::immediate());
1172 } else {
1173 net_.set_service(hostStation_.at(s.host), stepClass[i], s.svc.disabled ? Dist::immediate() : s.svc);
1174 }
1175 }
1176
1177 // A SetupTask is the Queue setup/delay-off pair at its host station
1178 std::set<std::size_t> warnedSetup;
1179 for (std::size_t i = 0; i < nsteps; ++i) {
1180 const Step& s = steps_[i];
1181 if (s.node != 0 || s.isthink || s.aidx == 0) continue;
1182 const std::size_t t = parent(s.aidx);
1183 if (t >= l.hassetup.size() || !l.hassetup[t] || !timed(l.setuptime[t])) continue;
1184 if (l.delayofftime[t].disabled) {
1185 if (warnedSetup.insert(t).second)
1186 warn("Setup of setup task " + nm(t) +
1187 " is not represented: it has no delay-off time, so its server never shuts down and "
1188 "never sets up again.");
1189 continue;
1190 }
1191 if (hostIsDelay_[s.host]) {
1192 if (warnedSetup.insert(t).second)
1193 warn("Setup of setup task " + nm(t) +
1194 " is not represented: its processor is an infinite server, which never shuts down.");
1195 continue;
1196 }
1197 net_.set_setup_delayoff(hostStation_.at(s.host), stepClass[i], l.setuptime[t], l.delayofftime[t]);
1198 }
1199
1200 // A reply signal is consumed at the caller's station, declared everywhere it may fall back
1201 for (std::size_t i = 0; i < nsteps; ++i) {
1202 if (stepSignal[i] == 0) continue;
1203 for (const std::pair<const Key, std::size_t>& hs : hostStation_)
1204 net_.set_service(hs.second, stepSignal[i], Dist::immediate());
1205 for (const std::pair<const Key, std::size_t>& tn : thinkNode_)
1206 net_.set_service(tn.second, stepSignal[i], Dist::immediate());
1207 }
1208
1209 // ---- cache read/hit/miss wiring, now that the classes exist
1210 for (const CacheWire& cw : cacheWiring_) {
1211 const std::size_t rc = stepClass[cw.readStep];
1212 {
1213 qn::CacheParam<T>& cp = net_.raw_struct().nodeparam.at(cw.node);
1214 resizeCacheParam(cp);
1215 // The item pmf of the entry, over the cache's items: an item beyond the entry's cardinality is never read
1216 std::vector<T> pmf(cp.nitems, zero);
1217 const std::vector<T>& ip = l.itemproc[cw.eidx];
1218 for (std::size_t k = 0; k < cp.nitems && k < ip.size(); ++k) pmf[k] = ip[k];
1219 cp.pread[rc - 1] = pmf;
1220 // setReadItemEntry copies the entry's DiscreteSampler, the law the JMT export needs beside the pmf
1221 cp.preadkind[rc - 1].type = lang::ProcessType::DISCRETESAMPLER;
1222 cp.preadkind[rc - 1].n = cp.nitems;
1223 cp.hitclass[rc - 1] = stepClass[cw.hitStep];
1224 cp.missclass[rc - 1] = stepClass[cw.missStep];
1225 }
1226 if (cw.fetch != 0) {
1227 // Service and routing of the retrieval system are read off the read class
1228 net_.set_service(cw.fetch, rc, cw.fetchSvc.disabled ? Dist::immediate() : cw.fetchSvc);
1229 net_.set_retrieval_system(cw.node, rc, stepClass[cw.missStep], std::vector<std::size_t>(1, cw.fetch));
1230 }
1231 }
1232 for (std::pair<const std::size_t, qn::CacheParam<T> >& np : net_.raw_struct().nodeparam)
1233 resizeCacheParam(np.second);
1234
1235 // ---- pass 4: routing
1236 qn::RoutingMatrix<T> P = net_.init_routing_matrix();
1237 for (const Flow& f : flow_) {
1238 const std::size_t si = stationOf(f.from), sj = stationOf(f.to);
1239 if (f.inTarget)
1240 P.set(stepClass[f.to], stepClass[f.to], si, sj, f.p);
1241 else if (f.fromSig)
1242 P.set(stepSignal[f.from], stepClass[f.to], si, sj, f.p);
1243 else
1244 P.set(stepClass[f.from], stepClass[f.to], si, sj, f.p);
1245 }
1246 for (const Reply& r : reply_) {
1247 // A nested call returns through its own reply signal, switching into the outer call site's signal
1248 const std::size_t src = r.viaSig ? stepSignal[r.exit] : stepClass[r.exit];
1249 P.set(src, stepSignal[r.owner], stationOf(r.exit), stationOf(r.owner), r.p);
1250 }
1251 // Retrieval systems: the read class circulates cache -> fetch -> cache
1252 for (const CacheWire& cw : cacheWiring_)
1253 if (cw.fetch != 0) {
1254 const std::size_t rc = stepClass[cw.readStep];
1255 P.set(rc, rc, cw.node, cw.fetch, one);
1256 P.set(rc, rc, cw.fetch, cw.node, one);
1257 }
1258 // Phase-2 chain ends: destroy the spawned token
1259 for (const Ph2Exit& pe : ph2Exits_) {
1260 if (pe.ref.first == 0) {
1261 const std::size_t ec = stepClass[pe.step];
1262 P.set(ec, ec, stationOf(pe.step), snkNode_, pe.p);
1263 } else {
1264 const std::size_t src = pe.sig ? stepSignal[pe.step] : stepClass[pe.step];
1265 P.set(src, ph2DestructorOf.at(pe.ref), stationOf(pe.step), ph2Dump, pe.p);
1266 }
1267 }
1268 // Open arrival wiring: Source into the first step, exits into the Sink
1269 for (const OpenWire& ow : openWiring) {
1270 const std::size_t fc = stepClass[ow.first];
1271 net_.set_arrival(srcNode_, fc, l.arrival[ow.eidx]);
1272 P.set(fc, fc, srcNode_, stationOf(ow.first), one);
1273 // Open chains carry no signals, so every exit is an ordinary class
1274 for (const Exit& ex : ow.exits) {
1275 const std::size_t ec = stepClass[ex.step];
1276 P.set(ec, ec, stationOf(ex.step), snkNode_, ex.prob);
1277 }
1278 }
1279 net_.link(P);
1280
1281 // ---- thread pools: one finite capacity region, one admission row per task replica
1282 std::set<Key> fcrSet;
1283 for (const Step& s : steps_) fcrSet.insert(s.tasks.begin(), s.tasks.end());
1284 if (!fcrSet.empty()) {
1285 const std::vector<Key> fcrList(fcrSet.begin(), fcrSet.end());
1286 const std::size_t K = net_.raw_struct().classes.size();
1287 Matrix<T> A(fcrList.size(), K, zero);
1288 std::vector<T> b(fcrList.size(), zero);
1289 std::vector<std::size_t> regionNodes;
1290 bool anyA = false;
1291 for (std::size_t ts = 0; ts < fcrList.size(); ++ts) {
1292 for (std::size_t i = 0; i < nsteps; ++i) {
1293 const std::vector<Key>& tk = steps_[i].tasks;
1294 if (std::find(tk.begin(), tk.end(), fcrList[ts]) == tk.end()) continue;
1295 A(ts, stepClass[i] - 1) = one;
1296 anyA = true;
1297 const std::size_t nd = stationOf(i);
1298 if (isStation(nd) && std::find(regionNodes.begin(), regionNodes.end(), nd) == regionNodes.end())
1299 regionNodes.push_back(nd);
1300 // The caller holds its thread for the whole fetch
1301 for (const CacheWire& cw : cacheWiring_)
1302 if (cw.readStep == i && cw.fetch != 0 &&
1303 std::find(regionNodes.begin(), regionNodes.end(), cw.fetch) == regionNodes.end())
1304 regionNodes.push_back(cw.fetch);
1305 }
1306 // A pooled task keeps one row whose bound covers all its replicas
1307 const std::size_t t = fcrList[ts].first;
1308 b[ts] = tnum(l.mult[t] * poolFactor(t));
1309 }
1310 if (anyA && !regionNodes.empty()) {
1311 const std::size_t rg = net_.add_region(regionNodes, std::vector<double>());
1312 net_.set_region_constraint(rg, A, b);
1313 }
1314 }
1315 }
1316
1317 bool isStation(std::size_t node) const {
1318 const qn::NetworkStruct<T>& sn = net_raw();
1319 return sn.nodes[node - 1].station != 0 && sn.nodes[node - 1].nodetype != lang::NodeType::Source;
1320 }
1321
1322 void resizeCacheParam(qn::CacheParam<T>& cp) {
1323 const std::size_t K = net_.raw_struct().classes.size();
1324 if (cp.pread.size() < K) cp.pread.resize(K);
1325 if (cp.preadkind.size() < K) cp.preadkind.resize(K);
1326 if (cp.hitclass.size() < K) cp.hitclass.resize(K, 0);
1327 if (cp.missclass.size() < K) cp.missclass.resize(K, 0);
1328 if (cp.classitem.size() < K) cp.classitem.resize(K, 0);
1329 }
1330};
1331
1332/** Entry/activity index pairs that reply, from the explicit and the inferred replies of each entry. */
1333template <class T>
1334std::set<std::pair<std::size_t, std::size_t> > lqn2qn_replies(const lqn::LqnModel<T>& m,
1335 const lqn::LqnStruct<T>& l) {
1336 std::map<std::string, std::size_t> ent, act;
1337 for (std::size_t i = 1; i <= l.nidx; ++i) {
1338 if (l.type[i] == lang::LqnElement::ENTRY) ent[l.names[i]] = i;
1339 if (l.type[i] == lang::LqnElement::ACTIVITY) act[l.names[i]] = i;
1340 }
1341 std::set<std::pair<std::size_t, std::size_t> > out;
1342 const std::map<std::string, std::vector<std::string> > rep = lqn::detail::reply_activities(m, l);
1343 for (const std::pair<const std::string, std::vector<std::string> >& kv : rep) {
1344 std::map<std::string, std::size_t>::const_iterator ei = ent.find(kv.first);
1345 if (ei == ent.end()) continue;
1346 for (const std::string& a : kv.second) {
1347 std::map<std::string, std::size_t>::const_iterator ai = act.find(a);
1348 if (ai != act.end()) out.insert(std::make_pair(ai->second, ei->second));
1349 }
1350 }
1351 return out;
1352}
1353
1354} // namespace detail
1355
1356/**
1357 * Port of MATLAB `LQN2QN(lqn, replication)`: flatten a layered network into a
1358 * queueing network whose synchronous calls block through REPLY signals.
1359 *
1360 * @param lqn the layered model, from lqn::LqnBuilder::model() or lqn::read_lqnx_model
1361 * @param replication 'auto' (default), 'materialize' or 'pool', as in the reference
1362 * @param name the layered model's name; the network is called `<name>-QN`
1363 * @param warnings when non-null, the reference's line_warning messages are appended here
1364 * instead of being printed to stderr
1365 * @return the queueing network, routed and ready for get_struct()
1366 */
1367template <class T>
1368qn::Network<T> lqn2qn(const lqn::LqnModel<T>& lqn, const std::string& replication = "auto",
1369 const std::string& name = "model", std::vector<std::string>* warnings = nullptr) {
1371 const std::set<std::pair<std::size_t, std::size_t> > replies = detail::lqn2qn_replies(lqn, l);
1372 detail::Lqn2Qn<T> conv(l, replies, replication, name, warnings);
1373 return conv.run();
1374}
1375
1376/**
1377 * LQN2QN over an already flattened LqnStruct.
1378 *
1379 * The struct carries no replygraph, so only the IMPLICIT replies (leaf activities) are known here;
1380 * a model with an explicit reply on a non-leaf activity (phase 2) must use the LqnModel overload.
1381 */
1382template <class T>
1383qn::Network<T> lqn2qn(const lqn::LqnStruct<T>& lsn, const std::string& replication = "auto",
1384 const std::string& name = "model", std::vector<std::string>* warnings = nullptr) {
1385 const std::set<std::pair<std::size_t, std::size_t> > replies =
1386 detail::lqn2qn_replies(lqn::LqnModel<T>(), lsn);
1387 detail::Lqn2Qn<T> conv(lsn, replies, replication, name, warnings);
1388 return conv.run();
1389}
1390
1391} // namespace io
1392} // namespace line
1393
1394#endif // LINE_IO_LQN2QN_H
InputError(const std::string &what)
Definition error.h:39
A queueing network under construction.
The exception types the port throws.
.lqnx -> LqnStruct, a port of matlab/src/lang/layered/@LayeredNetwork/parseXML.m followed by ....
LayeredNetworkStruct, the flattened description of a layered queueing network.
LqnModel -> .lqnx, a port of matlab/src/lang/layered/@LayeredNetwork/writeXML.m.
Dense matrix and non-owning view.
qn::Network< T > lqn2qn(const lqn::LqnModel< T > &lqn, const std::string &replication="auto", const std::string &name="model", std::vector< std::string > *warnings=nullptr)
Port of MATLAB LQN2QN(lqn, replication): flatten a layered network into a queueing network whose sync...
Definition lqn2qn.h:1368
@ NEGATIVE
removes a batch of jobs (Gelenbe's negative customer)
Definition lang_types.h:169
CallType
Call kinds, with the values of MATLAB CallType.
Definition lang_types.h:469
LqnStruct< T > lqn_finalize(const LqnModel< T > &m)
Port of @LayeredNetwork/getStruct.m: flatten the model into its struct.
Definition lqn_reader.h:444
Conservation laws of a layered queueing network, enumerated from its structure.
Definition aoi_dist2ph.h:52
lang::Distrib< double > Dist
Definition nodes.h:54
The Network constructor API: Queue, Delay, Source, Sink, Router, ClassSwitch, Cache,...
Number-type abstraction for the templated API port.
static Distrib disabled_dist()
Definition lang_types.h:988
static Distrib immediate()
The Immediate singleton.
Definition lang_types.h:977
static constexpr double FineTol
Definition lang_types.h:760
The intermediate model, and the second stage that flattens it.
Definition lqn_reader.h:409