Symbian platform (C++)
SDK native APIs, runtime and tooling
Loading...
Searching...
No Matches
native_task_owner.h
Go to the documentation of this file.
1// Copyright 2026 The Symbian SDK Authors.
2// Licensed under the Apache License, Version 2.0.
3
4#ifndef SYMBIAN_CONCURRENCY_NATIVE_TASK_OWNER_H_
5#define SYMBIAN_CONCURRENCY_NATIVE_TASK_OWNER_H_
6
7#include <atomic>
8#include <memory>
9#include <optional>
10
11#include "absl/time/time.h"
14
15namespace symbian::concurrency {
16
17// Starts one timer and one real property subscription under a shared owner.
18// Close requests cancellation; Finish settles only after both native adapters
19// have published their drained results. Call Start and Close on the event
20// thread. A worker may cancel the returned Task.
22 public:
24
27
29
30 Task Start(absl::Duration timer_delay,
31 absl::Time deadline = absl::InfiniteFuture()) {
32 if (closed_ || state_) {
33 return FailedTask(absl::FailedPreconditionError(
34 "Native task owner already started or closed"));
35 }
36 if (deadline != absl::InfiniteFuture() && deadline <= absl::Now()) {
37 return FailedTask(absl::DeadlineExceededError("Owner deadline expired"));
38 }
39 auto state = std::make_shared<State>();
40 state_ = state; // Establish owner state before either native submission.
41 state->completion.SetCancellationCallback(
42 [weak = std::weak_ptr<State>(state)] {
43 if (auto owned = weak.lock()) {
44 owned->RequestCancel();
45 }
46 });
47 state->timer = executor_.ScheduleAfter(timer_delay);
48 state->property = executor_.NextProperty();
49 TaskGroup children;
50 children.Add(state->timer);
51 children.Add(
52 Then(state->property,
53 [](const absl::StatusOr<int>& result) -> absl::StatusOr<Unit> {
54 if (!result.ok()) {
55 return result.status();
56 }
57 return Unit{};
58 }));
59 state->children = children.Finish();
60 if (deadline != absl::InfiniteFuture()) {
61 state->alarm = executor_.ScheduleAt(deadline);
62 }
63 state->timer.OnReady([weak = std::weak_ptr<State>(state)](
64 const absl::StatusOr<Unit>& result) {
65 if (!result.ok()) {
66 if (auto owned = weak.lock()) {
67 owned->property.Cancel();
68 }
69 }
70 });
71 state->property.OnReady([weak = std::weak_ptr<State>(state)](
72 const absl::StatusOr<int>& result) {
73 if (!result.ok()) {
74 if (auto owned = weak.lock()) {
75 owned->timer.Cancel();
76 }
77 }
78 });
79 if (state->alarm.valid()) {
80 state->alarm.OnReady([state](const absl::StatusOr<Unit>& result) {
81 state->alarm_result = result;
82 if (result.ok()) {
83 state->timed_out = true;
84 state->timer.Cancel();
85 state->property.Cancel();
86 } else if (result.status().code() != absl::StatusCode::kCancelled) {
87 state->timer.Cancel();
88 state->property.Cancel();
89 }
90 state->Check();
91 });
92 }
93 state->children.OnReady([state](const absl::StatusOr<Unit>& result) {
94 state->child_result = result;
95 if (state->alarm.valid() && !state->alarm.IsReady()) {
96 state->alarm.Cancel();
97 }
98 state->Check();
99 });
100 return state->joined;
101 }
102
104 return state_ ? state_->property : Future<int>();
105 }
106
107 Task Finish() const {
108 return state_ ? state_->joined
109 : FailedTask(absl::FailedPreconditionError(
110 "Native task owner was not started"));
111 }
112
113 void Close() {
114 if (closed_) {
115 return;
116 }
117 closed_ = true;
118 if (state_) {
119 state_->RequestCancel();
120 }
121 }
122
123 private:
124 struct State {
125 void RequestCancel() {
126 cancel_requested.store(true, std::memory_order_release);
127 timer.Cancel();
128 property.Cancel();
129 alarm.Cancel();
130 }
131
132 void Check() {
133 if (!child_result || (alarm.valid() && !alarm_result)) {
134 return;
135 }
136 if (timed_out) {
137 completion.SetError(
138 absl::DeadlineExceededError("Owner deadline expired"));
139 } else if (alarm_result && !alarm_result->ok() &&
140 alarm_result->status().code() !=
141 absl::StatusCode::kCancelled) {
142 completion.SetError(alarm_result->status());
143 } else if (!child_result->ok()) {
144 completion.SetError(child_result->status());
145 } else if (cancel_requested.load(std::memory_order_acquire)) {
146 completion.SetError(absl::CancelledError("Native task owner closed"));
147 } else {
148 completion.SetValue(Unit{});
149 }
150 }
151
152 Promise<Unit> completion;
153 Task joined = completion.future();
154 Task timer;
155 Future<int> property;
156 Task alarm;
157 Task children;
158 std::optional<absl::StatusOr<Unit>> child_result;
159 std::optional<absl::StatusOr<Unit>> alarm_result;
160 std::atomic<bool> cancel_requested{false};
161 bool timed_out = false;
162 };
163
164 EventExecutor& executor_;
165 std::shared_ptr<State> state_;
166 bool closed_ = false;
167};
168
169} // namespace symbian::concurrency
170
171#endif // SYMBIAN_CONCURRENCY_NATIVE_TASK_OWNER_H_
Definition bounded_channel.h:38
Definition event_executor.h:27
Task ScheduleAfter(absl::Duration delay)
Definition event_executor.h:68
Task ScheduleAt(absl::Time deadline)
Definition event_executor.h:72
Future< int > NextProperty()
Definition event_executor.h:86
void OnReady(std::function< void(const absl::StatusOr< T > &)> callback) const
Definition future.h:215
Definition native_task_owner.h:21
NativeTaskOwner(const NativeTaskOwner &)=delete
Task Start(absl::Duration timer_delay, absl::Time deadline=absl::InfiniteFuture())
Definition native_task_owner.h:30
~NativeTaskOwner()
Definition native_task_owner.h:28
NativeTaskOwner(EventExecutor &executor)
Definition native_task_owner.h:23
Future< int > property_result() const
Definition native_task_owner.h:103
Task Finish() const
Definition native_task_owner.h:107
NativeTaskOwner & operator=(const NativeTaskOwner &)=delete
void Close()
Definition native_task_owner.h:113
Definition task_group.h:18
bool Add(Task child)
Definition task_group.h:32
Definition bounded_channel.h:32
Task FailedTask(absl::Status error)
Definition future.h:276
Future< Unit > Task
Definition future.h:270
auto Then(const Future< T > &future, Fn transform) -> Future< typename std::invoke_result_t< Fn, const absl::StatusOr< T > & >::value_type >
Definition future.h:281
Definition future.h:33