evm_oracle_state/pending/
source.rs1use std::{future::Future, pin::Pin, time::SystemTime};
4
5use tokio::{sync::watch, task::JoinHandle};
6
7use super::{PendingOracleEvent, PendingOracleRuntime, PendingOracleSource, PendingOracleSourceId};
8
9pub type PendingOracleSourceFuture =
11 Pin<Box<dyn Future<Output = Result<(), PendingOracleSourceError>> + Send + 'static>>;
12
13#[derive(Clone, Debug, PartialEq, Eq)]
15pub struct PendingOracleSourceDescriptor {
16 pub source: PendingOracleSource,
18 pub id: PendingOracleSourceId,
20}
21
22impl PendingOracleSourceDescriptor {
23 pub fn new(source: PendingOracleSource, id: PendingOracleSourceId) -> Self {
25 Self { source, id }
26 }
27}
28
29#[non_exhaustive]
31#[derive(Clone, Copy, Debug, PartialEq, Eq)]
32pub enum PendingOracleSourceState {
33 Connecting,
35 Ready,
37 Degraded,
39 Stopped,
41}
42
43#[derive(Clone, Debug, PartialEq, Eq)]
45pub struct PendingOracleSourceHealth {
46 pub descriptor: PendingOracleSourceDescriptor,
48 pub state: PendingOracleSourceState,
50 pub last_transport_message_at: Option<SystemTime>,
52 pub last_candidate_at: Option<SystemTime>,
54 pub coverage_gap_count: u64,
56 pub last_error: Option<String>,
58}
59
60impl PendingOracleSourceHealth {
61 fn connecting(descriptor: PendingOracleSourceDescriptor) -> Self {
62 Self {
63 descriptor,
64 state: PendingOracleSourceState::Connecting,
65 last_transport_message_at: None,
66 last_candidate_at: None,
67 coverage_gap_count: 0,
68 last_error: None,
69 }
70 }
71}
72
73#[derive(Clone, Debug, PartialEq, Eq)]
75pub struct PendingOracleCoverageGap {
76 pub descriptor: PendingOracleSourceDescriptor,
78 pub observed_at: SystemTime,
80 pub reason: String,
82}
83
84#[non_exhaustive]
86#[derive(Clone, Debug, PartialEq, Eq, thiserror::Error)]
87pub enum PendingOracleSourceError {
88 #[error("pending oracle updates are not enabled on this runtime")]
90 PendingUpdatesDisabled,
91 #[error("pending source {transport:?} is not enabled")]
93 SourceDisabled {
94 transport: PendingOracleSource,
96 },
97 #[error("pending source id {id} is already running")]
99 DuplicateSource {
100 id: super::PendingOracleSourceId,
102 },
103 #[error("a Tokio runtime is required to start a pending oracle source")]
105 RuntimeUnavailable,
106 #[error("pending source transport failed: {0}")]
108 Transport(String),
109 #[error("pending source task failed: {0}")]
111 Task(String),
112}
113
114pub trait PendingOracleCandidateSource: Send + 'static {
116 fn descriptor(&self) -> PendingOracleSourceDescriptor;
118
119 fn run(
121 self: Box<Self>,
122 sink: PendingOracleSourceSink,
123 shutdown: watch::Receiver<bool>,
124 ) -> PendingOracleSourceFuture;
125}
126
127#[derive(Clone, Debug)]
129pub struct PendingOracleSourceSink {
130 runtime: PendingOracleRuntime,
131 descriptor: PendingOracleSourceDescriptor,
132}
133
134impl PendingOracleSourceSink {
135 pub(crate) fn new(
136 runtime: PendingOracleRuntime,
137 descriptor: PendingOracleSourceDescriptor,
138 ) -> Self {
139 Self {
140 runtime,
141 descriptor,
142 }
143 }
144
145 pub fn ready(&self) {
147 self.runtime
148 .update_source_health(&self.descriptor, |health| {
149 health.state = PendingOracleSourceState::Ready;
150 health.last_error = None;
151 });
152 }
153
154 pub fn transport_message(&self) {
156 self.runtime
157 .mutate_source_health(&self.descriptor, |health| {
158 health.last_transport_message_at = Some(SystemTime::now());
159 });
160 }
161
162 pub fn candidate(&self) {
164 self.runtime
165 .mutate_source_health(&self.descriptor, |health| {
166 health.last_candidate_at = Some(SystemTime::now());
167 });
168 }
169
170 pub fn coverage_gap(&self, reason: impl Into<String>) {
172 let reason = reason.into();
173 self.runtime
174 .update_source_health(&self.descriptor, |health| {
175 health.state = PendingOracleSourceState::Degraded;
176 health.coverage_gap_count = health.coverage_gap_count.saturating_add(1);
177 health.last_error = Some(reason.clone());
178 });
179 self.runtime
180 .publisher()
181 .publish(PendingOracleEvent::CoverageGap(PendingOracleCoverageGap {
182 descriptor: self.descriptor.clone(),
183 observed_at: SystemTime::now(),
184 reason,
185 }));
186 }
187
188 pub fn reconnecting(&self) {
190 self.runtime
191 .update_source_health(&self.descriptor, |health| {
192 health.state = PendingOracleSourceState::Connecting;
193 });
194 }
195
196 pub fn runtime(&self) -> &PendingOracleRuntime {
198 &self.runtime
199 }
200
201 pub fn descriptor(&self) -> &PendingOracleSourceDescriptor {
203 &self.descriptor
204 }
205
206 pub(crate) fn connecting(&self) {
207 self.runtime
208 .insert_source_health(PendingOracleSourceHealth::connecting(
209 self.descriptor.clone(),
210 ));
211 }
212
213 pub(crate) fn stopped(&self) {
214 self.runtime
215 .update_source_health(&self.descriptor, |health| {
216 health.state = PendingOracleSourceState::Stopped;
217 });
218 }
219
220 pub(crate) fn failed(&self, error: &PendingOracleSourceError) {
221 let reason = error.to_string();
222 self.runtime
223 .update_source_health(&self.descriptor, |health| {
224 health.state = PendingOracleSourceState::Degraded;
225 health.coverage_gap_count = health.coverage_gap_count.saturating_add(1);
226 health.last_error = Some(reason.clone());
227 });
228 self.runtime
229 .publisher()
230 .publish(PendingOracleEvent::CoverageGap(PendingOracleCoverageGap {
231 descriptor: self.descriptor.clone(),
232 observed_at: SystemTime::now(),
233 reason,
234 }));
235 }
236}
237
238#[derive(Debug)]
240pub struct PendingOracleSourceSession {
241 descriptor: PendingOracleSourceDescriptor,
242 shutdown: watch::Sender<bool>,
243 task: Option<JoinHandle<Result<(), PendingOracleSourceError>>>,
244}
245
246impl PendingOracleSourceSession {
247 pub(crate) fn new(
248 descriptor: PendingOracleSourceDescriptor,
249 shutdown: watch::Sender<bool>,
250 task: JoinHandle<Result<(), PendingOracleSourceError>>,
251 ) -> Self {
252 Self {
253 descriptor,
254 shutdown,
255 task: Some(task),
256 }
257 }
258
259 pub fn descriptor(&self) -> &PendingOracleSourceDescriptor {
261 &self.descriptor
262 }
263
264 pub fn stop(&self) {
266 let _ = self.shutdown.send(true);
267 }
268
269 pub async fn join(mut self) -> Result<(), PendingOracleSourceError> {
271 let Some(task) = self.task.take() else {
272 return Ok(());
273 };
274 task.await
275 .map_err(|error| PendingOracleSourceError::Task(error.to_string()))?
276 }
277}
278
279impl Drop for PendingOracleSourceSession {
280 fn drop(&mut self) {
281 let _ = self.shutdown.send(true);
282 if let Some(task) = self.task.take() {
283 task.abort();
284 }
285 }
286}