/home/runner/work/lyquor/lyquor/platform/node/src/availability.rs
Line | Count | Source |
1 | | //! Deployment image-availability certification worker. |
2 | | //! |
3 | | //! Registrations that enter `Pending` remain there until their content digest |
4 | | //! is admitted or a deployment deadline receives a certified negative verdict. |
5 | | //! This worker watches bartender's `AvailabilityPending` log topic (and periodically sweeps the |
6 | | //! registry so restarts and missed events self-heal), pulls each referenced |
7 | | //! image through the node's [`ImageResolver`], and proposes a digest-scoped |
8 | | //! availability certificate. One certificate admits every deployment that |
9 | | //! references the same immutable content. |
10 | | |
11 | | use std::sync::Arc; |
12 | | use std::time::Duration; |
13 | | |
14 | | use futures::{Stream, StreamExt}; |
15 | | use lyquor_api::{call::CallParams, log::LogClass, store::KVStoreError}; |
16 | | use lyquor_hosting::{ImageResolver, Lyquid}; |
17 | | use lyquor_primitives::{Address, AvailabilityPendingEvent, B256, LyteLog, encode_by_fields}; |
18 | | use lyquor_vm::scheduler::{RunOptions, RunSource}; |
19 | | use tokio_util::sync::CancellationToken; |
20 | | |
21 | | /// How often the worker re-checks the registry for deployments that are still |
22 | | /// pending (bounded re-probe of unretrievable images). |
23 | | const SWEEP_INTERVAL: Duration = Duration::from_secs(15); |
24 | | |
25 | | /// Run the availability watcher until its task group is cancelled. |
26 | 0 | pub async fn run_availability_watcher( |
27 | 0 | bartender: Lyquid, image_resolver: Arc<ImageResolver>, |
28 | 0 | log_stream: impl Stream<Item = Result<LyteLog, KVStoreError>> + Unpin, shutdown: CancellationToken, |
29 | 0 | ) { |
30 | 0 | tokio::select! { |
31 | | biased; |
32 | 0 | () = shutdown.cancelled() => {} |
33 | 0 | () = watch(bartender, image_resolver, log_stream) => {} |
34 | | } |
35 | 0 | } |
36 | | |
37 | 0 | async fn watch( |
38 | 0 | bartender: Lyquid, image_resolver: Arc<ImageResolver>, |
39 | 0 | mut log_stream: impl Stream<Item = Result<LyteLog, KVStoreError>> + Unpin, |
40 | 0 | ) { |
41 | 0 | let mut sweep = tokio::time::interval(SWEEP_INTERVAL); |
42 | 0 | sweep.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); |
43 | | loop { |
44 | 0 | tokio::select! { |
45 | 0 | _ = sweep.tick() => { |
46 | 0 | let pending = match call_bartender::<Vec<AvailabilityPendingEvent>>( |
47 | 0 | &bartender, |
48 | 0 | "get_pending_deployments", |
49 | 0 | encode_by_fields!().into(), |
50 | | ) |
51 | 0 | .await |
52 | | { |
53 | 0 | Ok(pending) => pending, |
54 | 0 | Err(e) => { |
55 | 0 | tracing::debug!(error = ?e, "availability sweep failed to list pending deployments"); |
56 | 0 | continue; |
57 | | } |
58 | | }; |
59 | 0 | for event in pending { |
60 | 0 | try_certify(&bartender, &image_resolver, &event).await; |
61 | | } |
62 | | } |
63 | 0 | log = log_stream.next() => { |
64 | 0 | let Some(log) = log else { break }; |
65 | 0 | let log = match log { |
66 | 0 | Ok(log) => log, |
67 | 0 | Err(error) => { |
68 | 0 | tracing::error!(class = %LogClass::CoreCorruption, error = ?error, "availability log stream failed; watcher stopped"); |
69 | 0 | break; |
70 | | } |
71 | | }; |
72 | 0 | let Some(event) = lyquor_primitives::decode_object::<AvailabilityPendingEvent>(&log.data) else { |
73 | 0 | continue; |
74 | | }; |
75 | 0 | try_certify(&bartender, &image_resolver, &event).await; |
76 | | } |
77 | | } |
78 | | } |
79 | 0 | } |
80 | | |
81 | | /// Pull and verify one pending deployment's image, record the local result, |
82 | | /// then propose the matching positive or negative certificate. |
83 | 0 | async fn try_certify(bartender: &Lyquid, image_resolver: &ImageResolver, event: &AvailabilityPendingEvent) { |
84 | | let AvailabilityPendingEvent { |
85 | 0 | id, |
86 | 0 | nth, |
87 | 0 | image_digest, |
88 | 0 | repo_hint, |
89 | 0 | } = event; |
90 | 0 | if let Some(hint) = repo_hint.as_deref() { |
91 | 0 | image_resolver.remember_repository_hint(*image_digest, hint); |
92 | 0 | } |
93 | | // `ensure_local` stores the pack content-addressed by digest, so success |
94 | | // means this node holds bytes matching the registered digest. |
95 | 0 | let available = match image_resolver.ensure_local(*image_digest).await { |
96 | 0 | Ok(()) => true, |
97 | 0 | Err(e) => { |
98 | 0 | tracing::debug!( |
99 | | lyquid_id = %id, |
100 | | deployment = nth, |
101 | | image_digest = %image_digest, |
102 | | error = ?e, |
103 | | "deployment image not retrievable yet" |
104 | | ); |
105 | 0 | false |
106 | | } |
107 | | }; |
108 | 0 | match call_bartender::<bool>( |
109 | 0 | bartender, |
110 | 0 | "note_image_probe", |
111 | 0 | encode_by_fields!(image_digest: B256 = *image_digest, available: bool = available).into(), |
112 | | ) |
113 | 0 | .await |
114 | | { |
115 | 0 | Ok(_) => {} |
116 | 0 | Err(e) => { |
117 | 0 | tracing::warn!( |
118 | | lyquid_id = %id, |
119 | | deployment = nth, |
120 | | image_digest = %image_digest, |
121 | | class = %LogClass::CoreRuntime, |
122 | | error = ?e, |
123 | | "failed to record deployment image probe" |
124 | | ); |
125 | 0 | return; |
126 | | } |
127 | | } |
128 | 0 | let method = if available { |
129 | 0 | "certify_availability" |
130 | | } else { |
131 | 0 | "certify_unavailability" |
132 | | }; |
133 | | // Every committee node proposes independently. Duplicate positive |
134 | | // proposals become no-ops after admission; negative proposals remain |
135 | | // deployment-scoped and are checked again at their sequenced position. |
136 | 0 | match call_bartender::<bool>( |
137 | 0 | bartender, |
138 | 0 | method, |
139 | 0 | encode_by_fields!(image_digest: B256 = *image_digest).into(), |
140 | | ) |
141 | 0 | .await |
142 | | { |
143 | 0 | Ok(true) => tracing::info!( |
144 | | lyquid_id = %id, |
145 | | deployment = nth, |
146 | | image_digest = %image_digest, |
147 | | available, |
148 | | "deployment image availability certified" |
149 | | ), |
150 | 0 | Ok(false) if available => { |
151 | 0 | tracing::debug!(image_digest = %image_digest, available, "deployment image certificate awaiting quorum"); |
152 | | } |
153 | 0 | Ok(false) => tracing::debug!( |
154 | | image_digest = %image_digest, |
155 | | available, |
156 | | "deployment image certificate awaiting deadline or quorum" |
157 | | ), |
158 | 0 | Err(e) => { |
159 | 0 | tracing::warn!( |
160 | | lyquid_id = %id, |
161 | | deployment = nth, |
162 | | image_digest = %image_digest, |
163 | | method, |
164 | | class = %LogClass::CoreRuntime, |
165 | | error = ?e, |
166 | | "deployment image certification call failed" |
167 | | ); |
168 | | } |
169 | | } |
170 | 0 | } |
171 | | |
172 | 0 | async fn call_bartender<T: for<'a> serde::Deserialize<'a> + Send + 'static>( |
173 | 0 | bartender: &Lyquid, method: &str, input: lyquor_primitives::Bytes, |
174 | 0 | ) -> Result<T, lyquor_vm::instance::Error> { |
175 | 0 | let call = { |
176 | 0 | let bar_ln = bartender.latest_number().await; |
177 | 0 | bartender |
178 | 0 | .call_instance_func_decoded( |
179 | 0 | bar_ln, |
180 | 0 | CallParams::builder() |
181 | 0 | .caller(Address::ZERO) |
182 | 0 | .method(method.into()) |
183 | 0 | .input(input) |
184 | 0 | .build(), |
185 | 0 | RunOptions::new(RunSource::InstanceCall), |
186 | 0 | ) |
187 | 0 | .await |
188 | | }; |
189 | 0 | call.await |
190 | 0 | } |