1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
25pub struct Options {
26 pub autostart: bool,
28 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
82fn 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
107pub 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 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 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}