1use std::time::{Duration, Instant};
7
8use serde::Deserialize;
9use serde_json::{Value, json};
10
11use super::{args, caller_name, obj};
12use crate::app::{AppSpec, Apps, Source};
13use crate::backup::{BackupSpec, Backups, Destination, RestoreRequest};
14use crate::error::{Error, Result};
15use crate::jobs::{JobSpec, Jobs, RunStore};
16use crate::org::OrgId;
17use crate::server::{Caller, Registry, Tool};
18
19#[derive(Clone)]
21pub struct Ctx {
22 pub apps: Apps,
23 pub jobs: Jobs,
24 pub backups: Backups,
25}
26
27fn org_of(a: &Value) -> Result<OrgId> {
28 super::arg_org(a)
29}
30
31fn take(a: &mut Value, keys: &[&str]) -> serde_json::Map<String, Value> {
33 let mut out = serde_json::Map::new();
34 if let Some(o) = a.as_object_mut() {
35 for k in keys {
36 if let Some(v) = o.remove(*k) {
37 out.insert((*k).to_string(), v);
38 }
39 }
40 }
41 out
42}
43
44fn trusted(c: &Caller) -> bool {
47 c.is_trusted() || c.principal().is_some_and(|p| p.platform_admin)
48}
49
50pub(super) fn unix_rfc3339(t: Option<i64>) -> Value {
51 match t {
52 Some(t) => json!(crate::cron::rfc3339(t)),
53 None => Value::Null,
54 }
55}
56
57pub(super) fn wait_run(store: &RunStore, id: u64, timeout: Duration) -> Result<crate::jobs::Run> {
59 let started = Instant::now();
60 loop {
61 let r = store.get(id)?;
62 if r.status.finished() || started.elapsed() >= timeout {
63 return Ok(r);
64 }
65 std::thread::sleep(Duration::from_millis(500));
66 }
67}
68
69pub(super) fn timeout_arg(t: &Option<String>, default: Duration) -> Result<Duration> {
70 match t {
71 Some(t) => crate::flex::parse_duration(t)
72 .map_err(Error::invalid)
73 .map(|d| d.min(Duration::from_secs(3600))),
74 None => Ok(default),
75 }
76}
77
78fn volume_backup_admin(x: &Ctx, org: &OrgId, name: &str, c: &Caller) -> Result<()> {
80 if x.backups
81 .get(org, name)
82 .is_ok_and(|b| b.spec.volume.is_some())
83 {
84 super::volumes::require_admin(c, org, "changing or running a volume backup")?;
85 }
86 Ok(())
87}
88
89pub fn database_json(org: &OrgId, a: &crate::app::App, password: Option<&str>) -> Value {
92 let mut v = super::apps::app_json(org, a);
93 if let Source::Database(db) = &a.spec.source {
94 v["connection"] = crate::app::database::connection(&a.spec, db, org, password);
95 }
96 v
97}
98
99fn backup_props() -> Value {
101 json!({
102 "name": {"type": "string"},
103 "database": {"type": "string", "description": "The database app (or give `volume`)."},
104 "volume": {"type": "string", "description": "Or a named volume in the org: its snapshot is exported (incus' tar) and streamed to the bucket; restore with volume_restore."},
105 "destination": {"type": "string"},
106 "schedule": {"type": "string", "description": "Cron: five fields (minute hour day-of-month month day-of-week) or @hourly, @daily, @weekly, @monthly, @yearly."},
107 "timezone": {"type": "string", "description": "UTC (default) or a fixed offset such as +02:00."},
108 "keep": {"type": "integer", "minimum": 1, "maximum": 1000, "description": "Backups kept in the bucket (default 7)."},
109 "compression": {"type": "string", "enum": ["gzip", "zstd", "none"]},
110 "enabled": {"type": "boolean"},
111 "missed_grace": {"type": "string", "description": "How late a slot missed while the daemon was down still runs (default 1h)."}
112 })
113}
114
115fn job_props() -> Value {
117 json!({
118 "name": {"type": "string"},
119 "schedule": {"type": "string", "description": "Cron: five fields or @hourly, @daily, @weekly, @monthly, @yearly."},
120 "timezone": {"type": "string", "description": "UTC (default) or a fixed offset such as +02:00."},
121 "target": {"type": "object", "description": "{app: NAME}, or {stack: NAME, service: NAME}."},
122 "mode": {"type": "string", "enum": ["exec", "run"], "description": "exec (default): in a running replica. run: in a fresh one-off instance from the service's image, env and secrets, deleted after."},
123 "command": {"description": "argv (a list), or a line split like a shell would (no shell runs unless you run one)."},
124 "timeout": {"type": "string", "description": "Kill after this long (default 10m, at most 24h)."},
125 "concurrency": {"type": "string", "enum": ["skip", "allow"], "description": "skip (default): a run due while one is going is skipped."},
126 "keep": {"type": "integer", "minimum": 1, "maximum": 1000, "description": "Runs kept (default 20)."},
127 "enabled": {"type": "boolean"},
128 "user": {"type": "string"},
129 "cwd": {"type": "string"},
130 "env": {"type": "object", "additionalProperties": {"type": "string"}},
131 "missed_grace": {"type": "string", "description": "How late a slot missed while the daemon was down still runs (default 1h)."}
132 })
133}
134
135struct Ann {
137 ro: Value,
138 destructive: Value,
139 write: Value,
140 write_open: Value,
142 destructive_open: Value,
143}
144
145#[derive(Deserialize)]
146#[serde(deny_unknown_fields)]
147struct Named {
148 name: String,
149 #[serde(default)]
150 #[allow(dead_code)]
151 org: Option<String>,
152}
153
154fn argv(a: &mut Value) -> Result<()> {
156 if let Some(Value::String(line)) = a.get("command") {
157 let v =
158 crate::flex::split_words(line).map_err(|e| Error::invalid(format!("command: {e}")))?;
159 a["command"] = json!(v);
160 }
161 Ok(())
162}
163
164fn job_json(j: &Jobs, org: &OrgId, job: &crate::jobs::Job) -> Value {
165 json!({
166 "job": job.spec,
167 "created_at": job.created_at,
168 "updated_at": job.updated_at,
169 "next_run": unix_rfc3339(j.next_run(job)),
170 "last_run": j.runs(org, &job.spec.name).last(),
171 })
172}
173
174pub fn register(r: &mut Registry, ctx: Ctx) -> Result<()> {
175 let ann = Ann {
176 ro: json!({"readOnlyHint": true, "openWorldHint": false}),
177 destructive: json!({"destructiveHint": true, "openWorldHint": false}),
178 write: json!({"destructiveHint": false, "openWorldHint": false}),
179 write_open: json!({"destructiveHint": false, "openWorldHint": true}),
180 destructive_open: json!({"destructiveHint": true, "openWorldHint": true}),
181 };
182
183 database_create_tool(r, &ctx, &ann)?;
186 database_list_tool(r, &ctx, &ann)?;
187 database_get_tool(r, &ctx, &ann)?;
188
189 backup_destination_create_tool(r, &ctx, &ann)?;
192 backup_destination_list_tool(r, &ctx, &ann)?;
193 backup_destination_delete_tool(r, &ctx, &ann)?;
194 backup_destination_test_tool(r, &ctx, &ann)?;
195
196 backup_create_tool(r, &ctx, &ann)?;
199 backup_update_tool(r, &ctx, &ann)?;
200 backup_list_tool(r, &ctx, &ann)?;
201 backup_delete_tool(r, &ctx, &ann)?;
202 backup_run_tool(r, &ctx, &ann)?;
203 backup_runs_tool(r, &ctx, &ann)?;
204 backup_run_log_tool(r, &ctx, &ann)?;
205 backup_restore_tool(r, &ctx, &ann)?;
206
207 job_create_tool(r, &ctx, &ann)?;
210 job_list_tool(r, &ctx, &ann)?;
211 job_get_tool(r, &ctx, &ann)?;
212 job_update_tool(r, &ctx, &ann)?;
213 job_delete_tool(r, &ctx, &ann)?;
214 job_run_tool(r, &ctx, &ann)?;
215 job_runs_tool(r, &ctx, &ann)?;
216 job_run_log_tool(r, &ctx, &ann)?;
217 Ok(())
218}
219
220fn database_create_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
221 tool!(
222 r,
223 ctx,
224 "database_create",
225 "Create a database",
226 "Create a database in a project's environment: Postgres, MySQL, MariaDB, MongoDB or Redis from the official image at `version`, its data on a named volume, one replica rolled out stop-first, with a health check. Credentials are generated and kept as org secrets (db.<name>.password; db.<name>.root-password for MySQL/MariaDB; db.<name>.url, the internal connection URL for apps: DATABASE_URL=${{secret.db.<name>.url}}). Other apps reach it at <name>.<project>-<env>. Not published outside the org unless `publish` is set. A database is an app: deploy, update, roll back and delete it with the app_* tools.",
227 obj(
228 json!({
229 "name": {"type": "string"},
230 "project": {"type": "string"},
231 "environment": {"type": "string", "description": "Default production."},
232 "engine": {"type": "string", "enum": ["postgres", "mysql", "mariadb", "mongodb", "redis"]},
233 "version": {"type": "string", "description": "Image tag (default: 17, 8.4, 11.4, 8.0, 7.4)."},
234 "database": {"type": "string", "description": "Database created on first start (default: the name with - as _). Not Redis."},
235 "user": {"type": "string", "description": "User created on first start (default: as database). Not Redis."},
236 "publish": {"type": "string", "description": "Publish the port on the host: [IP:]PORT (default address 127.0.0.1). Off by default."},
237 "env": {"description": "Extra environment (.env text or a map), e.g. POSTGRES_INITDB_ARGS."},
238 "resources": {"type": "object", "description": "{cpus, memory}."},
239 "deploy": {"type": "boolean", "description": "Deploy right away (default true)."},
240 "wait": {"type": "boolean", "description": "Wait until it is up (default false)."}
241 }),
242 &["name", "project", "engine"]
243 ),
244 ann.write,
245 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
246 let org = org_of(&a)?;
247 let t = take(
248 &mut a,
249 &[
250 "org", "engine", "version", "database", "user", "publish", "deploy", "wait",
251 ],
252 );
253 let engine = crate::app::Engine::parse(
254 t.get("engine").and_then(Value::as_str).unwrap_or_default(),
255 )?;
256 let mut db = json!({"engine": engine});
257 for k in ["version", "database", "user"] {
258 if let Some(v) = t.get(k) {
259 db[k] = v.clone();
260 }
261 }
262 a["source"] = json!({"database": db});
263 if let Some(p) = t.get("publish").and_then(Value::as_str) {
264 let (ip, port) = match p.rsplit_once(':') {
265 Some((ip, port)) => (ip.to_string(), port.to_string()),
266 None => ("127.0.0.1".to_string(), p.to_string()),
267 };
268 port.parse::<u16>()
269 .map_err(|_| Error::invalid(format!("publish {p:?}: [IP:]PORT")))?;
270 a["ports"] = json!([format!("{ip}:{port}:{}", engine.port())]);
271 }
272 let spec: AppSpec = args(a)?;
273 let (app, _) = x.apps.create(&org, spec)?;
274 let mut out = json!({"database": database_json(&org, &app, None)});
275 if t.get("deploy").and_then(Value::as_bool).unwrap_or(true) {
276 let trigger = if c.is_local() {
277 crate::app::deploy::Trigger::Manual
278 } else {
279 crate::app::deploy::Trigger::Api
280 };
281 let d = x
282 .apps
283 .deploy(&org, &app.spec.name, trigger, &caller_name(c), None)?;
284 let d = if t.get("wait").and_then(Value::as_bool).unwrap_or(false) {
285 x.apps
286 .wait(&org, &app.spec.name, d.id, Duration::from_secs(900))?
287 } else {
288 d
289 };
290 out["deployment"] = d.summary();
291 }
292 Ok(out)
293 }
294 );
295 Ok(())
296}
297
298fn database_list_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
299 tool!(
300 r,
301 ctx,
302 "database_list",
303 "List databases",
304 "An org's databases (apps with a database source), each with its engine, version, stack and connection details (password as a secret reference).",
305 obj(json!({"project": {"type": "string"}}), &[]),
306 ann.ro,
307 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
308 let org = org_of(&a)?;
309 let project = a.get("project").and_then(Value::as_str).map(String::from);
310 let dbs: Vec<Value> = x
311 .apps
312 .list(&org)?
313 .into_iter()
314 .filter(|d| matches!(d.spec.source, Source::Database(_)))
315 .filter(|d| project.as_ref().is_none_or(|p| *p == d.spec.project))
316 .map(|d| database_json(&org, &d, None))
317 .collect();
318 Ok(json!({"databases": dbs}))
319 }
320 );
321 Ok(())
322}
323
324fn database_get_tool(r: &mut Registry, ctx: &Ctx, _ann: &Ann) -> Result<()> {
325 tool!(
326 r,
327 ctx,
328 "database_get",
329 "Get a database",
330 "A database's settings and connection details: host (service name), port, user, database, the password as a reference to its org secret, and a URL with that reference. `reveal: true` adds the password and URL values (org members may read org secrets).",
331 obj(
332 json!({"name": {"type": "string"}, "reveal": {"type": "boolean"}}),
333 &["name"]
334 ),
335 json!({"readOnlyHint": true, "openWorldHint": false, "isbSecretReadArg": "reveal"}),
338 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
339 #[derive(Deserialize)]
340 #[serde(deny_unknown_fields)]
341 struct A {
342 name: String,
343 #[serde(default)]
344 reveal: bool,
345 #[serde(default)]
346 #[allow(dead_code)]
347 org: Option<String>,
348 }
349 let org = org_of(&a)?;
350 let a: A = args(a)?;
351 let app = x.apps.get(&org, &a.name)?;
352 if !matches!(app.spec.source, Source::Database(_)) {
353 return Err(Error::invalid(format!("app {} is not a database", a.name)));
354 }
355 let pw = if a.reveal {
356 let (v, _) = x
357 .apps
358 .secrets()
359 .get(&org, &crate::app::database::password_secret(&a.name))?;
360 Some(String::from_utf8_lossy(&v).trim().to_string())
361 } else {
362 None
363 };
364 Ok(database_json(&org, &app, pw.as_deref()))
365 }
366 );
367 Ok(())
368}
369
370fn backup_destination_create_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
371 tool!(
372 r,
373 ctx,
374 "backup_destination_create",
375 "Create a backup destination",
376 "An S3-compatible bucket for backups: endpoint (https://s3.<region>.amazonaws.com, an R2/B2/MinIO URL), region (default us-east-1), bucket, key prefix, path_style (true for MinIO and most self-hosted stores). The key pair is given as access_key/secret_key (stored as the org secrets backup.<name>.access-key/.secret-key) or as the names of existing secrets. Endpoints on this host (loopback) are for the local CLI and platform admins. create_bucket=true creates the bucket; test=true writes, reads back and deletes a small object.",
377 obj(
378 json!({
379 "name": {"type": "string"},
380 "endpoint": {"type": "string"},
381 "region": {"type": "string"},
382 "bucket": {"type": "string"},
383 "prefix": {"type": "string"},
384 "path_style": {"type": "boolean"},
385 "access_key": {"type": "string"},
386 "secret_key": {"type": "string"},
387 "access_key_secret": {"type": "string"},
388 "secret_key_secret": {"type": "string"},
389 "test": {"type": "boolean"},
390 "create_bucket": {"type": "boolean", "description": "Create the bucket first (self-hosted stores)."}
391 }),
392 &["name", "endpoint", "bucket"]
393 ),
394 ann.write_open,
395 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
396 let org = org_of(&a)?;
397 let t = take(
398 &mut a,
399 &["org", "access_key", "secret_key", "test", "create_bucket"],
400 );
401 let s = |k: &str| t.get(k).and_then(Value::as_str).map(String::from);
402 if let Some(o) = a.as_object_mut() {
403 o.entry("access_key_secret").or_insert(json!(""));
404 o.entry("secret_key_secret").or_insert(json!(""));
405 }
406 let d: Destination = args(a)?;
407 let d = x.backups.destination_create(
408 &org,
409 d,
410 s("access_key"),
411 s("secret_key"),
412 trusted(c),
413 )?;
414 let mut out = json!({"destination": d});
415 if t.get("create_bucket")
416 .and_then(Value::as_bool)
417 .unwrap_or(false)
418 {
419 x.backups.destination_create_bucket(&org, &d.name)?;
420 out["bucket_created"] = json!(true);
421 }
422 if t.get("test").and_then(Value::as_bool).unwrap_or(false) {
423 out["test"] = match x.backups.destination_test(&org, &d.name) {
424 Ok(v) => v,
425 Err(e) => json!({"ok": false, "error": e.to_string()}),
426 };
427 }
428 Ok(out)
429 }
430 );
431 Ok(())
432}
433
434fn backup_destination_list_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
435 tool!(
436 r,
437 ctx,
438 "backup_destination_list",
439 "List backup destinations",
440 "An org's backup destinations (key pairs as secret names).",
441 obj(json!({}), &[]),
442 ann.ro,
443 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
444 let org = org_of(&a)?;
445 Ok(json!({"destinations": x.backups.destination_list(&org)?}))
446 }
447 );
448 Ok(())
449}
450
451fn backup_destination_delete_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
452 tool!(
453 r,
454 ctx,
455 "backup_destination_delete",
456 "Delete a backup destination",
457 "Delete a destination no backup uses, with the key secrets isb stored for it. Objects in the bucket are kept.",
458 obj(json!({"name": {"type": "string"}}), &["name"]),
459 ann.destructive,
460 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
461 let org = org_of(&a)?;
462 let a: Named = args(a)?;
463 x.backups.destination_delete(&org, &a.name)?;
464 Ok(json!({"ok": true}))
465 }
466 );
467 Ok(())
468}
469
470fn backup_destination_test_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
471 tool!(
472 r,
473 ctx,
474 "backup_destination_test",
475 "Test a backup destination",
476 "Write a small object under the destination's prefix, check it with HEAD and delete it.",
477 obj(json!({"name": {"type": "string"}}), &["name"]),
478 ann.write_open,
479 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
480 let org = org_of(&a)?;
481 let a: Named = args(a)?;
482 x.backups.destination_test(&org, &a.name)
483 }
484 );
485 Ok(())
486}
487
488fn backup_create_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
489 tool!(
490 r,
491 ctx,
492 "backup_create",
493 "Schedule a backup",
494 "Back a database (or a named `volume`) up on a cron schedule to a destination. A database: the engine's own dump (pg_dump, mysqldump, mariadb-dump, mongodump, a Redis RDB) runs in the database's instance, is compressed and streamed to the bucket by the daemon, checked with HEAD, and the oldest beyond `keep` are deleted. Emits backup.succeeded / backup.failed events.",
495 obj(backup_props(), &["name", "destination", "schedule"]),
496 ann.write,
497 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
498 let org = org_of(&a)?;
499 take(&mut a, &["org"]);
500 let spec: BackupSpec = args(a)?;
501 if spec.volume.is_some() {
502 super::volumes::require_admin(c, &org, "backing up a volume")?;
503 }
504 let b = x.backups.create(&org, spec)?;
505 let next = x.backups.next_run(&b);
506 Ok(json!({"backup": b, "next_run": unix_rfc3339(next)}))
507 }
508 );
509 Ok(())
510}
511
512fn backup_update_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
513 tool!(
514 r,
515 ctx,
516 "backup_update",
517 "Update a backup",
518 "Change a backup's settings (a merge patch: schedule, timezone, destination, keep, compression, enabled, missed_grace). A changed schedule counts from now.",
519 obj(backup_props(), &["name"]),
520 ann.write,
521 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
522 let org = org_of(&a)?;
523 let name = a
524 .get("name")
525 .and_then(Value::as_str)
526 .unwrap_or_default()
527 .to_string();
528 volume_backup_admin(x, &org, &name, c)?;
529 take(&mut a, &["org", "name"]);
530 let b = x.backups.update(&org, &name, &a)?;
531 let next = x.backups.next_run(&b);
532 Ok(json!({"backup": b, "next_run": unix_rfc3339(next)}))
533 }
534 );
535 Ok(())
536}
537
538fn backup_list_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
539 tool!(
540 r,
541 ctx,
542 "backup_list",
543 "List backups",
544 "An org's backup schedules with their last run and next run. With `name`, that backup only, plus the backup files in its bucket (newest first): what backup_restore takes.",
545 obj(
546 json!({"name": {"type": "string"}, "database": {"type": "string"}}),
547 &[]
548 ),
549 ann.ro,
550 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
551 let org = org_of(&a)?;
552 let name = a.get("name").and_then(Value::as_str);
553 let database = a.get("database").and_then(Value::as_str);
554 let mut out = Vec::new();
555 for b in x.backups.list(&org)? {
556 if name.is_some_and(|n| n != b.spec.name)
557 || database.is_some_and(|d| d != b.spec.database)
558 {
559 continue;
560 }
561 let last = x.backups.runs(&org, &b.spec.name).last();
562 let next = x.backups.next_run(&b);
563 let mut v =
564 json!({"backup": b.spec, "last_run": last, "next_run": unix_rfc3339(next)});
565 if name.is_some() {
566 v["files"] = match x.backups.files(&org, &b.spec.name) {
567 Ok(f) => json!(f),
568 Err(e) => json!({"error": e.to_string()}),
569 };
570 }
571 out.push(v);
572 }
573 if let (Some(n), true) = (name, out.is_empty()) {
574 return Err(Error::NotFound(format!("backup {n} in org {org}")));
575 }
576 Ok(json!({"backups": out}))
577 }
578 );
579 Ok(())
580}
581
582fn backup_delete_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
583 tool!(
584 r,
585 ctx,
586 "backup_delete",
587 "Delete a backup",
588 "Delete a backup schedule and its run records. Its files stay in the bucket (restore them with destination and key).",
589 obj(json!({"name": {"type": "string"}}), &["name"]),
590 ann.destructive,
591 |x: &Ctx, a: Value, c: &Caller| -> Result<Value> {
592 let org = org_of(&a)?;
593 let a: Named = args(a)?;
594 volume_backup_admin(x, &org, &a.name, c)?;
595 x.backups.delete(&org, &a.name)?;
596 Ok(json!({"ok": true}))
597 }
598 );
599 Ok(())
600}
601
602fn backup_run_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
603 tool!(
604 r,
605 ctx,
606 "backup_run",
607 "Run a backup now",
608 "Back up now, outside the schedule. Returns the run; wait=true returns when it finishes (at most `timeout`, default 10m).",
609 obj(
610 json!({"name": {"type": "string"}, "wait": {"type": "boolean"}, "timeout": {"type": "string"}}),
611 &["name"]
612 ),
613 ann.write_open,
614 |x: &Ctx, a: Value, c: &Caller| -> Result<Value> {
615 #[derive(Deserialize)]
616 #[serde(deny_unknown_fields)]
617 struct A {
618 name: String,
619 #[serde(default)]
620 wait: bool,
621 #[serde(default)]
622 timeout: Option<String>,
623 #[serde(default)]
624 #[allow(dead_code)]
625 org: Option<String>,
626 }
627 let org = org_of(&a)?;
628 let a: A = args(a)?;
629 volume_backup_admin(x, &org, &a.name, c)?;
630 let r = x.backups.run_now(&org, &a.name, &caller_name(c))?;
631 let r = if a.wait {
632 wait_run(
633 &x.backups.runs(&org, &a.name),
634 r.id,
635 timeout_arg(&a.timeout, Duration::from_secs(600))?,
636 )?
637 } else {
638 r
639 };
640 Ok(json!({"run": r}))
641 }
642 );
643 Ok(())
644}
645
646fn backup_runs_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
647 tool!(
648 r,
649 ctx,
650 "backup_runs",
651 "Backup runs",
652 "A backup's runs, newest first (status, trigger, duration, object key and size). With restores=true instead, the org's restore runs.",
653 obj(
654 json!({"name": {"type": "string"}, "restores": {"type": "boolean"}, "limit": {"type": "integer", "minimum": 1, "maximum": 1000}}),
655 &[]
656 ),
657 ann.ro,
658 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
659 let org = org_of(&a)?;
660 let limit = a.get("limit").and_then(Value::as_u64).unwrap_or(50) as usize;
661 let store = if a.get("restores").and_then(Value::as_bool).unwrap_or(false) {
662 x.backups.restore_runs(&org)
663 } else {
664 let name = a
665 .get("name")
666 .and_then(Value::as_str)
667 .ok_or_else(|| Error::invalid("name (or restores: true) is required"))?;
668 x.backups.get(&org, name)?;
669 x.backups.runs(&org, name)
670 };
671 Ok(json!({"runs": store.list(limit)}))
672 }
673 );
674 Ok(())
675}
676
677fn backup_run_log_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
678 tool!(
679 r,
680 ctx,
681 "backup_run_log",
682 "Backup run log",
683 "One backup (or restore) run's log from byte `offset`; poll with the returned offset until finished.",
684 obj(
685 json!({
686 "name": {"type": "string", "description": "The backup (omit with restore=true)."},
687 "restore": {"type": "boolean"},
688 "run": {"type": "integer", "minimum": 1},
689 "offset": {"type": "integer", "minimum": 0}
690 }),
691 &["run"]
692 ),
693 ann.ro,
694 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
695 let org = org_of(&a)?;
696 let run = a.get("run").and_then(Value::as_u64).unwrap_or(0);
697 let offset = a.get("offset").and_then(Value::as_u64).unwrap_or(0);
698 let store = if a.get("restore").and_then(Value::as_bool).unwrap_or(false) {
699 x.backups.restore_runs(&org)
700 } else {
701 let name = a
702 .get("name")
703 .and_then(Value::as_str)
704 .ok_or_else(|| Error::invalid("name is required"))?;
705 x.backups.get(&org, name)?;
706 x.backups.runs(&org, name)
707 };
708 let (text, offset, finished) = store.log(run, offset)?;
709 Ok(
710 json!({"text": text, "offset": offset, "finished": finished, "run": store.get(run)?}),
711 )
712 }
713 );
714 Ok(())
715}
716
717fn backup_restore_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
718 tool!(
719 r,
720 ctx,
721 "backup_restore",
722 "Restore a backup",
723 "Restore a backup file into a database: `backup` (its newest file, or `key`) or `destination` + `key`; into `target`, an existing database of the same engine whose data is REPLACED (needs confirm: true), or `new`: {name, project?, environment?, version?}, a database created for it (by default beside the backed-up one). The file streams from the bucket through the daemon into the engine's restore tool (pg_restore --clean, mysql, mongorestore --drop, a Redis RDB swap). Emits restore.succeeded / restore.failed.",
724 obj(
725 json!({
726 "backup": {"type": "string"},
727 "destination": {"type": "string"},
728 "key": {"type": "string"},
729 "target": {"type": "string"},
730 "new": {"type": "object", "properties": {"name": {"type": "string"}, "project": {"type": "string"}, "environment": {"type": "string"}, "version": {"type": "string"}}, "required": ["name"]},
731 "confirm": {"type": "boolean"},
732 "wait": {"type": "boolean"},
733 "timeout": {"type": "string"}
734 }),
735 &[]
736 ),
737 ann.destructive_open,
738 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
739 let org = org_of(&a)?;
740 let t = take(&mut a, &["org", "wait", "timeout"]);
741 let req: RestoreRequest = args(a)?;
742 let r = x.backups.restore(&org, req, &caller_name(c))?;
743 let r = if t.get("wait").and_then(Value::as_bool).unwrap_or(false) {
744 let timeout = t.get("timeout").and_then(Value::as_str).map(String::from);
745 wait_run(
746 &x.backups.restore_runs(&org),
747 r.id,
748 timeout_arg(&timeout, Duration::from_secs(1800))?,
749 )?
750 } else {
751 r
752 };
753 Ok(json!({"run": r}))
754 }
755 );
756 Ok(())
757}
758
759fn job_create_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
760 tool!(
761 r,
762 ctx,
763 "job_create",
764 "Create a scheduled job",
765 "Run a command on a cron schedule against an app or a stack service: in a running replica (mode exec) or a fresh one-off instance from its image (mode run). Each run keeps its exit code, duration and output (bounded); job.succeeded / job.failed events.",
766 obj(job_props(), &["name", "schedule", "target", "command"]),
767 ann.write,
768 |x: &Ctx, mut a: Value, _c: &Caller| -> Result<Value> {
769 let org = org_of(&a)?;
770 take(&mut a, &["org"]);
771 argv(&mut a)?;
772 let spec: JobSpec = args(a)?;
773 let j = x.jobs.create(&org, spec)?;
774 Ok(job_json(&x.jobs, &org, &j))
775 }
776 );
777 Ok(())
778}
779
780fn job_list_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
781 tool!(
782 r,
783 ctx,
784 "job_list",
785 "List jobs",
786 "An org's jobs with their next and last run.",
787 obj(json!({}), &[]),
788 ann.ro,
789 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
790 let org = org_of(&a)?;
791 let jobs: Vec<Value> = x
792 .jobs
793 .list(&org)?
794 .iter()
795 .map(|j| job_json(&x.jobs, &org, j))
796 .collect();
797 Ok(json!({"jobs": jobs}))
798 }
799 );
800 Ok(())
801}
802
803fn job_get_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
804 tool!(
805 r,
806 ctx,
807 "job_get",
808 "Get a job",
809 "A job's settings, next run and last run.",
810 obj(json!({"name": {"type": "string"}}), &["name"]),
811 ann.ro,
812 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
813 let org = org_of(&a)?;
814 let a: Named = args(a)?;
815 Ok(job_json(&x.jobs, &org, &x.jobs.get(&org, &a.name)?))
816 }
817 );
818 Ok(())
819}
820
821fn job_update_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
822 tool!(
823 r,
824 ctx,
825 "job_update",
826 "Update a job",
827 "Change a job's settings (a merge patch; the name is fixed). A changed schedule counts from now.",
828 obj(job_props(), &["name"]),
829 ann.write,
830 |x: &Ctx, mut a: Value, _c: &Caller| -> Result<Value> {
831 let org = org_of(&a)?;
832 let name = a
833 .get("name")
834 .and_then(Value::as_str)
835 .unwrap_or_default()
836 .to_string();
837 take(&mut a, &["org", "name"]);
838 argv(&mut a)?;
839 let j = x.jobs.update(&org, &name, &a)?;
840 Ok(job_json(&x.jobs, &org, &j))
841 }
842 );
843 Ok(())
844}
845
846fn job_delete_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
847 tool!(
848 r,
849 ctx,
850 "job_delete",
851 "Delete a job",
852 "Delete a job and its run records.",
853 obj(json!({"name": {"type": "string"}}), &["name"]),
854 ann.destructive,
855 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
856 let org = org_of(&a)?;
857 let a: Named = args(a)?;
858 x.jobs.delete(&org, &a.name)?;
859 Ok(json!({"ok": true}))
860 }
861 );
862 Ok(())
863}
864
865fn job_run_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
866 tool!(
867 r,
868 ctx,
869 "job_run",
870 "Run a job now",
871 "Run a job now, outside its schedule (refused while a run is going under concurrency skip). wait=true returns when it finishes (at most `timeout`, default 10m).",
872 obj(
873 json!({"name": {"type": "string"}, "wait": {"type": "boolean"}, "timeout": {"type": "string"}}),
874 &["name"]
875 ),
876 ann.write,
877 |x: &Ctx, a: Value, c: &Caller| -> Result<Value> {
878 #[derive(Deserialize)]
879 #[serde(deny_unknown_fields)]
880 struct A {
881 name: String,
882 #[serde(default)]
883 wait: bool,
884 #[serde(default)]
885 timeout: Option<String>,
886 #[serde(default)]
887 #[allow(dead_code)]
888 org: Option<String>,
889 }
890 let org = org_of(&a)?;
891 let a: A = args(a)?;
892 let r = x.jobs.run_now(&org, &a.name, &caller_name(c))?;
893 let r = if a.wait {
894 wait_run(
895 &x.jobs.runs(&org, &a.name),
896 r.id,
897 timeout_arg(&a.timeout, Duration::from_secs(600))?,
898 )?
899 } else {
900 r
901 };
902 Ok(json!({"run": r}))
903 }
904 );
905 Ok(())
906}
907
908fn job_runs_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
909 tool!(
910 r,
911 ctx,
912 "job_runs",
913 "Job runs",
914 "A job's runs, newest first: trigger (schedule, missed, manual), status (running, succeeded, failed, skipped), exit code, duration, output size.",
915 obj(
916 json!({"name": {"type": "string"}, "limit": {"type": "integer", "minimum": 1, "maximum": 1000}}),
917 &["name"]
918 ),
919 ann.ro,
920 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
921 let org = org_of(&a)?;
922 let name = a
923 .get("name")
924 .and_then(Value::as_str)
925 .unwrap_or_default()
926 .to_string();
927 x.jobs.get(&org, &name)?;
928 let limit = a.get("limit").and_then(Value::as_u64).unwrap_or(50) as usize;
929 Ok(json!({"runs": x.jobs.runs(&org, &name).list(limit)}))
930 }
931 );
932 Ok(())
933}
934
935fn job_run_log_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
936 tool!(
937 r,
938 ctx,
939 "job_run_log",
940 "Job run log",
941 "One run's output from byte `offset` (the first 192 KiB and the last 64 KiB are kept); poll with the returned offset until finished.",
942 obj(
943 json!({"name": {"type": "string"}, "run": {"type": "integer", "minimum": 1}, "offset": {"type": "integer", "minimum": 0}}),
944 &["name", "run"]
945 ),
946 ann.ro,
947 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
948 let org = org_of(&a)?;
949 let name = a
950 .get("name")
951 .and_then(Value::as_str)
952 .unwrap_or_default()
953 .to_string();
954 x.jobs.get(&org, &name)?;
955 let run = a.get("run").and_then(Value::as_u64).unwrap_or(0);
956 let offset = a.get("offset").and_then(Value::as_u64).unwrap_or(0);
957 let store = x.jobs.runs(&org, &name);
958 let (text, offset, finished) = store.log(run, offset)?;
959 Ok(
960 json!({"text": text, "offset": offset, "finished": finished, "run": store.get(run)?}),
961 )
962 }
963 );
964 Ok(())
965}