|
A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
|
Namespaces | |
| namespace | actions |
| namespace | data |
| namespace | exception_guard |
| namespace | flow |
| namespace | internal |
| namespace | net |
| namespace | nodes |
| namespace | obs |
| namespace | percent |
| namespace | service |
| namespace | stores |
| namespace | utf8 |
| namespace | uv |
Classes | |
| class | Future |
| Shared handle to one asynchronous result. More... | |
| struct | InlinePumpState |
| Re-entry bookkeeping for a pump that may be driven from any thread. More... | |
| class | Promise |
| Move-only producer for a Future result. More... | |
| struct | Unit |
| Empty success value used by Future<Unit> operations that return no data. More... | |
Typedefs | |
| using | Task = Future< Unit > |
| Asynchronous operation whose only successful result is completion itself. | |
| using | Duration = absl::Duration |
A11 duration; an alias of absl::Duration (nanosecond-aware). | |
| using | Time = absl::Time |
A11 instant; an alias of absl::Time (nanosecond-aware). | |
Functions | |
| void | Schedule (absl::AnyInvocable< void() && > work, thread::TreeOptions tree_options={}) |
| Schedule work on A11's fiber pool without returning a completion handle. | |
| std::function< void()> | ScheduleCancelable (absl::AnyInvocable< void() && > work, thread::TreeOptions tree_options={}) |
| Schedule work and return an idempotent cooperative-cancellation function. | |
| template<typename T , typename Fn > | |
| auto | ThenAfterWaiting (Future< T > future, absl::Time deadline, Fn transform) -> Future< typename std::invoke_result_t< Fn, const absl::StatusOr< T > & >::value_type > |
Continue with transform when future completes or the deadline expires. | |
| template<typename T > | |
| Future< T > | SubmitWithCancellationHook (absl::AnyInvocable< absl::StatusOr< T >() && > work, std::function< void()> cancellation_hook, thread::TreeOptions tree_options={}) |
| Run work on A11's fiber pool with application-specific cancellation. | |
| template<typename T > | |
| Future< T > | Submit (absl::AnyInvocable< absl::StatusOr< T >() && > work, thread::TreeOptions tree_options) |
| Run status-returning work on the fiber pool and expose its Future. | |
| Task | SubmitTask (absl::AnyInvocable< absl::Status() && > work, thread::TreeOptions tree_options={}) |
| Run an operation that returns only a completion status. | |
| template<typename T > | |
| Future< T > | ReadyFuture (T value) |
Return an already-successful future containing value. | |
| template<typename T > | |
| Future< T > | CompletedFuture (absl::StatusOr< T > result) |
Return an already-completed future containing result. | |
| template<typename T > | |
| Future< T > | FailedFuture (absl::Status status) |
Return an already-failed future containing status. | |
| Task | ReadyTask () |
| Return an already-successful Task. | |
| Task | FailedTask (absl::Status status) |
| Return an already-failed Task. | |
| template<typename T , typename Fn > | |
| auto | Then (const Future< T > &future, Fn transform) -> Future< typename std::invoke_result_t< Fn, const absl::StatusOr< T > & >::value_type > |
Continue with transform once future completes, without a fiber. | |
| template<typename Once > | |
| void | DriveInline (thread::Mutex *absl_nonnull mu, InlinePumpState *absl_nonnull state, std::string_view name, Once &&once, size_t max_depth=4) |
Run once until the pump has nothing left to do without waiting. | |
| bool | PumpIsDriving (const InlinePumpState &state) |
| Whether a DriveInline() turn for this pump is running right now. | |
| absl::Status | RunAllToCompletion (std::vector< absl::AnyInvocable< absl::Status() && > > work, size_t stack_size=256 *1024) |
| Runs each callable on its own fiber and waits for all of them. | |
| template<typename T > | |
| std::vector< absl::StatusOr< T > > | AwaitAll (const std::vector< Future< T > > &futures, absl::Time deadline=absl::InfiniteFuture()) |
| Awaits futures that are already running, and returns every result. | |
| absl::StatusOr< Json > | ParseJson (std::string_view encoded, std::string_view what) |
| Parses JSON text, or explains why it is not JSON. | |
| bool | IsValidUtf8 (std::string_view text) |
Whether text is valid UTF-8. | |
| const Json * | FindUnencodableString (const Json &value) |
| absl::StatusOr< std::string > | DumpJson (const Json &value, std::string_view what) |
| std::string | DumpJsonLossy (const Json &value) |
| absl::StatusOr< std::string > | PackMsgpack (const Json &value, std::string_view what) |
| absl::StatusOr< Json > | UnpackMsgpack (std::string_view encoded, std::string_view what) |
| Decodes MessagePack bytes, or explains why they are not. | |
| const nlohmann::json * | FindUnencodableString (const nlohmann::json &value) |
The first string inside value that is not valid UTF-8, or nullptr. | |
| absl::StatusOr< std::string > | DumpJson (const nlohmann::json &value, std::string_view what) |
Serializes value, rejecting strings that are not valid UTF-8. | |
| std::string | DumpJsonLossy (const nlohmann::json &value) |
Serializes value, replacing anything that is not valid UTF-8. | |
| absl::StatusOr< std::string > | PackMsgpack (const nlohmann::json &value, std::string_view what) |
Encodes value as MessagePack, or says why it cannot be. | |
| absl::StatusCode | StatusCodeFromHttp (int http_code) |
| Maps an HTTP status code to a canonical status code. | |
| int | StatusCodeToHttp (absl::StatusCode code) |
| Maps a canonical status code to its HTTP status code. | |
| absl::StatusCode | StatusCodeFromWebSocket (std::uint16_t close_code) |
| Maps a WebSocket close code to a canonical status code. | |
| std::uint16_t | StatusCodeToWebSocket (absl::StatusCode code) |
| Maps a canonical status code to its WebSocket close code. | |
| absl::Status | MakeStatus (absl::StatusCode code, const std::string &message, const nlohmann::json &details) |
| Builds a status carrying an A11 structured-details payload. | |
| nlohmann::json | StatusDetails (const absl::Status &status) |
| Extracts the structured-details payload from a status. | |
| absl::StatusOr< nlohmann::json > | StatusToJson (const absl::Status &status) |
| Serializes a status (code, message, details) to JSON. | |
| nlohmann::json | StatusToJsonOrEmptyDetails (const absl::Status &status) |
| Serializes a status to JSON, never failing. | |
| absl::StatusOr< absl::Status > | StatusFromJson (const nlohmann::json &value) |
| Reconstructs a status from its JSON representation. | |
| Duration | ZeroDuration () |
| Returns the zero duration. | |
| Duration | InfiniteDuration () |
| Returns the duration greater than every finite duration. | |
| Time | Now () |
| Returns the current wall-clock time from Abseil's clock. | |
| Time | InfiniteFuture () |
| Returns the instant later than every finite time. | |
| Time | InfinitePast () |
| Returns the instant earlier than every finite time. | |
| absl::StatusOr< std::int64_t > | DurationNanoseconds (Duration duration) |
| Converts a duration to whole nanoseconds. | |
| absl::StatusOr< std::int64_t > | TimeNanosecondsSinceEpoch (Time time) |
| Converts an instant to nanoseconds since the Unix epoch. | |
| std::uint64_t | RandomUint64 () |
| 64 random bits from a thread-local generator. | |
| std::string | NewUuid () |
| Generate a random RFC 4122 version 4 UUID in canonical text form. | |
| std::string | NewStreamId (std::string_view prefix="") |
Generate a stream identifier: prefix and 128 random bits as hex. | |
| std::string | NewShortId (int hex_digits=12) |
| Generate a short random identifier as lowercase hexadecimal. | |
Variables | |
| constexpr std::string_view | kStatusDetailsPayloadUrl |
| Type URL identifying A11's status-details payload. | |
| using a11::Duration = typedef absl::Duration |
A11 duration; an alias of absl::Duration (nanosecond-aware).
Asynchronous operation whose only successful result is completion itself.
| using a11::Time = typedef absl::Time |
A11 instant; an alias of absl::Time (nanosecond-aware).
| std::vector< absl::StatusOr< T > > a11::AwaitAll | ( | const std::vector< Future< T > > & | futures, |
| absl::Time | deadline = absl::InfiniteFuture() |
||
| ) |
Awaits futures that are already running, and returns every result.
This function starts no work. It awaits every future, including those after a failed result, and returns results in input order.
| Future< T > a11::CompletedFuture | ( | absl::StatusOr< T > | result | ) |
Return an already-completed future containing result.
| void a11::DriveInline | ( | thread::Mutex *absl_nonnull | mu, |
| InlinePumpState *absl_nonnull | state, | ||
| std::string_view | name, | ||
| Once && | once, | ||
| size_t | max_depth = 4 |
||
| ) |
Run once until the pump has nothing left to do without waiting.
Recursion is bounded rather than forbidden: a call arriving over max_depth asks the turns already running for another pass instead of adding a frame. The cap counts all live drives, so a genuinely concurrent driver can be turned away too. That costs it a pass it need not have made, which is far cheaper than turning away every concurrent caller and handing its work to whoever happens to be inside – an inline drive exists precisely so that a caller does its own work.
Deciding to leave and dropping out of the count happen under one hold of mu, and a call that is turned away sets again under the same lock, so one of the two always observes the other and the requested pass is retained. once is wrapped so an escaping exception cannot leak the count and strand the pump.
| mu | The pump's mutex, guarding state. |
| state | The pump's re-entry bookkeeping. |
| name | Pump name for the diagnostic when once raises. |
| once | One turn of the state machine. Must not block. |
| max_depth | How many live drives to allow before folding further calls into them. |
| absl::StatusOr< std::string > a11::DumpJson | ( | const Json & | value, |
| std::string_view | what | ||
| ) |
| absl::StatusOr< std::string > a11::DumpJson | ( | const nlohmann::json & | value, |
| std::string_view | what | ||
| ) |
Serializes value, rejecting strings that are not valid UTF-8.
Strict on purpose: JSON is defined over text, and a chunk holding arbitrary bytes has to be encoded (base64, as data/json.cc does) rather than smuggled into a string field where a peer's parser would reject it. This is the one caller of nlohmann that wants the error rather than a replacement character.
nlohmann turns its throw into std::abort() in every translation unit compiled -fno-exceptions, and most of A11 is. This file is compiled with exceptions on so that the try below means something – but dump() is a template, and whichever instantiation the linker keeps decides whether a bad string raises or aborts. An aborting one from another TU is a legal choice, and it was the one being made: a read_file over a binary file killed the process with no output at all.
So the strings are checked here, where the answer depends on the bytes rather than on how something was compiled, and the try stays as a second line for everything else dump() can object to. A unit test cannot cover the difference: test binaries are built with exceptions on, so the version under test is the one that throws while the shipped library is the one that aborts.
| std::string a11::DumpJsonLossy | ( | const Json & | value | ) |
| std::string a11::DumpJsonLossy | ( | const nlohmann::json & | value | ) |
Serializes value, replacing anything that is not valid UTF-8.
For a log line or a span attribute, where a lost byte is better than a lost message and there is no peer to reject it.
| absl::StatusOr< std::int64_t > a11::DurationNanoseconds | ( | Duration | duration | ) |
Converts a duration to whole nanoseconds.
| duration | Duration to convert. |
duration is infinite and therefore not representable. | Future< T > a11::FailedFuture | ( | absl::Status | status | ) |
Return an already-failed future containing status.
|
inline |
Return an already-failed Task.
| const Json * a11::FindUnencodableString | ( | const Json & | value | ) |
| const nlohmann::json * a11::FindUnencodableString | ( | const nlohmann::json & | value | ) |
The first string inside value that is not valid UTF-8, or nullptr.
The strings that come from outside are not only at the top level: a response header's value and a directory entry's name are both fields of a record, and both are exactly where something outside this process can put arbitrary bytes.
Iterative, not recursive: this runs on whatever fiber is serializing, and A11's fiber stacks are fixed and small, so the depth of the document must not decide how much stack the check needs.
| Duration a11::InfiniteDuration | ( | ) |
Returns the duration greater than every finite duration.
| Time a11::InfiniteFuture | ( | ) |
Returns the instant later than every finite time.
| Time a11::InfinitePast | ( | ) |
Returns the instant earlier than every finite time.
| bool a11::IsValidUtf8 | ( | std::string_view | text | ) |
Whether text is valid UTF-8.
Public because the answer has to be available before nlohmann is asked. See DumpJson for why asking nlohmann is not safe.
Strict by the definition every other A11 language's string type enforces, which is stricter than "the bytes are shaped like UTF-8": overlong encodings, surrogate halves and anything above U+10FFFF are all rejected. Python's bytes.decode("utf-8"), Kotlin's String(bytes) and JavaScript's TextDecoder with fatal all refuse them, so a chunk carrying them is one a peer cannot read – and a value this side refuses to write is a much better outcome than a session that dies one hop away.
| absl::Status a11::MakeStatus | ( | absl::StatusCode | code, |
| const std::string & | message, | ||
| const nlohmann::json & | details | ||
| ) |
Builds a status carrying an A11 structured-details payload.
The details JSON is attached to the returned absl::Status so it survives the JSON and MessagePack bridges used across the wire.
| code | Canonical status code (use kOk for success). |
| message | Human-readable message. |
| details | Arbitrary structured detail payload. |
code, message and details. | std::string a11::NewShortId | ( | int | hex_digits = 12 | ) |
Generate a short random identifier as lowercase hexadecimal.
Where NewUuid() names something a human may have to correlate across systems, this names something that only has to be unique inside one – an action, and through it every node id derived from it. Twelve digits is 48 random bits: a one-in-a-million chance of a collision needs some 24 000 live ids in one keyspace, far above what a session holds. Eight digits is 32 bits and reaches the same odds at ninety-two, which is why the minimum here is not the minimum you should ask for.
Every character is a hex digit, so the result is always a valid a11::data::ValidateName name and contains neither - nor # – which matters because a node id is <action id>#<port> and every parser of one splits at the single #.
Not cryptographically secure, for the same reason NewUuid() is not.
| hex_digits | Length of the result, clamped to [8, 16]. |
hex_digits characters. | std::string a11::NewStreamId | ( | std::string_view | prefix = "" | ) |
Generate a stream identifier: prefix and 128 random bits as hex.
What every transport names a stream, a session or a data channel with. The prefix says which transport minted it, which is the difference between two ids in one log being confusing and being informative.
Not cryptographically secure, for the same reason NewUuid() is not.
| prefix | Written in front of the digits, e.g. "ws-". May be empty. |
| std::string a11::NewUuid | ( | ) |
Generate a random RFC 4122 version 4 UUID in canonical text form.
The result is 36 characters, lowercase, 8-4-4-4-12 hyphenated, with the version and variant bits set. Randomness comes from a thread-local absl::BitGen, so this is cheap to call in a loop but is not cryptographically secure: use it for identity, never for secrets.
| Time a11::Now | ( | ) |
Returns the current wall-clock time from Abseil's clock.
| absl::StatusOr< std::string > a11::PackMsgpack | ( | const Json & | value, |
| std::string_view | what | ||
| ) |
| absl::StatusOr< std::string > a11::PackMsgpack | ( | const nlohmann::json & | value, |
| std::string_view | what | ||
| ) |
Encodes value as MessagePack, or says why it cannot be.
| absl::StatusOr< nlohmann::json > a11::ParseJson | ( | std::string_view | encoded, |
| std::string_view | what | ||
| ) |
Parses JSON text, or explains why it is not JSON.
| encoded | JSON text to parse. |
| what | Names the document in the error, e.g. "WireMessage JSON". |
|
inline |
Whether a DriveInline() turn for this pump is running right now.
For a completion that has just run inline and has more to do, handing the next pass to the active turn avoids posting one to a worker. Posting would race the caller's next drive and take the work from it. Call this with mu held.
| std::uint64_t a11::RandomUint64 | ( | ) |
64 random bits from a thread-local generator.
For building identifiers. Every id in A11 – session, stream, data channel, WebSocket masking key – should come from here rather than from a locally constructed absl::BitGen: constructing one seeds Randen from the operating system's entropy source, which costs microseconds and, at a site called per connection or per call, costs them repeatedly for no benefit. One generator per thread needs no lock and no seeding after the first call.
Not cryptographically secure, for the same reason NewUuid() is not: use it for identity, never for secrets.
| Future< T > a11::ReadyFuture | ( | T | value | ) |
Return an already-successful future containing value.
|
inline |
Return an already-successful Task.
| absl::Status a11::RunAllToCompletion | ( | std::vector< absl::AnyInvocable< absl::Status() && > > | work, |
| size_t | stack_size = 256 *1024 |
||
| ) |
Runs each callable on its own fiber and waits for all of them.
Each callable runs to completion even after another fails; the function returns the first error. Use it for independent ordered chains such as a sequence of writes followed by a close. stack_size sets each fiber's stack.
| void a11::Schedule | ( | absl::AnyInvocable< void() && > | work, |
| thread::TreeOptions | tree_options = {} |
||
| ) |
Schedule work on A11's fiber pool without returning a completion handle.
Work may await another A11 Future without consuming an OS worker thread.
| std::function< void()> a11::ScheduleCancelable | ( | absl::AnyInvocable< void() && > | work, |
| thread::TreeOptions | tree_options = {} |
||
| ) |
Schedule work and return an idempotent cooperative-cancellation function.
The scheduler retains the root fiber until it has been joined.
| absl::StatusCode a11::StatusCodeFromHttp | ( | int | http_code | ) |
Maps an HTTP status code to a canonical status code.
| http_code | HTTP status code. |
absl::StatusCode; unknown codes map to kUnknown rather than an invented status. | absl::StatusCode a11::StatusCodeFromWebSocket | ( | std::uint16_t | close_code | ) |
Maps a WebSocket close code to a canonical status code.
| close_code | WebSocket close code. |
absl::StatusCode; unknown codes map to kUnknown. | int a11::StatusCodeToHttp | ( | absl::StatusCode | code | ) |
Maps a canonical status code to its HTTP status code.
| code | Canonical status code. |
| std::uint16_t a11::StatusCodeToWebSocket | ( | absl::StatusCode | code | ) |
Maps a canonical status code to its WebSocket close code.
| code | Canonical status code. |
| nlohmann::json a11::StatusDetails | ( | const absl::Status & | status | ) |
Extracts the structured-details payload from a status.
| status | Status to inspect. |
| absl::StatusOr< absl::Status > a11::StatusFromJson | ( | const nlohmann::json & | value | ) |
Reconstructs a status from its JSON representation.
| value | JSON produced by StatusToJson. |
value is not a valid Status document. | absl::StatusOr< nlohmann::json > a11::StatusToJson | ( | const absl::Status & | status | ) |
Serializes a status (code, message, details) to JSON.
| status | Status to serialize. |
| nlohmann::json a11::StatusToJsonOrEmptyDetails | ( | const absl::Status & | status | ) |
Serializes a status to JSON, never failing.
For diagnostics that embed a status in a larger document and have nothing useful to do with an encoding failure. Falls back to the same fields with empty details, so the layout is the one StatusToJson produces either way.
| status | Status to serialize. |
| Future< T > a11::Submit | ( | absl::AnyInvocable< absl::StatusOr< T >() && > | work, |
| thread::TreeOptions | tree_options | ||
| ) |
Run status-returning work on the fiber pool and expose its Future.
|
inline |
Run an operation that returns only a completion status.
| Future< T > a11::SubmitWithCancellationHook | ( | absl::AnyInvocable< absl::StatusOr< T >() && > | work, |
| std::function< void()> | cancellation_hook, | ||
| thread::TreeOptions | tree_options = {} |
||
| ) |
Run work on A11's fiber pool with application-specific cancellation.
The returned Future requests both cancellation_hook and cancellation of the scheduled fiber when Future::Cancel() is called. Use this for operations that must also interrupt an external SDK or transport; ordinary cooperative A11 work can use Submit(). Cancellation remains a request, so consumers must still observe the Future's eventual result.
| auto a11::Then | ( | const Future< T > & | future, |
| Fn | transform | ||
| ) | -> Future< typename std::invoke_result_t<Fn, const absl::StatusOr<T>&>::value_type> |
Continue with transform once future completes, without a fiber.
Submit([f]{ return g(f.Await()); }) needs a worker only because Await() blocks; transform does not. This runs it on whichever thread completes future, or immediately on this one when future is already complete, so the continuation costs no handoff to the worker pool, which is otherwise the dominant cost of the layers built out of it.
Two obligations come with that. transform must not block: it may run on a pooled fiber with a small stack or on a transport's own thread, and it runs before the completing side gets on with its work. And it must not assume a thread; see a11/concurrency/inline_pump.h for why.
| future | The operation to continue from. |
| transform | Called with future's result, returning the continued result. |
transform's result, cancellable through to future. | auto a11::ThenAfterWaiting | ( | Future< T > | future, |
| absl::Time | deadline, | ||
| Fn | transform | ||
| ) | -> Future<typename std::invoke_result_t< Fn, const absl::StatusOr<T>&>::value_type> |
Continue with transform when future completes or the deadline expires.
A ready future is transformed inline; otherwise a fiber waits until deadline. Use Then() when the operation has no deadline.
| future | The operation to continue from. |
| deadline | How long the fiber may wait when future is not already complete. |
| transform | Called with future's result. Runs inline in the ready case and must not block. |
| absl::StatusOr< std::int64_t > a11::TimeNanosecondsSinceEpoch | ( | Time | time | ) |
Converts an instant to nanoseconds since the Unix epoch.
| time | Instant to convert. |
time is infinite and therefore not representable. | absl::StatusOr< nlohmann::json > a11::UnpackMsgpack | ( | std::string_view | encoded, |
| std::string_view | what | ||
| ) |
Decodes MessagePack bytes, or explains why they are not.
| Duration a11::ZeroDuration | ( | ) |
Returns the zero duration.
|
constexpr |
Type URL identifying A11's status-details payload.