A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
callback_scheduler.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_CONCURRENCY_CALLBACK_SCHEDULER_H_
18#define A11_CONCURRENCY_CALLBACK_SCHEDULER_H_
19
20#include <cstddef>
21#include <deque>
22
23#include <absl/functional/any_invocable.h>
24
25#include "thread/boost_primitives.h"
26
27namespace a11::internal {
28
29// A fair, stackless callback pump. Instances must have process lifetime because
30// posted callbacks retain a raw pointer to the scheduler, while queued work
31// owns the state it operates on.
32class CallbackScheduler {
33 public:
34 static constexpr size_t kDefaultMaxCallbacksPerTurn = 64;
35 static constexpr size_t kDefaultMaxConcurrentTurns = 2;
36
45 explicit CallbackScheduler(
46 size_t max_callbacks_per_turn = kDefaultMaxCallbacksPerTurn,
47 size_t max_concurrent_turns = kDefaultMaxConcurrentTurns)
48 : max_callbacks_per_turn_(max_callbacks_per_turn),
49 max_concurrent_turns_(max_concurrent_turns < 2 ? 2
50 : max_concurrent_turns) {
51 }
52
53 CallbackScheduler(const CallbackScheduler&) = delete;
54 CallbackScheduler& operator=(const CallbackScheduler&) = delete;
55
56 void Schedule(absl::AnyInvocable<void() &&> callback);
57
58 private:
59 void Run();
60
61 const size_t max_callbacks_per_turn_;
62 const size_t max_concurrent_turns_;
63 thread::Mutex mu_;
64 std::deque<absl::AnyInvocable<void() &&>> callbacks_ ABSL_GUARDED_BY(mu_);
66 size_t active_turns_ ABSL_GUARDED_BY(mu_) = 0;
67};
68
69} // namespace a11::internal
70
71#endif // A11_CONCURRENCY_CALLBACK_SCHEDULER_H_
absl::StatusOr< RunOutcome > Run(const Source &source, const RunOptions &options)
Compiles source and runs its entry flow to completion.
Definition interpreter.cc:214
Definition callback_scheduler.cc:24
void Schedule(absl::AnyInvocable< void() && > work, thread::TreeOptions tree_options)
Schedule work on A11's fiber pool without returning a completion handle.
Definition executor.cc:47