A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
a11 Namespace Reference

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.
 

Typedef Documentation

◆ Duration

using a11::Duration = typedef absl::Duration

A11 duration; an alias of absl::Duration (nanosecond-aware).

◆ Task

using a11::Task = typedef Future<Unit>

Asynchronous operation whose only successful result is completion itself.

◆ Time

using a11::Time = typedef absl::Time

A11 instant; an alias of absl::Time (nanosecond-aware).

Function Documentation

◆ AwaitAll()

template<typename T >
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.

◆ CompletedFuture()

template<typename T >
Future< T > a11::CompletedFuture ( absl::StatusOr< T >  result)

Return an already-completed future containing result.

◆ DriveInline()

template<typename Once >
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.

Parameters
muThe pump's mutex, guarding state.
stateThe pump's re-entry bookkeeping.
namePump name for the diagnostic when once raises.
onceOne turn of the state machine. Must not block.
max_depthHow many live drives to allow before folding further calls into them.

◆ DumpJson() [1/2]

absl::StatusOr< std::string > a11::DumpJson ( const Json &  value,
std::string_view  what 
)

◆ DumpJson() [2/2]

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.

Why it checks rather than only catching

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.

◆ DumpJsonLossy() [1/2]

std::string a11::DumpJsonLossy ( const Json &  value)

◆ DumpJsonLossy() [2/2]

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.

◆ DurationNanoseconds()

absl::StatusOr< std::int64_t > a11::DurationNanoseconds ( Duration  duration)

Converts a duration to whole nanoseconds.

Parameters
durationDuration to convert.
Returns
The nanosecond count, or an error status when duration is infinite and therefore not representable.

◆ FailedFuture()

template<typename T >
Future< T > a11::FailedFuture ( absl::Status  status)

Return an already-failed future containing status.

◆ FailedTask()

Task a11::FailedTask ( absl::Status  status)
inline

Return an already-failed Task.

◆ FindUnencodableString() [1/2]

const Json * a11::FindUnencodableString ( const Json &  value)

◆ FindUnencodableString() [2/2]

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.

◆ InfiniteDuration()

Duration a11::InfiniteDuration ( )

Returns the duration greater than every finite duration.

◆ InfiniteFuture()

Time a11::InfiniteFuture ( )

Returns the instant later than every finite time.

◆ InfinitePast()

Time a11::InfinitePast ( )

Returns the instant earlier than every finite time.

◆ IsValidUtf8()

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.

◆ MakeStatus()

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.

Parameters
codeCanonical status code (use kOk for success).
messageHuman-readable message.
detailsArbitrary structured detail payload.
Returns
A status combining code, message and details.

◆ NewShortId()

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.

Parameters
hex_digitsLength of the result, clamped to [8, 16].
Returns
A freshly generated identifier of hex_digits characters.

◆ NewStreamId()

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.

Parameters
prefixWritten in front of the digits, e.g. "ws-". May be empty.
Returns
A freshly generated identifier.

◆ NewUuid()

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.

Returns
A freshly generated UUID string.

◆ Now()

Time a11::Now ( )

Returns the current wall-clock time from Abseil's clock.

◆ PackMsgpack() [1/2]

absl::StatusOr< std::string > a11::PackMsgpack ( const Json &  value,
std::string_view  what 
)

◆ PackMsgpack() [2/2]

absl::StatusOr< std::string > a11::PackMsgpack ( const nlohmann::json &  value,
std::string_view  what 
)

Encodes value as MessagePack, or says why it cannot be.

◆ ParseJson()

absl::StatusOr< nlohmann::json > a11::ParseJson ( std::string_view  encoded,
std::string_view  what 
)

Parses JSON text, or explains why it is not JSON.

Parameters
encodedJSON text to parse.
whatNames the document in the error, e.g. "WireMessage JSON".

◆ PumpIsDriving()

bool a11::PumpIsDriving ( const InlinePumpState &  state)
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.

◆ RandomUint64()

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.

Returns
64 uniformly distributed random bits.

◆ ReadyFuture()

template<typename T >
Future< T > a11::ReadyFuture ( T  value)

Return an already-successful future containing value.

◆ ReadyTask()

Task a11::ReadyTask ( )
inline

Return an already-successful Task.

◆ RunAllToCompletion()

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.

◆ Schedule()

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.

◆ ScheduleCancelable()

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.

◆ StatusCodeFromHttp()

absl::StatusCode a11::StatusCodeFromHttp ( int  http_code)

Maps an HTTP status code to a canonical status code.

Parameters
http_codeHTTP status code.
Returns
The corresponding absl::StatusCode; unknown codes map to kUnknown rather than an invented status.

◆ StatusCodeFromWebSocket()

absl::StatusCode a11::StatusCodeFromWebSocket ( std::uint16_t  close_code)

Maps a WebSocket close code to a canonical status code.

Parameters
close_codeWebSocket close code.
Returns
The corresponding absl::StatusCode; unknown codes map to kUnknown.

◆ StatusCodeToHttp()

int a11::StatusCodeToHttp ( absl::StatusCode  code)

Maps a canonical status code to its HTTP status code.

Parameters
codeCanonical status code.
Returns
The corresponding HTTP status code.

◆ StatusCodeToWebSocket()

std::uint16_t a11::StatusCodeToWebSocket ( absl::StatusCode  code)

Maps a canonical status code to its WebSocket close code.

Parameters
codeCanonical status code.
Returns
The corresponding WebSocket close code.

◆ StatusDetails()

nlohmann::json a11::StatusDetails ( const absl::Status &  status)

Extracts the structured-details payload from a status.

Parameters
statusStatus to inspect.
Returns
The attached details JSON, or a null/empty value when none is set.

◆ StatusFromJson()

absl::StatusOr< absl::Status > a11::StatusFromJson ( const nlohmann::json &  value)

Reconstructs a status from its JSON representation.

Parameters
valueJSON produced by StatusToJson.
Returns
The reconstructed status, or an error status when value is not a valid Status document.

◆ StatusToJson()

absl::StatusOr< nlohmann::json > a11::StatusToJson ( const absl::Status &  status)

Serializes a status (code, message, details) to JSON.

Parameters
statusStatus to serialize.
Returns
The JSON representation, or an error status on failure.

◆ StatusToJsonOrEmptyDetails()

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.

Parameters
statusStatus to serialize.
Returns
The JSON representation.

◆ Submit()

template<typename T >
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.

◆ SubmitTask()

Task a11::SubmitTask ( absl::AnyInvocable< absl::Status() && >  work,
thread::TreeOptions  tree_options = {} 
)
inline

Run an operation that returns only a completion status.

◆ SubmitWithCancellationHook()

template<typename T >
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.

◆ Then()

template<typename T , typename Fn >
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.

Parameters
futureThe operation to continue from.
transformCalled with future's result, returning the continued result.
Returns
A Future for transform's result, cancellable through to future.

◆ ThenAfterWaiting()

template<typename T , typename Fn >
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.

Parameters
futureThe operation to continue from.
deadlineHow long the fiber may wait when future is not already complete.
transformCalled with future's result. Runs inline in the ready case and must not block.

◆ TimeNanosecondsSinceEpoch()

absl::StatusOr< std::int64_t > a11::TimeNanosecondsSinceEpoch ( Time  time)

Converts an instant to nanoseconds since the Unix epoch.

Parameters
timeInstant to convert.
Returns
The nanoseconds since epoch, or an error status when time is infinite and therefore not representable.

◆ UnpackMsgpack()

absl::StatusOr< nlohmann::json > a11::UnpackMsgpack ( std::string_view  encoded,
std::string_view  what 
)

Decodes MessagePack bytes, or explains why they are not.

◆ ZeroDuration()

Duration a11::ZeroDuration ( )

Returns the zero duration.

Variable Documentation

◆ kStatusDetailsPayloadUrl

constexpr std::string_view a11::kStatusDetailsPayloadUrl
constexpr
Initial value:
=
"type.a11.dev/status-details+json"

Type URL identifying A11's status-details payload.