1use std::collections::{BTreeMap, VecDeque};
6use std::path::PathBuf;
7use std::sync::{Arc, Condvar, Mutex};
8use std::time::{Duration, Instant};
9
10use serde::Deserialize;
11use serde_json::{Value, json};
12
13use crate::build::{BuildRequest, Builder, BuiltImage};
14use crate::client::Client;
15use crate::error::{Error, Result};
16use crate::org::OrgId;
17use crate::server::{Caller, Registry, Tool};
18use crate::stack::Controller;
19
20use super::policy::RemotePolicy;
21
22const MAX_LINES: usize = 20_000;
24const KEEP_FINISHED: usize = 50;
26
27struct Job {
28 org: OrgId,
29 app: String,
30 tag: String,
31 started: u64,
32 state: Mutex<JobState>,
33 changed: Condvar,
34}
35
36#[derive(Default)]
37struct JobState {
38 first: usize,
40 lines: VecDeque<String>,
41 result: Option<std::result::Result<BuiltImage, String>>,
42 finished: Option<u64>,
43}
44
45impl Job {
46 fn push(&self, line: &str) {
47 let mut s = self.state.lock().unwrap();
48 s.lines.push_back(line.to_string());
49 if s.lines.len() > MAX_LINES {
50 s.lines.pop_front();
51 s.first += 1;
52 }
53 self.changed.notify_all();
54 }
55
56 fn summary(&self, id: &str) -> Value {
57 let s = self.state.lock().unwrap();
58 let mut v = json!({
59 "id": id, "org": self.org, "app": self.app, "tag": self.tag,
60 "started_at": self.started, "finished_at": s.finished,
61 "state": match &s.result { None => "running", Some(Ok(_)) => "succeeded", Some(Err(_)) => "failed" },
62 });
63 match &s.result {
64 Some(Ok(b)) => {
65 v["image"] = json!(b.image);
66 v["digest"] = json!(b.digest);
67 }
68 Some(Err(e)) => v["error"] = json!(e),
69 None => {}
70 }
71 v
72 }
73}
74
75#[derive(Default)]
76struct Jobs {
77 by_id: Mutex<BTreeMap<String, Arc<Job>>>,
78}
79
80impl Jobs {
81 fn get(&self, id: &str, org: &OrgId) -> Result<Arc<Job>> {
82 self.by_id
83 .lock()
84 .unwrap()
85 .get(id)
86 .filter(|j| j.org == *org)
87 .cloned()
88 .ok_or_else(|| Error::NotFound(format!("build {id} in org {org}")))
89 }
90
91 fn add(&self, id: String, j: Arc<Job>) {
92 let mut m = self.by_id.lock().unwrap();
93 m.insert(id, j);
94 let finished: Vec<String> = m
96 .iter()
97 .filter(|(_, j)| j.state.lock().unwrap().result.is_some())
98 .map(|(k, _)| k.clone())
99 .collect();
100 if finished.len() > KEEP_FINISHED {
101 for k in finished.iter().take(finished.len() - KEEP_FINISHED) {
102 m.remove(k);
103 }
104 }
105 }
106}
107
108pub struct Ctx {
110 pub client: Client,
111 pub policy: RemotePolicy,
112 pub ctl: Controller,
113}
114
115fn args<T: serde::de::DeserializeOwned>(v: Value) -> Result<T> {
116 serde_json::from_value(v).map_err(|e| Error::invalid(format!("bad arguments: {e}")))
117}
118
119fn org_of(o: &Option<String>) -> Result<OrgId> {
120 match o {
121 Some(o) => OrgId::new(o.clone()),
122 None => Ok(OrgId::default_org()),
123 }
124}
125
126fn obj(mut props: Value, required: &[&str]) -> Value {
127 props["org"] =
128 json!({"type": "string", "description": "The org to act in (default: default)."});
129 json!({"type": "object", "properties": props, "required": required, "additionalProperties": false})
130}
131
132fn new_build_id() -> String {
135 format!(
136 "b{}-{}{}",
137 crate::stack::now_secs(),
138 crate::stack::new_id(),
139 crate::stack::new_id()
140 )
141}
142
143#[expect(
145 clippy::too_many_lines,
146 reason = "predates the lint ratchet; split it when next changed"
147)]
148pub fn register(r: &mut Registry, ctx: Ctx) -> Result<()> {
149 let ctx = Arc::new(ctx);
150 let jobs = Arc::new(Jobs::default());
151 let ro = json!({"readOnlyHint": true, "openWorldHint": false});
152 let write = json!({"destructiveHint": false, "openWorldHint": true});
153 let destructive = json!({"destructiveHint": true, "openWorldHint": false});
154
155 let (c2, j2) = (ctx.clone(), jobs.clone());
156 r.register(
157 Tool::new(
158 "build_run",
159 "Start a build of a source directory into an image in the org's local registry. It runs in a fresh sandbox in the org (a VM with untrusted=true), never on the host, and is pushed as <org>/<app>:<tag>. Returns an id at once: follow it with build_logs. The result's `image` (registry:APP:TAG@DIGEST) goes in a compose `image:`.",
160 obj(
161 json!({
162 "app": {"type": "string", "description": "The app: names the repository (<org>/<app>) and the build cache. [a-z0-9][a-z0-9._-]*, up to 40 characters."},
163 "context": {"type": "string", "description": "Absolute path of the source directory on the host (read, never written). Remote callers: under a --bind-root."},
164 "subdir": {"type": "string", "description": "Build from this subdirectory of the context."},
165 "builder": {"type": "string", "enum": ["railpack", "nixpacks", "dockerfile"], "description": "How to build (default railpack, or dockerfile when `dockerfile` is set)."},
166 "dockerfile": {"type": "string", "description": "Dockerfile path relative to the context (default Dockerfile)."},
167 "target": {"type": "string", "description": "Dockerfile stage to build."},
168 "args": {"type": "object", "additionalProperties": {"type": "string"}, "description": "Build arguments (Dockerfile ARGs; environment for railpack and nixpacks)."},
169 "tag": {"type": "string", "description": "The tag to push (default latest)."},
170 "untrusted": {"type": "boolean", "description": "Build in a VM (its own kernel)."},
171 "timeout": {"type": "string", "description": "Longest the build may take, e.g. 30m (the default)."}
172 }),
173 &["app", "context"],
174 ),
175 move |a, c| build_run(&c2, &j2, a, c),
176 )
177 .title("Start a build")
178 .annotations(write.clone()),
179 )?;
180
181 let j2 = jobs.clone();
182 r.register(
183 Tool::new(
184 "build_logs",
185 "A build's state and its log lines from `since` (a line number; 0 for all). `wait` (seconds, at most 30) holds the call until there are new lines or the build ends. Done when `state` is succeeded (then `image` and `digest`) or failed (`error`).",
186 obj(
187 json!({
188 "id": {"type": "string"},
189 "since": {"type": "integer", "minimum": 0},
190 "wait": {"type": "integer", "minimum": 0, "maximum": 30}
191 }),
192 &["id"],
193 ),
194 move |a, _c| {
195 #[derive(Deserialize)]
196 #[serde(deny_unknown_fields)]
197 struct A {
198 #[serde(default)]
199 org: Option<String>,
200 id: String,
201 #[serde(default)]
202 since: usize,
203 #[serde(default)]
204 wait: u64,
205 }
206 let a: A = args(a)?;
207 let org = org_of(&a.org)?;
208 let job = j2.get(&a.id, &org)?;
209 let until = Instant::now() + Duration::from_secs(a.wait.min(30));
210 let mut s = job.state.lock().unwrap();
211 while s.result.is_none() && s.first + s.lines.len() <= a.since {
212 let left = until.saturating_duration_since(Instant::now());
213 if left.is_zero() {
214 break;
215 }
216 s = job.changed.wait_timeout(s, left).unwrap().0;
217 }
218 let from = a.since.max(s.first);
219 let lines: Vec<&String> = s.lines.iter().skip(from - s.first).collect();
220 let next = s.first + s.lines.len();
221 let dropped = a.since < s.first;
222 let lines = json!(lines);
223 drop(s);
224 let mut v = job.summary(&a.id);
225 v["lines"] = lines;
226 v["next"] = json!(next);
227 if dropped {
228 v["truncated"] = json!(true);
229 }
230 Ok(v)
231 },
232 )
233 .title("Follow a build")
234 .annotations(ro.clone()),
235 )?;
236
237 let j2 = jobs.clone();
238 r.register(
239 Tool::new(
240 "build_list",
241 "The org's recent builds (running and finished), newest first.",
242 obj(json!({}), &[]),
243 move |a, _c| {
244 #[derive(Deserialize)]
245 #[serde(deny_unknown_fields)]
246 struct A {
247 #[serde(default)]
248 org: Option<String>,
249 }
250 let a: A = args(a)?;
251 let org = org_of(&a.org)?;
252 let m = j2.by_id.lock().unwrap();
253 let mut out: Vec<Value> = m
254 .iter()
255 .filter(|(_, j)| j.org == org)
256 .map(|(id, j)| j.summary(id))
257 .collect();
258 out.reverse();
259 Ok(json!({"builds": out}))
260 },
261 )
262 .title("List builds")
263 .annotations(ro.clone()),
264 )?;
265
266 let c2 = ctx.clone();
267 r.register(
268 Tool::new(
269 "registry_list",
270 "The org's images in the local registry: each app's tags with their digests and push times, newest first. Run one with image: registry:APP:TAG (or @DIGEST) in a compose file.",
271 obj(json!({}), &[]),
272 move |a, _c| {
273 #[derive(Deserialize)]
274 #[serde(deny_unknown_fields)]
275 struct A {
276 #[serde(default)]
277 org: Option<String>,
278 }
279 let a: A = args(a)?;
280 let org = org_of(&a.org)?;
281 let reg = crate::registry::Registry::shared(&c2.client)?;
282 Ok(json!({
283 "registry": crate::registry::status_json(®),
284 "repositories": reg.list(Some(&org))?,
285 }))
286 },
287 )
288 .title("List images")
289 .annotations(ro.clone()),
290 )?;
291
292 let c2 = ctx.clone();
293 r.register(
294 Tool::new(
295 "registry_gc",
296 "Retention for the whole local registry (platform admins): per repository keep the newest `keep` tags (default 10) and every image a deployed stack runs or would roll back to, delete the rest, and reclaim their storage. dry_run reports only.",
297 obj(
298 json!({
299 "keep": {"type": "integer", "minimum": 0},
300 "dry_run": {"type": "boolean"}
301 }),
302 &[],
303 ),
304 move |a, _c| {
305 #[derive(Deserialize)]
306 #[serde(deny_unknown_fields)]
307 struct A {
308 #[serde(default)]
309 #[allow(dead_code)]
310 org: Option<String>,
311 keep: Option<usize>,
312 #[serde(default)]
313 dry_run: bool,
314 }
315 let a: A = args(a)?;
316 let reg = crate::registry::Registry::shared(&c2.client)?;
317 let protected = crate::registry::protected_by(&c2.ctl.definitions());
318 let mut lines = Vec::new();
319 let rep = reg.gc(
320 a.keep.unwrap_or(crate::registry::DEFAULT_KEEP),
321 &protected,
322 a.dry_run,
323 &mut |l| lines.push(l.to_string()),
324 )?;
325 let mut v = serde_json::to_value(rep)?;
326 v["log"] = json!(lines);
327 Ok(v)
328 },
329 )
330 .title("Registry retention")
331 .annotations(destructive),
332 )?;
333 Ok(())
334}
335
336#[expect(
337 clippy::too_many_lines,
338 reason = "predates the lint ratchet; split it when next changed"
339)]
340fn build_run(ctx: &Ctx, jobs: &Arc<Jobs>, a: Value, c: &Caller) -> Result<Value> {
341 #[derive(Deserialize)]
342 #[serde(deny_unknown_fields)]
343 struct A {
344 #[serde(default)]
345 org: Option<String>,
346 app: String,
347 context: PathBuf,
348 #[serde(default)]
349 subdir: Option<String>,
350 #[serde(default)]
351 builder: Option<String>,
352 #[serde(default)]
353 dockerfile: Option<String>,
354 #[serde(default)]
355 target: Option<String>,
356 #[serde(default)]
357 args: BTreeMap<String, String>,
358 #[serde(default)]
359 tag: Option<String>,
360 #[serde(default)]
361 untrusted: bool,
362 #[serde(default)]
363 timeout: Option<String>,
364 }
365 let a: A = args(a)?;
366 let org = org_of(&a.org)?;
367 if !a.context.is_absolute() {
368 return Err(Error::invalid("context must be an absolute path"));
369 }
370 if !c.is_trusted() {
371 ctx.policy.check_base_dir(&a.context)?;
372 }
373 let builder = match (a.builder.as_deref(), &a.dockerfile) {
374 (None, Some(_)) | (Some("dockerfile"), _) => Builder::Dockerfile {
375 path: a.dockerfile.clone().unwrap_or_else(|| "Dockerfile".into()),
376 target: a.target.clone(),
377 },
378 (None, None) | (Some("railpack"), None) => Builder::Railpack,
379 (Some("nixpacks"), None) => Builder::Nixpacks,
380 (Some("buildpacks"), None) => Builder::Buildpacks { builder: None },
381 (Some(b), _) => {
382 return Err(Error::invalid(format!(
383 "builder {b:?}: railpack, nixpacks or dockerfile (a dockerfile path needs the dockerfile builder)"
384 )));
385 }
386 };
387 let req = BuildRequest {
388 org: org.clone(),
389 app: a.app.clone(),
390 context: a.context.clone(),
391 subdir: a.subdir.clone(),
392 builder,
393 args: a.args.into_iter().collect(),
394 tag: a.tag.clone().unwrap_or_else(|| "latest".into()),
395 untrusted: a.untrusted,
396 cache: None,
397 };
398 let mut opts = crate::build::BuildOptions::default();
399 if let Some(t) = &a.timeout {
400 opts.timeout = crate::flex::parse_duration(t).map_err(Error::invalid)?;
401 }
402 crate::registry::Registry::shared(&ctx.client)?;
404 let id = new_build_id();
405 let job = Arc::new(Job {
406 org: org.clone(),
407 app: req.app.clone(),
408 tag: req.tag.clone(),
409 started: crate::stack::now_secs(),
410 state: Mutex::new(JobState::default()),
411 changed: Condvar::new(),
412 });
413 jobs.add(id.clone(), job.clone());
414 let client = ctx.client.clone();
415 let ctl = ctx.ctl.clone();
416 let who = c.to_string();
417 let stack = crate::stack::qualified(&org, &format!("build:{}", req.app));
418 ctl.note(
419 "info",
420 &stack,
421 format!("build {id} of {}:{} started by {who}", req.app, req.tag),
422 );
423 std::thread::Builder::new()
424 .name(format!("isb-build-{id}"))
425 .spawn(move || {
426 let r = crate::build::run_with(&client, &req, &opts, &mut |l| job.push(l));
427 let (level, msg) = match &r {
428 Ok(b) => (
429 "info",
430 format!("build of {}:{} succeeded: {}", req.app, req.tag, b.image),
431 ),
432 Err(e) => (
433 "error",
434 format!("build of {}:{} failed: {e}", req.app, req.tag),
435 ),
436 };
437 if let Err(e) = &r {
438 job.push(&format!("error: {e}"));
439 }
440 ctl.note(level, &stack, msg);
441 let mut s = job.state.lock().unwrap();
442 s.result = Some(r.map_err(|e| e.to_string()));
443 s.finished = Some(crate::stack::now_secs());
444 job.changed.notify_all();
445 })
446 .map_err(Error::Io)?;
447 Ok(json!({"id": id, "org": org, "app": a.app, "tag": a.tag.unwrap_or_else(|| "latest".into())}))
448}