1use serde::{Deserialize, Serialize};
7use std::env;
8use std::path::PathBuf;
9use std::time::Duration;
10
11fn timeout() -> Duration {
14 std::env::var("PACKSET_TIMEOUT_MS")
15 .ok()
16 .and_then(|v| v.trim().parse::<u64>().ok())
17 .filter(|ms| *ms > 0)
18 .map_or(Duration::from_secs(30), Duration::from_millis)
19}
20
21fn path_seg(id: &str) -> String {
22 let mut out = String::with_capacity(id.len());
23 for b in id.bytes() {
24 match b {
25 b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
26 out.push(b as char)
27 }
28 _ => out.push_str(&format!("%{b:02X}")),
29 }
30 }
31 out
32}
33
34pub const DEFAULT_PORT: u16 = 8761;
37
38pub fn load_seat_env() {
41 let Some(home) = env::var_os("HOME") else {
42 return;
43 };
44 let path = PathBuf::from(home).join(".config/ljos/env");
45 let Ok(text) = std::fs::read_to_string(path) else {
46 return;
47 };
48 for line in text.lines() {
49 let line = line.trim();
50 if line.is_empty() || line.starts_with('#') {
51 continue;
52 }
53 let Some((k, v)) = line.split_once('=') else {
54 continue;
55 };
56 let k = k.trim();
57 if k.is_empty() || env::var_os(k).is_some() {
58 continue;
59 }
60 env::set_var(k, v.trim());
61 }
62}
63
64#[must_use]
69pub fn resolved_workspace() -> String {
70 load_seat_env();
71 env::var("PACKSET_WORKSPACE")
72 .ok()
73 .map(|w| w.trim().to_string())
74 .filter(|w| !w.is_empty())
75 .unwrap_or_else(|| "seat".to_string())
76}
77
78#[must_use]
80pub fn default_port() -> u16 {
81 env::var("PACKSET_PORT")
82 .or_else(|_| env::var("GROK_MEM_PORT"))
83 .ok()
84 .and_then(|raw| raw.trim().parse().ok())
85 .unwrap_or(DEFAULT_PORT)
86}
87
88#[derive(Debug, thiserror::Error)]
89pub enum Error {
90 #[error("packset url missing")]
91 NoUrl,
92 #[error("http: {0}")]
93 Http(#[from] Box<ureq::Error>),
94 #[error("io: {0}")]
95 Io(#[from] std::io::Error),
96 #[error("json: {0}")]
97 Json(#[from] serde_json::Error),
98 #[error("bad response: {0}")]
99 Bad(String),
100}
101
102#[derive(Debug, Clone)]
103pub struct PacksetClient {
104 base: String,
105 workspace: Option<String>,
108}
109
110#[derive(Debug, Clone, Serialize, Deserialize)]
111pub struct Hit {
112 pub id: Option<String>,
113 pub text: String,
114 #[serde(default)]
115 pub score: f64,
116 #[serde(default)]
117 pub kind: String,
118 #[serde(default)]
121 pub ts: Option<String>,
122 #[serde(default)]
127 pub entities: Vec<String>,
128 #[serde(default)]
129 pub ballots: Option<u32>,
130 #[serde(default)]
131 pub of: Option<u32>,
132}
133
134fn refused(url: &str, e: ureq::Error) -> Error {
136 match e {
137 ureq::Error::Status(code, response) => {
138 let text = response.into_string().unwrap_or_default();
139 let reason = serde_json::from_str::<serde_json::Value>(&text)
140 .ok()
141 .and_then(|v| v.get("error").and_then(|r| r.as_str()).map(str::to_string))
142 .unwrap_or(text);
143 let reason = reason.trim();
144 if reason.is_empty() {
145 Error::Bad(format!("{url}: status code {code}"))
146 } else {
147 Error::Bad(format!("{url}: {code}: {reason}"))
148 }
149 }
150 other => Error::Http(Box::new(other)),
151 }
152}
153
154impl PacksetClient {
155 pub fn new(base: impl Into<String>) -> Self {
156 let mut base = base.into();
157 while base.ends_with('/') {
158 base.pop();
159 }
160 Self {
161 base,
162 workspace: None,
163 }
164 }
165
166 #[must_use]
170 pub fn with_workspace(mut self, workspace: impl Into<String>) -> Self {
171 let workspace = workspace.into();
172 self.workspace = (!workspace.is_empty()).then_some(workspace);
173 self
174 }
175
176 pub fn from_env() -> Result<Self, Error> {
181 load_seat_env();
182 let url = env::var("PACKSET_URL")
183 .or_else(|_| env::var("INSIDE_MEMORY_URL"))
184 .ok()
185 .filter(|url| !url.is_empty());
186 match url {
187 Some(url) if url == "off" => Err(Error::NoUrl),
188 Some(url) => Ok(Self::new(url)),
189 None => Ok(Self::new(format!("http://127.0.0.1:{}", default_port()))),
190 }
191 }
192
193 pub fn base(&self) -> &str {
194 &self.base
195 }
196
197 pub fn workspace(&self) -> String {
198 if let Some(w) = &self.workspace {
199 return w.clone();
200 }
201 if let Ok(w) = env::var("PACKSET_WORKSPACE") {
202 if !w.is_empty() {
203 return w;
204 }
205 }
206 let cwd = env::var("GROKOS_WORKSPACE")
207 .ok()
208 .map(|s| s.trim().to_string())
209 .filter(|s| !s.is_empty())
210 .map(std::path::PathBuf::from)
211 .or_else(|| env::current_dir().ok())
212 .unwrap_or_else(|| std::path::PathBuf::from("."));
213 self.workspace_for_cwd(&cwd)
214 }
215
216 pub fn workspace_for_cwd(&self, cwd: &std::path::Path) -> String {
218 let abs = cwd.canonicalize().unwrap_or_else(|_| cwd.to_path_buf());
219 let url = format!("{}/v1/identity", self.base);
220 let body = ureq::get(&url)
221 .query("cwd", abs.to_string_lossy().as_ref())
222 .timeout(timeout())
223 .call()
224 .ok()
225 .and_then(|r| r.into_string().ok());
226 if let Some(body) = body {
227 if let Ok(val) = serde_json::from_str::<serde_json::Value>(&body) {
228 if let Some(ws) = val.get("workspace").and_then(|v| v.as_str()) {
229 if !ws.is_empty() {
230 return ws.to_string();
231 }
232 }
233 }
234 }
235 format!("dir:{}", abs.display())
236 }
237
238 pub fn health(&self) -> Result<String, Error> {
239 let url = format!("{}/health", self.base);
240 let body = ureq::get(&url)
241 .timeout(timeout())
242 .call()
243 .map_err(|e| refused(&url, e))?
244 .into_string()?;
245 Ok(body)
246 }
247
248 pub fn get_atom(&self, workspace: &str, id: &str) -> Result<serde_json::Value, Error> {
249 let encoded = path_seg(id);
250 let url = format!("{}/v1/atoms/{encoded}", self.base);
251 let resp = match ureq::get(&url)
252 .query("workspace", workspace)
253 .timeout(timeout())
254 .call()
255 {
256 Ok(resp) => resp,
257 Err(ureq::Error::Status(404, _)) => {
258 return Err(Error::Bad(format!("no atom {id}")));
259 }
260 Err(e) => return Err(Error::Http(Box::new(e))),
261 };
262 Ok(resp.into_json()?)
263 }
264
265 pub fn list_atoms(&self, workspace: &str) -> Result<Vec<serde_json::Value>, Error> {
266 self.atoms_as_of(workspace, None)
267 }
268
269 pub fn atoms_as_of(
275 &self,
276 workspace: &str,
277 as_of: Option<&str>,
278 ) -> Result<Vec<serde_json::Value>, Error> {
279 let url = format!("{}/v1/atoms", self.base);
280 let mut req = ureq::get(&url)
281 .query("workspace", workspace)
282 .timeout(timeout());
283 if let Some(at) = as_of {
284 req = req.query("as_of", at);
285 }
286 let body: serde_json::Value = req.call().map_err(|e| refused(&url, e))?.into_json()?;
287 let atoms = body
288 .get("atoms")
289 .cloned()
290 .unwrap_or(serde_json::Value::Array(vec![]));
291 Ok(serde_json::from_value(atoms)?)
292 }
293
294 pub fn search(&self, workspace: &str, q: &str, limit: u32) -> Result<Vec<Hit>, Error> {
295 self.search_as_of(workspace, q, limit, None)
296 }
297
298 pub fn search_as_of(
304 &self,
305 workspace: &str,
306 q: &str,
307 limit: u32,
308 as_of: Option<&str>,
309 ) -> Result<Vec<Hit>, Error> {
310 self.search_opts(workspace, q, limit, as_of, false)
311 }
312
313 pub fn search_opts(
323 &self,
324 workspace: &str,
325 q: &str,
326 limit: u32,
327 as_of: Option<&str>,
328 rerank: bool,
329 ) -> Result<Vec<Hit>, Error> {
330 let url = format!("{}/v1/search", self.base);
331 let budget = if rerank {
332 Duration::from_secs(60).max(timeout())
333 } else {
334 timeout()
335 };
336 let mut req = ureq::get(&url)
337 .query("workspace", workspace)
338 .query("q", q)
339 .query("limit", &limit.to_string())
340 .timeout(budget);
341 if let Some(at) = as_of {
342 req = req.query("as_of", at);
343 }
344 if rerank {
345 req = req.query("rerank", "1");
346 }
347 let body: serde_json::Value = req.call().map_err(|e| refused(&url, e))?.into_json()?;
348 let hits = body
349 .get("hits")
350 .cloned()
351 .unwrap_or(serde_json::Value::Array(vec![]));
352 Ok(serde_json::from_value(hits)?)
353 }
354
355 pub fn status(&self, workspace: Option<&str>) -> Result<serde_json::Value, Error> {
361 let url = format!("{}/v1/status", self.base);
362 let mut req = ureq::get(&url).timeout(timeout());
363 if let Some(workspace) = workspace {
364 req = req.query("workspace", workspace);
365 }
366 Ok(req.call().map_err(|e| refused(&url, e))?.into_json()?)
367 }
368
369 pub fn pin(&self, workspace: &str) -> Result<serde_json::Value, Error> {
375 let url = format!("{}/v1/pin", self.base);
376 Ok(ureq::get(&url)
377 .query("workspace", workspace)
378 .timeout(timeout())
379 .call()
380 .map_err(|e| refused(&url, e))?
381 .into_json()?)
382 }
383
384 pub fn set_pin(&self, workspace: &str, name: &str) -> Result<serde_json::Value, Error> {
390 let url = format!("{}/v1/pin", self.base);
391 Ok(ureq::put(&url)
392 .timeout(timeout())
393 .send_json(serde_json::json!({ "workspace": workspace, "name": name }))
394 .map_err(|e| refused(&url, e))?
395 .into_json()?)
396 }
397
398 pub fn accessions(&self, workspace: &str) -> Result<Vec<String>, Error> {
407 let url = format!("{}/v1/accessions", self.base);
408 let body: serde_json::Value = ureq::get(&url)
409 .query("workspace", workspace)
410 .timeout(timeout())
411 .call()
412 .map_err(|e| refused(&url, e))?
413 .into_json()?;
414 let found = body
415 .get("accessions")
416 .cloned()
417 .unwrap_or(serde_json::Value::Array(vec![]));
418 Ok(serde_json::from_value(found)?)
419 }
420 pub fn atoms(&self, workspace: &str) -> Result<Vec<serde_json::Value>, Error> {
429 self.atoms_as_of(workspace, None)
430 }
431
432 pub fn atoms_of_kind(
439 &self,
440 workspace: &str,
441 kind: &str,
442 ) -> Result<Vec<serde_json::Value>, Error> {
443 let url = format!("{}/v1/atoms", self.base);
444 let body: serde_json::Value = ureq::get(&url)
445 .query("workspace", workspace)
446 .query("kind", kind)
447 .timeout(timeout())
448 .call()
449 .map_err(|e| refused(&url, e))?
450 .into_json()?;
451 Ok(body
452 .get("atoms")
453 .and_then(serde_json::Value::as_array)
454 .cloned()
455 .unwrap_or_default())
456 }
457
458 pub fn citers(
464 &self,
465 workspace: &str,
466 accession: &str,
467 ) -> Result<Vec<serde_json::Value>, Error> {
468 let url = format!("{}/v1/citers", self.base);
469 let body: serde_json::Value = ureq::get(&url)
470 .query("workspace", workspace)
471 .query("accession", accession)
472 .timeout(timeout())
473 .call()
474 .map_err(|e| refused(&url, e))?
475 .into_json()?;
476 let found = body
477 .get("atoms")
478 .cloned()
479 .unwrap_or(serde_json::Value::Array(vec![]));
480 Ok(serde_json::from_value(found)?)
481 }
482
483 pub fn delete_atom(
496 &self,
497 workspace: &str,
498 id: &str,
499 why: Option<&str>,
500 ) -> Result<serde_json::Value, Error> {
501 let url = format!("{}/v1/atoms/delete", self.base);
502 let mut body = serde_json::json!({
503 "workspace": workspace,
504 "id": id,
505 });
506 if let Some(accession) = why {
507 body["why"] = serde_json::Value::String(accession.to_string());
508 }
509 let resp = match ureq::post(&url).timeout(timeout()).send_json(body) {
510 Ok(resp) => resp,
511 Err(ureq::Error::Status(404, _)) => {
512 return Err(Error::Bad(format!("no atom {id}")));
513 }
514 Err(e) => return Err(refused(&url, e)),
515 };
516 Ok(resp.into_json()?)
517 }
518
519 pub fn grade(
521 &self,
522 workspace: &str,
523 id: &str,
524 recalled: bool,
525 ) -> Result<serde_json::Value, Error> {
526 let url = format!("{}/v1/grade", self.base);
527 let body: serde_json::Value = ureq::post(&url)
528 .timeout(timeout())
529 .send_json(serde_json::json!({
530 "workspace": workspace,
531 "id": id,
532 "recalled": recalled,
533 }))
534 .map_err(|e| refused(&url, e))?
535 .into_json()?;
536 Ok(body)
537 }
538
539 pub fn hubs(&self, workspace: &str, limit: usize) -> Result<serde_json::Value, Error> {
542 let url = format!("{}/v1/hubs", self.base);
543 let body: serde_json::Value = ureq::get(&url)
544 .query("workspace", workspace)
545 .query("limit", &limit.to_string())
546 .timeout(timeout())
547 .call()
548 .map_err(|e| refused(&url, e))?
549 .into_json()?;
550 Ok(body)
551 }
552
553 pub fn islands(&self, workspace: &str) -> Result<serde_json::Value, Error> {
554 let url = format!("{}/v1/islands", self.base);
555 let body: serde_json::Value = ureq::get(&url)
556 .query("workspace", workspace)
557 .timeout(timeout())
558 .call()
559 .map_err(|e| refused(&url, e))?
560 .into_json()?;
561 Ok(body)
562 }
563
564 pub fn fire(&self, workspace: &str, ids: &[String]) -> Result<serde_json::Value, Error> {
566 self.fire_as(workspace, ids, None)
567 }
568
569 pub fn fire_as(
576 &self,
577 workspace: &str,
578 ids: &[String],
579 lens: Option<&str>,
580 ) -> Result<serde_json::Value, Error> {
581 let url = format!("{}/v1/fire", self.base);
582 let body: serde_json::Value = ureq::post(&url)
583 .timeout(timeout())
584 .send_json(
585 serde_json::json!({"workspace": workspace, "ids": ids, "as": lens.unwrap_or("")}),
586 )
587 .map_err(|e| refused(&url, e))?
588 .into_json()?;
589 Ok(body)
590 }
591
592 pub fn consolidate(&self, workspace: &str, apply: bool) -> Result<serde_json::Value, Error> {
596 let url = format!("{}/v1/consolidate", self.base);
597 let body: serde_json::Value = ureq::post(&url)
598 .timeout(timeout())
599 .send_json(serde_json::json!({"workspace": workspace, "apply": apply}))
600 .map_err(|e| refused(&url, e))?
601 .into_json()?;
602 Ok(body)
603 }
604
605 pub fn sweep(&self, workspace: &str) -> Result<serde_json::Value, Error> {
612 let url = format!("{}/v1/sweep", self.base);
613 let body: serde_json::Value = ureq::post(&url)
614 .timeout(timeout())
615 .send_json(serde_json::json!({"workspace": workspace}))
616 .map_err(|e| refused(&url, e))?
617 .into_json()?;
618 Ok(body)
619 }
620
621 pub fn activate(
624 &self,
625 workspace: &str,
626 q: &str,
627 limit: u32,
628 fire: bool,
629 ) -> Result<serde_json::Value, Error> {
630 self.activate_as(workspace, q, limit, fire, None)
631 }
632
633 pub fn activate_as(
640 &self,
641 workspace: &str,
642 q: &str,
643 limit: u32,
644 fire: bool,
645 lens: Option<&str>,
646 ) -> Result<serde_json::Value, Error> {
647 let url = format!("{}/v1/activate", self.base);
648 let mut req = ureq::get(&url)
649 .query("workspace", workspace)
650 .query("q", q)
651 .query("limit", &limit.to_string())
652 .query("fire", if fire { "1" } else { "0" })
653 .timeout(timeout());
654 if let Some(name) = lens.filter(|n| !n.trim().is_empty()) {
655 req = req.query("as", name);
656 }
657 let body: serde_json::Value = req.call().map_err(|e| refused(&url, e))?.into_json()?;
658 Ok(body)
659 }
660
661 pub fn atoms_in_set(
668 &self,
669 workspace: &str,
670 set: &str,
671 ) -> Result<Vec<serde_json::Value>, Error> {
672 let url = format!("{}/v1/pack", self.base);
673 let body: serde_json::Value = ureq::get(&url)
674 .query("workspace", workspace)
675 .query("set", set)
676 .timeout(timeout())
677 .call()
678 .map_err(|e| refused(&url, e))?
679 .into_json()?;
680 Ok(body
681 .get("atoms")
682 .and_then(serde_json::Value::as_array)
683 .cloned()
684 .unwrap_or_default())
685 }
686
687 pub fn post_atom(&self, atom: &serde_json::Value) -> Result<serde_json::Value, Error> {
688 let url = format!("{}/v1/atoms", self.base);
689 let body: serde_json::Value = ureq::post(&url)
690 .timeout(timeout())
691 .send_json(atom.clone())
692 .map_err(|e| refused(&url, e))?
693 .into_json()?;
694 Ok(body)
695 }
696}
697
698#[cfg(test)]
699mod tests {
700 use super::*;
701
702 #[test]
703 fn resolved_workspace_reads_ljos_env_not_default() {
704 let dir = std::env::temp_dir().join(format!("packset-ljos-env-{}", std::process::id()));
705 std::fs::create_dir_all(dir.join(".config/ljos")).unwrap();
706 std::fs::write(
707 dir.join(".config/ljos/env"),
708 "PACKSET_WORKSPACE=git:example.com/seat/notes\n",
709 )
710 .unwrap();
711 let old_home = env::var("HOME").ok();
712 let old_ws = env::var("PACKSET_WORKSPACE").ok();
713 unsafe {
714 env::remove_var("PACKSET_WORKSPACE");
715 env::set_var("HOME", &dir);
716 }
717 let got = resolved_workspace();
718 unsafe {
719 match old_home {
720 Some(h) => env::set_var("HOME", h),
721 None => env::remove_var("HOME"),
722 }
723 match old_ws {
724 Some(w) => env::set_var("PACKSET_WORKSPACE", w),
725 None => env::remove_var("PACKSET_WORKSPACE"),
726 }
727 }
728 assert_eq!(got, "git:example.com/seat/notes");
729 }
730
731 #[test]
732 fn resolved_workspace_without_env_is_seat_not_default() {
733 let dir = std::env::temp_dir().join(format!("packset-no-ljos-env-{}", std::process::id()));
734 std::fs::create_dir_all(&dir).unwrap();
735 let old_home = env::var("HOME").ok();
736 let old_ws = env::var("PACKSET_WORKSPACE").ok();
737 unsafe {
738 env::remove_var("PACKSET_WORKSPACE");
739 env::set_var("HOME", &dir);
740 }
741 let got = resolved_workspace();
742 unsafe {
743 match old_home {
744 Some(h) => env::set_var("HOME", h),
745 None => env::remove_var("HOME"),
746 }
747 match old_ws {
748 Some(w) => env::set_var("PACKSET_WORKSPACE", w),
749 None => env::remove_var("PACKSET_WORKSPACE"),
750 }
751 }
752 assert_eq!(got, "seat");
753 assert_ne!(got, "default");
754 }
755}