Skip to main content

isml_lsp/
sync.rs

1//! Uploader state via LSP progress, Zed's only status channel for an extension.
2//! Only uploading and failed are shown: a progress item that never ends reads as a stuck spinner.
3
4use crossbeam_channel::{SendError, Sender};
5use std::path::{Path, PathBuf};
6use std::time::Duration;
7
8use lsp_server::Message;
9use lsp_types::notification::{Notification, Progress, ShowMessage};
10use lsp_types::request::{Request as RequestTrait, WorkDoneProgressCreate};
11use lsp_types::{
12    MessageType, NumberOrString, ProgressParams, ProgressParamsValue, ShowMessageParams,
13    WorkDoneProgress, WorkDoneProgressBegin, WorkDoneProgressCreateParams, WorkDoneProgressEnd,
14    WorkDoneProgressReport,
15};
16
17const TOKEN: &str = "sfcc/sync";
18const POLL: Duration = Duration::from_millis(1500);
19
20pub use sfcc_core::state::upload::{State, Status};
21use sfcc_core::state::{is_within, now_seconds, upload};
22
23/// `upload` in the initialization options.
24#[derive(Debug, Clone, Copy, PartialEq, Eq)]
25pub struct Options {
26    /// Start `sfcc-upload` for a workspace with a dw.json and no watcher. Off by default.
27    pub autostart: bool,
28    /// A notification when an upload fails, and when it recovers. Never for one that works.
29    pub notify: bool,
30}
31
32impl Options {
33    pub fn from_settings(settings: Option<&serde_json::Value>) -> Options {
34        let upload = settings.and_then(|settings| settings.get("upload"));
35        let flag = |name: &str, default: bool| {
36            upload
37                .and_then(|upload| upload.get(name))
38                .and_then(serde_json::Value::as_bool)
39                .unwrap_or(default)
40        };
41        Options {
42            autostart: flag("autostart", false),
43            notify: flag("notify", true),
44        }
45    }
46}
47
48pub fn report(roots: Vec<PathBuf>, options: Options, sender: Sender<Message>) {
49    std::thread::spawn(move || {
50        if create_token(&sender).is_err() {
51            return;
52        }
53        if options.autostart && current(&roots).is_none() {
54            for (level, text) in autostart(&roots, "sfcc-upload", "upload") {
55                if show(&sender, level, text).is_err() {
56                    return;
57                }
58            }
59        }
60        let mut shown: (Option<State>, String) = (None, String::new());
61        let mut failing = false;
62        loop {
63            let status = current(&roots);
64            let now = (
65                status.as_ref().map(|status| status.state),
66                status.as_ref().map(describe).unwrap_or_default(),
67            );
68            if now != shown && announce(&sender, shown.0, now.0, now.1.clone()).is_err() {
69                return;
70            }
71            if let Some((level, text)) = toast(status.as_ref(), &mut failing) {
72                if options.notify && show(&sender, level, text).is_err() {
73                    return;
74                }
75            }
76            shown = now;
77            std::thread::sleep(POLL);
78        }
79    });
80}
81
82/// Only the edges: into a failure, and out of it. Retries in between stay in the status bar.
83fn toast(status: Option<&Status>, failing: &mut bool) -> Option<(MessageType, String)> {
84    let status = match status {
85        Some(status) if status.state != State::Stopped => status,
86        _ => {
87            *failing = false;
88            return None;
89        }
90    };
91    match (status.state, *failing) {
92        (State::Failed, false) => {
93            *failing = true;
94            Some((MessageType::ERROR, format!("SFCC {}", describe(status))))
95        }
96        (State::Synced, true) => {
97            *failing = false;
98            Some((
99                MessageType::INFO,
100                format!("SFCC upload back in sync with {}", status.hostname),
101            ))
102        }
103        _ => None,
104    }
105}
106
107/// `<program> start` in every root with a dw.json where the tools look for one. Says nothing
108/// when it started, or was already running. `setting` is the option that asked for it.
109pub fn autostart(roots: &[PathBuf], program: &str, setting: &str) -> Vec<(MessageType, String)> {
110    let mut said = Vec::new();
111    for root in roots {
112        if !root.join("dw.json").is_file() && !root.join("source").join("dw.json").is_file() {
113            continue;
114        }
115        let mut command = std::process::Command::new(program);
116        command
117            .arg("start")
118            .current_dir(root)
119            .stdin(std::process::Stdio::null())
120            .stdout(std::process::Stdio::null())
121            .stderr(std::process::Stdio::piped());
122        #[cfg(windows)]
123        {
124            use std::os::windows::process::CommandExt;
125            // CREATE_NO_WINDOW: no console flashing up from an editor.
126            command.creation_flags(0x0800_0000);
127        }
128        match command.output() {
129            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
130                said.push((
131                    MessageType::WARNING,
132                    format!("SFCC: {setting}.autostart is on, but {program} is not installed"),
133                ));
134                break;
135            }
136            Err(error) => said.push((
137                MessageType::WARNING,
138                format!("SFCC: could not start {program}: {error}"),
139            )),
140            Ok(output) if !output.status.success() => {
141                let stderr = String::from_utf8_lossy(&output.stderr);
142                if !stderr.contains("already running") {
143                    let reason = stderr
144                        .lines()
145                        .rev()
146                        .find(|line| !line.trim().is_empty())
147                        .unwrap_or("it exited with an error");
148                    said.push((
149                        MessageType::WARNING,
150                        format!("SFCC: {program} did not start: {}", reason.trim()),
151                    ));
152                }
153            }
154            Ok(_) => {}
155        }
156    }
157    said
158}
159
160pub fn show(
161    sender: &Sender<Message>,
162    level: MessageType,
163    text: String,
164) -> Result<(), SendError<Message>> {
165    sender.send(Message::Notification(lsp_server::Notification::new(
166        ShowMessage::METHOD.to_string(),
167        ShowMessageParams {
168            typ: level,
169            message: text,
170        },
171    )))
172}
173
174pub fn current(roots: &[PathBuf]) -> Option<Status> {
175    let mut best: Option<Status> = None;
176    for status in upload::all() {
177        if !covers(&status, roots) || is_stale(&status) {
178            continue;
179        }
180        if best.as_ref().is_none_or(|found| found.at < status.at) {
181            best = Some(status);
182        }
183    }
184    best
185}
186
187fn covers(status: &Status, roots: &[PathBuf]) -> bool {
188    let watched = Path::new(&status.cartridges);
189    roots.iter().any(|root| is_within(watched, root))
190}
191
192fn is_stale(status: &Status) -> bool {
193    now_seconds() - status.at > sfcc_core::state::STALE_SECONDS
194}
195
196fn describe(status: &Status) -> String {
197    match status.state {
198        State::Uploading => match status.files {
199            0 => format!("uploading to {}", status.hostname),
200            1 => format!("uploading 1 file to {}", status.hostname),
201            files => format!("uploading {files} files to {}", status.hostname),
202        },
203        State::Failed => match &status.detail {
204            Some(detail) => format!("upload failed — {detail}"),
205            None => "upload failed".to_string(),
206        },
207        State::Synced | State::Stopped => String::new(),
208    }
209}
210
211fn worth_showing(state: Option<State>) -> bool {
212    matches!(state, Some(State::Uploading) | Some(State::Failed))
213}
214
215fn announce(
216    sender: &Sender<Message>,
217    was: Option<State>,
218    now: Option<State>,
219    message: String,
220) -> Result<(), SendError<Message>> {
221    let value = match (worth_showing(was), worth_showing(now)) {
222        (false, true) => WorkDoneProgress::Begin(WorkDoneProgressBegin {
223            title: "SFCC".to_string(),
224            message: Some(message),
225            ..Default::default()
226        }),
227        (true, true) => WorkDoneProgress::Report(WorkDoneProgressReport {
228            message: Some(message),
229            ..Default::default()
230        }),
231        (true, false) => WorkDoneProgress::End(WorkDoneProgressEnd { message: None }),
232        (false, false) => return Ok(()),
233    };
234    sender.send(Message::Notification(lsp_server::Notification::new(
235        Progress::METHOD.to_string(),
236        ProgressParams {
237            token: NumberOrString::String(TOKEN.to_string()),
238            value: ProgressParamsValue::WorkDone(value),
239        },
240    )))
241}
242
243fn create_token(sender: &Sender<Message>) -> Result<(), SendError<Message>> {
244    sender.send(Message::Request(lsp_server::Request::new(
245        lsp_server::RequestId::from(TOKEN.to_string()),
246        WorkDoneProgressCreate::METHOD.to_string(),
247        WorkDoneProgressCreateParams {
248            token: NumberOrString::String(TOKEN.to_string()),
249        },
250    )))
251}
252
253#[cfg(test)]
254mod tests {
255    use super::*;
256
257    fn status(state: State, at: i64) -> Status {
258        Status {
259            state,
260            cartridges: "/repo/source/cartridges".into(),
261            hostname: "sbx-001.example.com".into(),
262            code_version: "version1".into(),
263            files: 7,
264            detail: None,
265            at,
266        }
267    }
268
269    #[test]
270    fn reads_what_the_uploader_writes() {
271        let body = r#"{"state":"uploading","cartridges":"/repo/cartridges",
272                       "hostname":"host","code_version":"version1","files":3,"at":1700000000}"#;
273        let found: Status = serde_json::from_str(body).unwrap();
274        assert_eq!(found.state, State::Uploading);
275        assert_eq!(found.files, 3);
276    }
277
278    #[test]
279    fn claims_only_a_watcher_inside_an_open_folder() {
280        let found = status(State::Synced, 0);
281        assert!(covers(&found, &[PathBuf::from("/repo")]));
282        assert!(!covers(&found, &[PathBuf::from("/elsewhere")]));
283    }
284
285    #[test]
286    fn treats_a_watcher_that_stopped_beating_as_gone() {
287        assert!(is_stale(&status(State::Synced, 0)));
288        assert!(!is_stale(&status(State::Synced, now_seconds())));
289    }
290
291    #[test]
292    fn shows_only_what_is_worth_interrupting_for() {
293        assert!(worth_showing(Some(State::Uploading)));
294        assert!(worth_showing(Some(State::Failed)));
295        assert!(!worth_showing(Some(State::Synced)));
296        assert!(!worth_showing(None));
297    }
298
299    #[test]
300    fn counts_files_in_the_message_it_shows() {
301        assert!(describe(&status(State::Uploading, 0)).starts_with("uploading 7 files"));
302        let mut one = status(State::Uploading, 0);
303        one.files = 1;
304        assert!(describe(&one).starts_with("uploading 1 file to"));
305        assert!(describe(&status(State::Synced, 0)).is_empty());
306    }
307
308    #[test]
309    fn notifies_only_on_the_way_into_a_failure_and_out_of_it() {
310        let mut failing = false;
311        let mut failed = status(State::Failed, 0);
312        failed.detail = Some("HTTP 503".into());
313
314        assert!(toast(Some(&status(State::Uploading, 0)), &mut failing).is_none());
315        assert!(toast(Some(&status(State::Synced, 0)), &mut failing).is_none());
316
317        let (level, text) = toast(Some(&failed), &mut failing).unwrap();
318        assert_eq!(level, MessageType::ERROR);
319        assert_eq!(text, "SFCC upload failed — HTTP 503");
320        // Every retry after it is the same failure.
321        assert!(toast(Some(&failed), &mut failing).is_none());
322        assert!(toast(Some(&status(State::Uploading, 0)), &mut failing).is_none());
323        assert!(toast(Some(&failed), &mut failing).is_none());
324
325        let (level, text) = toast(Some(&status(State::Synced, 0)), &mut failing).unwrap();
326        assert_eq!(level, MessageType::INFO);
327        assert!(text.contains("back in sync with sbx-001.example.com"));
328        assert!(toast(Some(&status(State::Synced, 0)), &mut failing).is_none());
329    }
330
331    #[test]
332    fn a_watcher_that_went_away_is_not_a_recovery() {
333        let mut failing = false;
334        toast(Some(&status(State::Failed, 0)), &mut failing);
335        assert!(toast(None, &mut failing).is_none());
336        assert!(toast(Some(&status(State::Synced, 0)), &mut failing).is_none());
337    }
338
339    #[test]
340    fn reads_the_upload_options() {
341        let on = serde_json::json!({ "upload": { "autostart": true } });
342        assert_eq!(
343            Options::from_settings(Some(&on)),
344            Options {
345                autostart: true,
346                notify: true
347            }
348        );
349        let quiet = serde_json::json!({ "upload": { "notify": false } });
350        assert_eq!(
351            Options::from_settings(Some(&quiet)),
352            Options {
353                autostart: false,
354                notify: false
355            }
356        );
357        assert_eq!(
358            Options::from_settings(None),
359            Options {
360                autostart: false,
361                notify: true
362            }
363        );
364    }
365
366    #[test]
367    fn puts_the_reason_in_a_failure() {
368        let mut failed = status(State::Failed, 0);
369        failed.detail = Some("502 Bad Gateway".into());
370        assert_eq!(describe(&failed), "upload failed — 502 Bad Gateway");
371    }
372}