Coverage Report

Created: 2026-10-07 21:30

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}