A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
graph.h
Go to the documentation of this file.
1/*
2 * Copyright 2026 The A11 Authors
3 *
4 * Licensed under the Apache License, Version 2.0 (the "License");
5 * you may not use this file except in compliance with the License.
6 * You may obtain a copy of the License at
7 *
8 * http://www.apache.org/licenses/LICENSE-2.0
9 *
10 * Unless required by applicable law or agreed to in writing, software
11 * distributed under the License is distributed on an "AS IS" BASIS,
12 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13 * See the License for the specific language governing permissions and
14 * limitations under the License.
15 */
16
17#ifndef A11_FLOW_GRAPH_H_
18#define A11_FLOW_GRAPH_H_
19
20#include <cstddef>
21#include <optional>
22#include <string>
23#include <string_view>
24#include <vector>
25
26#include <absl/base/nullability.h>
27#include <absl/container/flat_hash_map.h>
28#include <absl/container/flat_hash_set.h>
29#include <absl/time/time.h>
30
31#include "a11/flow/syntax.h"
32#include "a11/flow/vocabulary.h"
33
34namespace a11::flow::graph {
35
41using RefId = size_t;
42using StepId = size_t;
43using BodyId = size_t;
44using ExprId = size_t;
45
47inline constexpr size_t kNone = static_cast<size_t>(-1);
48
50enum class RefKind {
56 kNode,
59 kNodeId,
61 kStatus,
63 kHeader,
65 kExpr,
69 kZip,
72 kMerge,
76 kWinner,
79 kBound,
80};
81
94struct LogTail {
96 std::string level;
98 std::string format;
99 bool has_format = false;
101 std::vector<ExprId> arguments;
105 int line = 0;
106};
107
108struct Stage {
109 std::string name;
112 long long count = 0;
114 std::string text;
122 bool indexed = false;
123 long long index = 0;
127 bool named_or_indexed = false;
130 absl::Duration duration;
132 bool descending = false;
133 // `fold`: what the first pass carries, the name it is bound to, and the ref
134 // that name resolves to inside the fold's expression.
141 std::string carried;
145 bool tolerant = false;
149 int parallel = 1;
152 bool ordered = true;
153};
154
163struct Ref {
166 std::string label;
171 bool writable = false;
186 bool unary = false;
193 long long skip = 0;
194
196 std::string name;
203 std::string node_map;
209 std::string header;
211 bool has_fallback = false;
218 std::string role;
220 // `kZip`: the streams read in step, in the order they were written, which is
221 // the order their values appear in each tuple.
226 std::vector<RefId> sources;
227};
228
230enum class StepKind {
231 kCall,
232 kPipe,
233 kSkip,
234 kWait,
235 kDrain,
236 kCancel,
238 kAbort,
239 kFail,
240 kLog,
243 kCapture,
244 kForEach,
245 kRepeat,
246 kIf,
248 kBlock,
249};
250
251std::string_view StepKindName(StepKind kind);
252
263
265struct Step {
269 std::string label;
272 std::vector<StepId> after;
274
276 std::string action;
277 std::string mode;
278 std::string node_map;
279 std::optional<absl::Duration> timeout;
280 bool tee = false;
281 bool tolerant = false;
282 std::vector<std::pair<std::string, ExprId>> headers;
283 std::vector<std::string> forward;
286 absl::flat_hash_map<std::string, RefId> ports;
289
302 bool discard = false;
305 std::optional<long long> count;
307 std::string slot;
308
312 // `kWait`: the outcomes of a `wait first of` / `wait all of`, in the order
313 // they were written.
318 std::vector<RefId> subjects;
320 bool race = false;
324
327
337 std::string code_name;
338
341
344 std::vector<BodyId> bodies;
348 int parallel = 1;
356 std::optional<int> max_iterations;
358 bool stop_when = true;
359};
360
362struct Body {
363 std::string label;
366 std::vector<StepId> steps;
367};
368
376struct Expr {
377 const syntax::Node* absl_nullable node = nullptr;
378 std::vector<std::pair<const syntax::Node* absl_nonnull, RefId>> bound;
380 std::vector<RefId> refs;
381};
382
390struct FlowGraph {
391 std::string name;
392 std::vector<Ref> refs;
393 std::vector<Step> steps;
394 std::vector<Body> bodies;
395 std::vector<Expr> exprs;
398
400 [[nodiscard]] std::vector<RefId> Upstreams(RefId ref) const;
402 [[nodiscard]] std::vector<RefId> ValueRefs(RefId ref) const;
403
405 [[nodiscard]] std::vector<RefId> Sources(StepId step) const;
407 [[nodiscard]] std::vector<RefId> ValueSources(StepId step) const;
409 [[nodiscard]] std::vector<RefId> Destinations(StepId step) const;
410 // The refs the *stages* a step reads write to: a `try ...
418 [[nodiscard]] std::vector<RefId> StageDestinations(StepId step) const;
420 [[nodiscard]] std::vector<RefId> Observed(StepId step) const;
422 [[nodiscard]] std::vector<BodyId> NestedBodies(BodyId body) const;
423};
424
430struct Analysis {
433 std::vector<RefId> refs;
435 absl::flat_hash_map<RefId, int> readers;
436 // The ones read from inside a nested body, which are buffered once and
437 // replayed to each reader:
440 absl::flat_hash_set<RefId> materialise;
442 absl::flat_hash_map<RefId, int> writers;
444 std::vector<RefId> destinations;
446 absl::flat_hash_map<StepId, std::vector<RefId>> held;
448 std::vector<RefId> nodes;
449};
450
452Analysis Analyse(const FlowGraph& flow, BodyId body);
453
460 public:
461 explicit GraphBuilder(FlowGraph& flow) : flow_(&flow) {}
462
463 FlowGraph& flow() { return *flow_; }
464
466 StepId owner_step = kNone) {
467 Body body;
468 body.label = std::move(label);
469 body.parent = parent;
470 body.owner_step = owner_step;
471 flow_->bodies.push_back(std::move(body));
472 return flow_->bodies.size() - 1;
473 }
474
476 [[nodiscard]] bool Carries(RefId ref) const {
477 return ref != kNone && ref < flow_->refs.size() && flow_->refs[ref].unary;
478 }
479
486 static bool StageMakesOne(const Stage& stage) {
487 if (vocabulary::ReducingStages().contains(stage.name)) {
488 {
489 return true;
490 }
491 }
492 return stage.name == "first" && stage.count == 1;
493 }
494
506 static bool StagePreservesCount(const Stage& stage) {
507 static const auto* const kPerValue = new absl::flat_hash_set<std::string>{
508 "map", "at", "truncate", "text", "json", "packb", "strformat", "where",
509 "mime", "distinct", "first", "last", "drop", "log", "logf", "scan",
510 // `sort` reorders and keeps every value; `timeout` and `pace` say
511 // *when* a value may pass and change nothing about which do.
512 "sort", "timeout", "pace"};
513 return kPerValue->contains(stage.name);
514 }
515
524 switch (ref.kind) {
525 case RefKind::kHeader:
526 case RefKind::kNodeId:
527 case RefKind::kStatus:
528 case RefKind::kWinner:
529 case RefKind::kExpr:
530 case RefKind::kBound:
531 // One header, one id, one status record, one winner, one evaluated
532 // expression, and one value bound per pass of a loop.
533 ref.unary = true;
534 break;
535 case RefKind::kNode:
536 // A node is a stream the flow may write from anywhere, including from
537 // inside a loop, so nothing about it is provable from the plan.
538 ref.unary = false;
539 break;
541 ref.unary = StageMakesOne(ref.stage) ||
542 (Carries(ref.source) && StagePreservesCount(ref.stage));
543 break;
544 case RefKind::kZip:
545 // A tuple per round, and the rounds run until every source has ended,
546 // so one round is only certain when every source has at most one value.
547 ref.unary = !ref.sources.empty();
548 for (const RefId source : ref.sources) {
549 if (!Carries(source)) {
550 {
551 ref.unary = false;
552 }
553 }
554 }
555 break;
556 case RefKind::kMerge:
557 // Every value of every source, so one value only when they all have at
558 // most one *and* there is one of them: two unary sources interleave
559 // into two values.
560 ref.unary = ref.sources.size() == 1 && Carries(ref.sources.front());
561 break;
564 // Whatever the declaration said, which only the caller knows.
565 break;
566 }
567 flow_->refs.push_back(std::move(ref));
568 return flow_->refs.size() - 1;
569 }
570
573 const BodyId body = step.body;
574 flow_->steps.push_back(std::move(step));
575 const StepId id = flow_->steps.size() - 1;
576 if (body != kNone) {
577 {
578 flow_->bodies[body].steps.push_back(id);
579 }
580 }
581 return id;
582 }
583
585 flow_->exprs.push_back(std::move(expr));
586 return flow_->exprs.size() - 1;
587 }
588
589 Ref& ref(RefId id) { return flow_->refs[id]; }
590
591 Step& step(StepId id) { return flow_->steps[id]; }
592
593 Expr& expr(ExprId id) { return flow_->exprs[id]; }
594
595 private:
596 FlowGraph* absl_nonnull flow_;
597};
598
599} // namespace a11::flow::graph
600
601#endif // A11_FLOW_GRAPH_H_
Appends to a [FlowGraph] while the resolver walks a flow.
Definition graph.h:459
bool Carries(RefId ref) const
Whether the ref already in the graph carries at most one value.
Definition graph.h:476
StepId AddStep(Step step)
Append a step, and record it in the body it belongs to.
Definition graph.h:572
GraphBuilder(FlowGraph &flow)
Definition graph.h:461
static bool StagePreservesCount(const Stage &stage)
Whether a stage yields one value per value it was given.
Definition graph.h:506
ExprId AddExpr(Expr expr)
Definition graph.h:584
static bool StageMakesOne(const Stage &stage)
Whether a stage yields exactly one value however many it was given.
Definition graph.h:486
Step & step(StepId id)
Definition graph.h:591
Ref & ref(RefId id)
Definition graph.h:589
Expr & expr(ExprId id)
Definition graph.h:593
RefId AddRef(Ref ref)
Append a ref, working out what it carries.
Definition graph.h:523
BodyId AddBody(std::string label, BodyId parent=kNone, StepId owner_step=kNone)
Definition graph.h:465
FlowGraph & flow()
Definition graph.h:463
TokenKind kind
Definition format.cc:913
Definition graph.cc:21
constexpr size_t kNone
No such thing – the absent id, for the many fields only some kinds use.
Definition graph.h:47
size_t BodyId
Definition graph.h:43
std::string_view StepKindName(StepKind kind)
Definition graph.cc:38
size_t RefId
Everything in a graph is named by index into the graph that owns it.
Definition graph.h:41
Analysis Analyse(const FlowGraph &flow, BodyId body)
Work out who reads and writes what in one body.
Definition graph.cc:354
size_t StepId
Definition graph.h:42
bool RecordsOutcome(StepKind kind)
Whether a step records an outcome of its own for a name to read.
Definition graph.h:259
RefKind
What a stream in the plan is.
Definition graph.h:50
@ kNodeId
A node's id, as one value – what a flow hands to an action that expects to be told where to write.
@ kMerge
Several streams read at once, as one stream of their values in the order they arrive: interleave(a,...
@ kBound
A stream the runtime binds per pass: a loop's value, its index, or a repeat's carry.
@ kZip
Several streams read in step, as one stream of tuples: zip(a, b).
@ kStatus
The outcome of a call, a node, a port or a barrier, as a status record.
@ kCallPort
A port of an action this flow calls: x.out.
@ kWinner
Which subject of a wait first of finished first, counted from zero, as one value.
@ kFlowPort
A declared port of this flow.
@ kHeader
One header of the call running this flow, as a single value.
@ kDerived
One stage applied to another stream.
@ kNode
A node of the flow's own: a stream it can write and read back.
@ kExpr
A stream of one value: an expression evaluated once.
size_t ExprId
Definition graph.h:44
StepKind
What a statement became.
Definition graph.h:230
@ kAbort
abort node ..: end a node with a failure rather than with an end.
@ kBlock
[try] { ... }: a body run as one step, whose outcome is its own.
@ kCapture
Remember a stream's first value for the loop that owns this body: what <- and an until condition comp...
PortDirection
Which side of the descriptor a port lands on, spelled as the plan spells it.
Definition syntax.h:728
const absl::flat_hash_set< std::string_view > & ReducingStages()
The stages that read a whole stream and yield exactly one value.
Definition vocabulary.cc:2037
StageArgument
What a stage takes after its name.
Definition vocabulary.h:51
@ kNone
Nothing: | collect, | count.
std::string label
Definition resolve.cc:684
graph::StepId step
The graph step it is, so <- and until can fill it in.
Definition resolve.cc:3763
const Scope * parent
Definition resolve.cc:730
std::string ref
Definition sqlite_chunk_store.cc:193
Who reads and who writes each ref a body owns.
Definition graph.h:430
std::vector< RefId > nodes
The nodes of its own this body names, wherever it names them.
Definition graph.h:448
std::vector< RefId > refs
Every ref this body owns.
Definition graph.h:433
absl::flat_hash_map< RefId, int > writers
How many writers each written ref has.
Definition graph.h:442
BodyId body
Definition graph.h:431
absl::flat_hash_map< StepId, std::vector< RefId > > held
Which destinations each step holds open until it finishes.
Definition graph.h:446
absl::flat_hash_set< RefId > materialise
Refs read inside a nested body.
Definition graph.h:440
std::vector< RefId > destinations
Every ref this body writes, whether or not anything reads it back.
Definition graph.h:444
absl::flat_hash_map< RefId, int > readers
How many readers each of them has.
Definition graph.h:435
A block of steps: a flow's top level, or a loop or branch body.
Definition graph.h:362
StepId owner_step
Definition graph.h:365
BodyId parent
Definition graph.h:364
std::vector< StepId > steps
Definition graph.h:366
std::string label
Definition graph.h:363
One resolved expression, and the streams it reads.
Definition graph.h:376
const syntax::Node *absl_nullable node
Definition graph.h:377
std::vector< std::pair< const syntax::Node *absl_nonnull, RefId > > bound
Definition graph.h:378
std::vector< RefId > refs
The same refs, in the order they were first mentioned, for the analysis.
Definition graph.h:380
The executable graph of one flow.
Definition graph.h:390
std::vector< RefId > ValueRefs(RefId ref) const
The refs this one reads for their first value to produce itself.
Definition graph.cc:97
std::vector< RefId > ValueSources(StepId step) const
The refs a step reads for one value each.
Definition graph.cc:163
std::vector< RefId > StageDestinations(StepId step) const
The refs the stages a step reads write to: a try ... into failures.
Definition graph.cc:199
std::vector< Body > bodies
Definition graph.h:394
std::vector< Step > steps
Definition graph.h:393
std::string name
Definition graph.h:391
std::vector< Expr > exprs
Definition graph.h:395
std::vector< BodyId > NestedBodies(BodyId body) const
Every body inside this one, at any depth.
Definition graph.cc:242
BodyId root
The flow's top level.
Definition graph.h:397
std::vector< RefId > Destinations(StepId step) const
The refs a step writes, one entry per independent writer.
Definition graph.cc:190
std::vector< RefId > Upstreams(RefId ref) const
The refs this one reads as streams to produce itself.
Definition graph.cc:74
std::vector< RefId > Sources(StepId step) const
The refs a step reads as streams, one entry per independent read.
Definition graph.cc:130
std::vector< RefId > Observed(StepId step) const
Written refs a step only watches, without writing them itself.
Definition graph.cc:224
std::vector< Ref > refs
Definition graph.h:392
One resolved pipeline stage.
Definition graph.h:94
std::vector< ExprId > arguments
What to log, or what fills the format. Empty in a stage means it.
Definition graph.h:101
std::string level
The level, canonically spelled, or empty for the default.
Definition graph.h:96
std::string format
logf's format. Empty and has_format false for a log.
Definition graph.h:98
bool has_format
Definition graph.h:99
int line
The line it was written on, which the log carries so a consumer can point at it.
Definition graph.h:105
One stream in the plan.
Definition graph.h:163
std::string label
How a reader would say it: search.hits | truncate 200.
Definition graph.h:166
BodyId owner
The body this belongs to, which is where it is materialised.
Definition graph.h:168
StepId subject_step
kStatus: the call or barrier it is the outcome of.
Definition graph.h:207
StepId bound_by
Definition graph.h:219
ExprId expr
kExpr: the expression evaluated once.
Definition graph.h:213
long long skip
How many of this stream's first values skip n has spoken for.
Definition graph.h:193
RefId subject
kNodeId: the node. kStatus: the node or port it is the outcome of.
Definition graph.h:205
std::string header
kHeader: the header name, and the default when it was not sent.
Definition graph.h:209
StepId call
kCallPort: the call it belongs to.
Definition graph.h:200
Stage stage
Definition graph.h:216
std::vector< RefId > sources
kZip: the streams read in step, in the order they were written, which is the order their values appea...
Definition graph.h:226
ExprId id_expr
kNode: the id expression it attaches to, and the map it lands in.
Definition graph.h:202
syntax::Constant fallback
Definition graph.h:210
std::string node_map
Definition graph.h:203
std::string name
kFlowPort, kCallPort, kNode: the port or node name.
Definition graph.h:196
bool has_fallback
Definition graph.h:211
bool unary
Whether this stream provably carries at most one value.
Definition graph.h:186
RefId source
kDerived: what it reads, and the stage applied to it.
Definition graph.h:215
bool writable
Whether this flow could ever write it, which is what makes it a destination rather than something fin...
Definition graph.h:171
RefKind kind
Definition graph.h:164
std::string role
kBound: which of a loop's streams this is – item, index, carry.
Definition graph.h:218
syntax::PortDirection direction
kFlowPort, kCallPort: which side it is on.
Definition graph.h:198
Definition graph.h:108
ExprId expr
kExpression: the expression, with it bound.
Definition graph.h:116
bool named_or_indexed
Whether at should try its name and then its index, which is what a destructuring let needs: let name,...
Definition graph.h:127
vocabulary::StageArgument takes
Definition graph.h:110
bool indexed
Whether at was written as an index rather than a field name.
Definition graph.h:122
bool tolerant
try: a value this stage cannot do is dropped rather than ending the pipeline.
Definition graph.h:145
int parallel
parallel n: how many values this stage may work on at once.
Definition graph.h:149
std::string text
kString/kOptionalString: the text. Also the field name for at.
Definition graph.h:114
absl::Duration duration
kDuration: how long timeout waits for the next value, or pace leaves between two of them.
Definition graph.h:130
std::string name
Definition graph.h:109
RefId failures
try ... into ref: where a tolerated failure goes, as a status record.
Definition graph.h:147
RefId stream
kStream: the stream then reads after this one.
Definition graph.h:118
std::string carried
Definition graph.h:141
syntax::Constant start
fold: what the first pass carries, the name it is bound to, and the ref that name resolves to inside ...
Definition graph.h:140
long long index
Definition graph.h:123
bool descending
sort: whether the order is reversed.
Definition graph.h:132
long long count
kNumber: the count.
Definition graph.h:112
RefId carry
Definition graph.h:142
LogTail log
kLog/kLogFormat: what was written after the stage name.
Definition graph.h:120
bool ordered
Whether values leave in the order they arrived.
Definition graph.h:152
One resolved statement.
Definition graph.h:265
bool tee
Definition graph.h:280
std::optional< absl::Duration > timeout
Definition graph.h:279
RefId item
kForEach: the value and the index each pass binds.
Definition graph.h:346
std::vector< std::pair< std::string, ExprId > > headers
Definition graph.h:282
ExprId condition
Definition graph.h:357
bool discard
kPipe: -> _, which reads the stream and keeps nothing.
Definition graph.h:302
RefId status
Its outcome, once something has asked for one.
Definition graph.h:288
ExprId message
Definition graph.h:330
std::string code_name
The canonical code fail named outright: fail not_found "...".
Definition graph.h:337
std::string node_map
Definition graph.h:278
std::vector< RefId > subjects
kWait: the outcomes of a wait first of / wait all of, in the order they were written.
Definition graph.h:318
RefId index
Definition graph.h:347
StepKind kind
Definition graph.h:266
std::optional< long long > count
kSkip: skip n port, which claims no reader slot – the count is applied where the stream is produced a...
Definition graph.h:305
bool stop_when
Definition graph.h:358
syntax::Constant start
Definition graph.h:353
std::string slot
kCapture: which slot of the enclosing loop it fills.
Definition graph.h:307
bool race
kWait: whether the first subject to finish is enough.
Definition graph.h:320
StepId target
kCancel: the call to stop.
Definition graph.h:326
RefId winner
kWait: the value a race is – which subject won – for whoever reads it.
Definition graph.h:323
std::vector< StepId > after
What it waits for.
Definition graph.h:272
RefId destination
Definition graph.h:292
LogTail log
kLog: the level, the format and what fills it.
Definition graph.h:340
RefId carry_source
Definition graph.h:352
RefId source
kPipe reads and writes these refs. kCapture reads source.
Definition graph.h:291
absl::flat_hash_map< std::string, RefId > ports
Every port of this call that the flow wired up, by direction:name.
Definition graph.h:286
BodyId body
Definition graph.h:270
ExprId code
kFail.
Definition graph.h:329
syntax::Location location
Definition graph.h:273
std::vector< std::string > forward
Definition graph.h:283
RefId outcome
kWait/kDrain: the outcome read, and whether a bad one is this flow's business or the subject's.
Definition graph.h:311
std::string action
kCall.
Definition graph.h:276
RefId carry
kRepeat: what it carries, where the next pass's value comes from, and when it stops.
Definition graph.h:351
bool tolerant
Definition graph.h:281
int parallel
Definition graph.h:348
std::optional< int > max_iterations
max n, where one was written; nothing means the loop is bounded only by its condition.
Definition graph.h:356
ExprId action_id
Definition graph.h:284
std::string label
The name this step is known by: a bound name, or action/for/if with a #2 after it where one label is ...
Definition graph.h:269
std::string mode
Definition graph.h:277
std::vector< BodyId > bodies
kForEach, kRepeat, kIf: the bodies nested here, in reading order.
Definition graph.h:344
A value the language can write out in full: what a literal is, and what folding a literal expression ...
Definition syntax.h:80
Where a piece of syntax starts, reduced to what every reader of it needs.
Definition syntax.h:49
Base of every syntax node: what it is, and where it started.
Definition syntax.h:204