A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
inline_pump.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
32#ifndef A11_CONCURRENCY_INLINE_PUMP_H_
33#define A11_CONCURRENCY_INLINE_PUMP_H_
34
35#include <cstddef>
36#include <exception>
37#include <string_view>
38#include <utility>
39
40#include <absl/log/log.h>
41#include <absl/strings/str_cat.h>
42
43#include "a11/exception_guard.h"
44#include "thread/boost_primitives.h"
45
46namespace a11 {
47
58 size_t depth = 0;
60 bool again = false;
61};
62
92template <typename Once>
93void DriveInline(thread::Mutex* absl_nonnull mu,
94 InlinePumpState* absl_nonnull state, std::string_view name,
95 Once&& once, size_t max_depth = 4) {
96 {
97 thread::MutexLock lock(mu);
98 if (state->depth >= max_depth) {
99 state->again = true;
100 return;
101 }
102 ++state->depth;
103 }
104 while (true) {
105 // Every pump body is A11's own -- the store writers and the session pump --
106 // so inside A11 this is a plain call.
107 const absl::Status raised =
108 exception_guard::Attempt([&] { once(); }, absl::StrCat(name, " pump"));
109 if (!raised.ok()) {
110 LOG(ERROR) << raised.message();
111 }
112 thread::MutexLock lock(mu);
113 if (!state->again) {
114 --state->depth;
115 return;
116 }
117 state->again = false;
118 }
119}
120
130inline bool PumpIsDriving(const InlinePumpState& state) {
131 return state.depth > 0;
132}
133
134} // namespace a11
135
136#endif // A11_CONCURRENCY_INLINE_PUMP_H_
Turns what a caller's callable throws into a Status, at the boundary.
thread::Mutex mu
Definition executor.cc:32
std::string name
The name and its colon, which travel together because they always do.
Definition format.cc:49
absl::Status Attempt(Callable &&callable, std::string_view what)
Runs callable, reporting anything it throws as a Status.
Definition exception_guard.h:114
Definition action.cc:63
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.
Definition inline_pump.h:93
bool PumpIsDriving(const InlinePumpState &state)
Whether a DriveInline() turn for this pump is running right now.
Definition inline_pump.h:130
Re-entry bookkeeping for a pump that may be driven from any thread.
Definition inline_pump.h:56
size_t depth
Live DriveInline() calls for this pump, across all threads.
Definition inline_pump.h:58
bool again
Set by a call turned away at the cap; a running turn owes it a pass.
Definition inline_pump.h:60