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", "description": "Default true. false: the job is kept but never runs on its schedule (job_run still runs it) until enabled."},
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 {
167 let mut v = json!(job.spec);
168 v["created_at"] = json!(job.created_at);
169 v["updated_at"] = json!(job.updated_at);
170 v["next_run"] = json!(unix_rfc3339(j.next_run(job)));
171 v["last_run"] = json!(j.runs(org, &job.spec.name).last());
172 v
173}
174
175pub fn register(r: &mut Registry, ctx: Ctx) -> Result<()> {
176 let ann = Ann {
177 ro: json!({"readOnlyHint": true, "openWorldHint": false}),
178 destructive: json!({"destructiveHint": true, "openWorldHint": false}),
179 write: json!({"destructiveHint": false, "openWorldHint": false}),
180 write_open: json!({"destructiveHint": false, "openWorldHint": true}),
181 destructive_open: json!({"destructiveHint": true, "openWorldHint": true}),
182 };
183
184 database_create_tool(r, &ctx, &ann)?;
187 database_list_tool(r, &ctx, &ann)?;
188 database_get_tool(r, &ctx, &ann)?;
189
190 backup_destination_create_tool(r, &ctx, &ann)?;
193 backup_destination_list_tool(r, &ctx, &ann)?;
194 backup_destination_delete_tool(r, &ctx, &ann)?;
195 backup_destination_test_tool(r, &ctx, &ann)?;
196
197 backup_create_tool(r, &ctx, &ann)?;
200 backup_update_tool(r, &ctx, &ann)?;
201 backup_list_tool(r, &ctx, &ann)?;
202 backup_delete_tool(r, &ctx, &ann)?;
203 backup_run_tool(r, &ctx, &ann)?;
204 backup_runs_tool(r, &ctx, &ann)?;
205 backup_run_log_tool(r, &ctx, &ann)?;
206 backup_restore_tool(r, &ctx, &ann)?;
207
208 job_create_tool(r, &ctx, &ann)?;
211 job_list_tool(r, &ctx, &ann)?;
212 job_get_tool(r, &ctx, &ann)?;
213 job_update_tool(r, &ctx, &ann)?;
214 job_delete_tool(r, &ctx, &ann)?;
215 job_run_tool(r, &ctx, &ann)?;
216 job_runs_tool(r, &ctx, &ann)?;
217 job_run_log_tool(r, &ctx, &ann)?;
218 Ok(())
219}
220
221fn database_create_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
222 tool!(
223 r,
224 ctx,
225 "database_create",
226 "Create a database",
227 "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}}; `urls` keeps more such secrets with driver options). Setting db.<name>.password changes the password inside the running database first, then the URL secrets. 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.",
228 obj(
229 json!({
230 "name": {"type": "string"},
231 "project": {"type": "string"},
232 "environment": {"type": "string", "description": "Default production."},
233 "engine": {"type": "string", "enum": ["postgres", "mysql", "mariadb", "mongodb", "redis"]},
234 "version": {"type": "string", "description": "Image tag (default: 17, 8.4, 11.4, 8.0, 7.4)."},
235 "database": {"type": "string", "description": "Database created on first start (default: the name with - as _). Not Redis."},
236 "user": {"type": "string", "description": "User created on first start (default: as database). Not Redis."},
237 "urls": {"type": "object", "additionalProperties": {"type": "string"}, "description": "More secrets isb keeps holding the internal URL, each with a query string for a driver's options ({\"dsn.main-db.web\": \"sslmode=disable\"}, \"\" for none). Written at deploy and again whenever the password changes, so no app holds a stale copy of it."},
238 "publish": {"type": "string", "description": "Publish the port on the host: [IP:]PORT (default address 127.0.0.1). Off by default."},
239 "env": {"description": "Extra environment (.env text or a map), e.g. POSTGRES_INITDB_ARGS."},
240 "resources": crate::app::resources_schema(),
241 "deploy": {"type": "boolean", "description": "Deploy right away (default true)."},
242 "wait": {"type": "boolean", "description": "Wait until it is up (default false)."}
243 }),
244 &["name", "project", "engine"]
245 ),
246 ann.write,
247 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
248 let org = org_of(&a)?;
249 let t = take(
250 &mut a,
251 &[
252 "org", "engine", "version", "database", "user", "urls", "publish", "deploy",
253 "wait",
254 ],
255 );
256 let engine = crate::app::Engine::parse(
257 t.get("engine").and_then(Value::as_str).unwrap_or_default(),
258 )?;
259 let mut db = json!({"engine": engine});
260 for k in ["version", "database", "user", "urls"] {
261 if let Some(v) = t.get(k) {
262 db[k] = v.clone();
263 }
264 }
265 a["source"] = json!({"database": db});
266 if let Some(p) = t.get("publish").and_then(Value::as_str) {
267 let (ip, port) = match p.rsplit_once(':') {
268 Some((ip, port)) => (ip.to_string(), port.to_string()),
269 None => ("127.0.0.1".to_string(), p.to_string()),
270 };
271 port.parse::<u16>()
272 .map_err(|_| Error::invalid(format!("publish {p:?}: [IP:]PORT")))?;
273 a["ports"] = json!([format!("{ip}:{port}:{}", engine.port())]);
274 }
275 let spec: AppSpec = args(a)?;
276 let (app, _) = x.apps.create(&org, spec)?;
277 let mut out = json!({"database": database_json(&org, &app, None)});
278 if t.get("deploy").and_then(Value::as_bool).unwrap_or(true) {
279 let trigger = if c.is_local() {
280 crate::app::deploy::Trigger::Manual
281 } else {
282 crate::app::deploy::Trigger::Api
283 };
284 let d = x
285 .apps
286 .deploy(&org, &app.spec.name, trigger, &caller_name(c), None)?;
287 let d = if t.get("wait").and_then(Value::as_bool).unwrap_or(false) {
288 x.apps
289 .wait(&org, &app.spec.name, d.id, Duration::from_secs(900))?
290 } else {
291 d
292 };
293 out["deployment"] = d.summary();
294 }
295 Ok(out)
296 }
297 );
298 Ok(())
299}
300
301fn database_list_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
302 tool!(
303 r,
304 ctx,
305 "database_list",
306 "List databases",
307 "An org's databases (apps with a database source), each with its engine, version, stack and connection details (password as a secret reference).",
308 obj(json!({"project": {"type": "string"}}), &[]),
309 ann.ro,
310 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
311 let org = org_of(&a)?;
312 let project = a.get("project").and_then(Value::as_str).map(String::from);
313 let dbs: Vec<Value> = x
314 .apps
315 .list(&org)?
316 .into_iter()
317 .filter(|d| matches!(d.spec.source, Source::Database(_)))
318 .filter(|d| project.as_ref().is_none_or(|p| *p == d.spec.project))
319 .map(|d| database_json(&org, &d, None))
320 .collect();
321 Ok(json!({"databases": dbs}))
322 }
323 );
324 Ok(())
325}
326
327fn database_get_tool(r: &mut Registry, ctx: &Ctx, _ann: &Ann) -> Result<()> {
328 tool!(
329 r,
330 ctx,
331 "database_get",
332 "Get a database",
333 "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).",
334 obj(
335 json!({"name": {"type": "string"}, "reveal": {"type": "boolean"}}),
336 &["name"]
337 ),
338 json!({"readOnlyHint": true, "openWorldHint": false, "isbSecretReadArg": "reveal"}),
341 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
342 #[derive(Deserialize)]
343 #[serde(deny_unknown_fields)]
344 struct A {
345 name: String,
346 #[serde(default)]
347 reveal: bool,
348 #[serde(default)]
349 #[allow(dead_code)]
350 org: Option<String>,
351 }
352 let org = org_of(&a)?;
353 let a: A = args(a)?;
354 let app = x.apps.get(&org, &a.name)?;
355 if !matches!(app.spec.source, Source::Database(_)) {
356 return Err(Error::invalid(format!("app {} is not a database", a.name)));
357 }
358 let pw = if a.reveal {
359 let (v, _) = x
360 .apps
361 .secrets()
362 .get(&org, &crate::app::database::password_secret(&a.name))?;
363 Some(String::from_utf8_lossy(&v).trim().to_string())
364 } else {
365 None
366 };
367 Ok(database_json(&org, &app, pw.as_deref()))
368 }
369 );
370 Ok(())
371}
372
373fn backup_destination_create_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
374 tool!(
375 r,
376 ctx,
377 "backup_destination_create",
378 "Create a backup destination",
379 "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.",
380 obj(
381 json!({
382 "name": {"type": "string"},
383 "endpoint": {"type": "string"},
384 "region": {"type": "string"},
385 "bucket": {"type": "string"},
386 "prefix": {"type": "string"},
387 "path_style": {"type": "boolean"},
388 "access_key": {"type": "string"},
389 "secret_key": {"type": "string"},
390 "access_key_secret": {"type": "string"},
391 "secret_key_secret": {"type": "string"},
392 "test": {"type": "boolean"},
393 "create_bucket": {"type": "boolean", "description": "Create the bucket first (self-hosted stores)."}
394 }),
395 &["name", "endpoint", "bucket"]
396 ),
397 ann.write_open,
398 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
399 let org = org_of(&a)?;
400 let t = take(
401 &mut a,
402 &["org", "access_key", "secret_key", "test", "create_bucket"],
403 );
404 let s = |k: &str| t.get(k).and_then(Value::as_str).map(String::from);
405 if let Some(o) = a.as_object_mut() {
406 o.entry("access_key_secret").or_insert(json!(""));
407 o.entry("secret_key_secret").or_insert(json!(""));
408 }
409 let d: Destination = args(a)?;
410 let d = x.backups.destination_create(
411 &org,
412 d,
413 s("access_key"),
414 s("secret_key"),
415 trusted(c),
416 )?;
417 let mut out = json!({"destination": d});
418 if t.get("create_bucket")
419 .and_then(Value::as_bool)
420 .unwrap_or(false)
421 {
422 x.backups.destination_create_bucket(&org, &d.name)?;
423 out["bucket_created"] = json!(true);
424 }
425 if t.get("test").and_then(Value::as_bool).unwrap_or(false) {
426 out["test"] = match x.backups.destination_test(&org, &d.name) {
427 Ok(v) => v,
428 Err(e) => json!({"ok": false, "error": e.to_string()}),
429 };
430 }
431 Ok(out)
432 }
433 );
434 Ok(())
435}
436
437fn backup_destination_list_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
438 tool!(
439 r,
440 ctx,
441 "backup_destination_list",
442 "List backup destinations",
443 "An org's backup destinations (key pairs as secret names).",
444 obj(json!({}), &[]),
445 ann.ro,
446 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
447 let org = org_of(&a)?;
448 Ok(json!({"destinations": x.backups.destination_list(&org)?}))
449 }
450 );
451 Ok(())
452}
453
454fn backup_destination_delete_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
455 tool!(
456 r,
457 ctx,
458 "backup_destination_delete",
459 "Delete a backup destination",
460 "Delete a destination no backup uses, with the key secrets isb stored for it. Objects in the bucket are kept.",
461 obj(json!({"name": {"type": "string"}}), &["name"]),
462 ann.destructive,
463 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
464 let org = org_of(&a)?;
465 let a: Named = args(a)?;
466 x.backups.destination_delete(&org, &a.name)?;
467 Ok(json!({"ok": true}))
468 }
469 );
470 Ok(())
471}
472
473fn backup_destination_test_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
474 tool!(
475 r,
476 ctx,
477 "backup_destination_test",
478 "Test a backup destination",
479 "Write a small object under the destination's prefix, check it with HEAD and delete it.",
480 obj(json!({"name": {"type": "string"}}), &["name"]),
481 ann.write_open,
482 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
483 let org = org_of(&a)?;
484 let a: Named = args(a)?;
485 x.backups.destination_test(&org, &a.name)
486 }
487 );
488 Ok(())
489}
490
491fn backup_create_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
492 tool!(
493 r,
494 ctx,
495 "backup_create",
496 "Schedule a backup",
497 "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.",
498 obj(backup_props(), &["name", "destination", "schedule"]),
499 ann.write,
500 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
501 let org = org_of(&a)?;
502 take(&mut a, &["org"]);
503 let spec: BackupSpec = args(a)?;
504 if spec.volume.is_some() {
505 super::volumes::require_admin(c, &org, "backing up a volume")?;
506 }
507 let b = x.backups.create(&org, spec)?;
508 let next = x.backups.next_run(&b);
509 Ok(json!({"backup": b, "next_run": unix_rfc3339(next)}))
510 }
511 );
512 Ok(())
513}
514
515fn backup_update_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
516 tool!(
517 r,
518 ctx,
519 "backup_update",
520 "Update a backup",
521 "Change a backup's settings (a merge patch: schedule, timezone, destination, keep, compression, enabled, missed_grace). A changed schedule counts from now.",
522 obj(backup_props(), &["name"]),
523 ann.write,
524 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
525 let org = org_of(&a)?;
526 let name = a
527 .get("name")
528 .and_then(Value::as_str)
529 .unwrap_or_default()
530 .to_string();
531 volume_backup_admin(x, &org, &name, c)?;
532 take(&mut a, &["org", "name"]);
533 let b = x.backups.update(&org, &name, &a)?;
534 let next = x.backups.next_run(&b);
535 Ok(json!({"backup": b, "next_run": unix_rfc3339(next)}))
536 }
537 );
538 Ok(())
539}
540
541fn backup_list_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
542 tool!(
543 r,
544 ctx,
545 "backup_list",
546 "List backups",
547 "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.",
548 obj(
549 json!({"name": {"type": "string"}, "database": {"type": "string"}}),
550 &[]
551 ),
552 ann.ro,
553 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
554 let org = org_of(&a)?;
555 let name = a.get("name").and_then(Value::as_str);
556 let database = a.get("database").and_then(Value::as_str);
557 let mut out = Vec::new();
558 for b in x.backups.list(&org)? {
559 if name.is_some_and(|n| n != b.spec.name)
560 || database.is_some_and(|d| d != b.spec.database)
561 {
562 continue;
563 }
564 let last = x.backups.runs(&org, &b.spec.name).last();
565 let next = x.backups.next_run(&b);
566 let mut v =
567 json!({"backup": b.spec, "last_run": last, "next_run": unix_rfc3339(next)});
568 if name.is_some() {
569 v["files"] = match x.backups.files(&org, &b.spec.name) {
570 Ok(f) => json!(f),
571 Err(e) => json!({"error": e.to_string()}),
572 };
573 }
574 out.push(v);
575 }
576 if let (Some(n), true) = (name, out.is_empty()) {
577 return Err(Error::NotFound(format!("backup {n} in org {org}")));
578 }
579 Ok(json!({"backups": out}))
580 }
581 );
582 Ok(())
583}
584
585fn backup_delete_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
586 tool!(
587 r,
588 ctx,
589 "backup_delete",
590 "Delete a backup",
591 "Delete a backup schedule and its run records. Its files stay in the bucket (restore them with destination and key).",
592 obj(json!({"name": {"type": "string"}}), &["name"]),
593 ann.destructive,
594 |x: &Ctx, a: Value, c: &Caller| -> Result<Value> {
595 let org = org_of(&a)?;
596 let a: Named = args(a)?;
597 volume_backup_admin(x, &org, &a.name, c)?;
598 x.backups.delete(&org, &a.name)?;
599 Ok(json!({"ok": true}))
600 }
601 );
602 Ok(())
603}
604
605fn backup_run_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
606 tool!(
607 r,
608 ctx,
609 "backup_run",
610 "Run a backup now",
611 "Back up now, outside the schedule. Returns the run; wait=true returns when it finishes (at most `timeout`, default 10m).",
612 obj(
613 json!({"name": {"type": "string"}, "wait": {"type": "boolean"}, "timeout": {"type": "string"}}),
614 &["name"]
615 ),
616 ann.write_open,
617 |x: &Ctx, a: Value, c: &Caller| -> Result<Value> {
618 #[derive(Deserialize)]
619 #[serde(deny_unknown_fields)]
620 struct A {
621 name: String,
622 #[serde(default)]
623 wait: bool,
624 #[serde(default)]
625 timeout: Option<String>,
626 #[serde(default)]
627 #[allow(dead_code)]
628 org: Option<String>,
629 }
630 let org = org_of(&a)?;
631 let a: A = args(a)?;
632 volume_backup_admin(x, &org, &a.name, c)?;
633 let r = x.backups.run_now(&org, &a.name, &caller_name(c))?;
634 let r = if a.wait {
635 wait_run(
636 &x.backups.runs(&org, &a.name),
637 r.id,
638 timeout_arg(&a.timeout, Duration::from_secs(600))?,
639 )?
640 } else {
641 r
642 };
643 Ok(json!({"run": r}))
644 }
645 );
646 Ok(())
647}
648
649fn backup_runs_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
650 tool!(
651 r,
652 ctx,
653 "backup_runs",
654 "Backup runs",
655 "A backup's runs, newest first (status, trigger, duration, object key and size). With restores=true instead, the org's restore runs.",
656 obj(
657 json!({"name": {"type": "string"}, "restores": {"type": "boolean"}, "limit": {"type": "integer", "minimum": 1, "maximum": 1000}}),
658 &[]
659 ),
660 ann.ro,
661 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
662 let org = org_of(&a)?;
663 let limit = a.get("limit").and_then(Value::as_u64).unwrap_or(50) as usize;
664 let store = if a.get("restores").and_then(Value::as_bool).unwrap_or(false) {
665 x.backups.restore_runs(&org)
666 } else {
667 let name = a
668 .get("name")
669 .and_then(Value::as_str)
670 .ok_or_else(|| Error::invalid("name (or restores: true) is required"))?;
671 x.backups.get(&org, name)?;
672 x.backups.runs(&org, name)
673 };
674 Ok(json!({"runs": store.list(limit)}))
675 }
676 );
677 Ok(())
678}
679
680fn backup_run_log_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
681 tool!(
682 r,
683 ctx,
684 "backup_run_log",
685 "Backup run log",
686 "One backup (or restore) run's log from byte `offset`; poll with the returned offset until finished.",
687 obj(
688 json!({
689 "name": {"type": "string", "description": "The backup (omit with restore=true)."},
690 "restore": {"type": "boolean"},
691 "run": {"type": "integer", "minimum": 1},
692 "offset": {"type": "integer", "minimum": 0}
693 }),
694 &["run"]
695 ),
696 ann.ro,
697 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
698 let org = org_of(&a)?;
699 let run = a.get("run").and_then(Value::as_u64).unwrap_or(0);
700 let offset = a.get("offset").and_then(Value::as_u64).unwrap_or(0);
701 let store = if a.get("restore").and_then(Value::as_bool).unwrap_or(false) {
702 x.backups.restore_runs(&org)
703 } else {
704 let name = a
705 .get("name")
706 .and_then(Value::as_str)
707 .ok_or_else(|| Error::invalid("name is required"))?;
708 x.backups.get(&org, name)?;
709 x.backups.runs(&org, name)
710 };
711 let (text, offset, finished) = store.log(run, offset)?;
712 Ok(
713 json!({"text": text, "offset": offset, "finished": finished, "run": store.get(run)?}),
714 )
715 }
716 );
717 Ok(())
718}
719
720fn backup_restore_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
721 tool!(
722 r,
723 ctx,
724 "backup_restore",
725 "Restore a backup",
726 "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.",
727 obj(
728 json!({
729 "backup": {"type": "string"},
730 "destination": {"type": "string"},
731 "key": {"type": "string"},
732 "target": {"type": "string"},
733 "new": {"type": "object", "properties": {"name": {"type": "string"}, "project": {"type": "string"}, "environment": {"type": "string"}, "version": {"type": "string"}}, "required": ["name"]},
734 "confirm": {"type": "boolean"},
735 "wait": {"type": "boolean"},
736 "timeout": {"type": "string"}
737 }),
738 &[]
739 ),
740 ann.destructive_open,
741 |x: &Ctx, mut a: Value, c: &Caller| -> Result<Value> {
742 let org = org_of(&a)?;
743 let t = take(&mut a, &["org", "wait", "timeout"]);
744 let req: RestoreRequest = args(a)?;
745 let r = x.backups.restore(&org, req, &caller_name(c))?;
746 let r = if t.get("wait").and_then(Value::as_bool).unwrap_or(false) {
747 let timeout = t.get("timeout").and_then(Value::as_str).map(String::from);
748 wait_run(
749 &x.backups.restore_runs(&org),
750 r.id,
751 timeout_arg(&timeout, Duration::from_secs(1800))?,
752 )?
753 } else {
754 r
755 };
756 Ok(json!({"run": r}))
757 }
758 );
759 Ok(())
760}
761
762fn job_create_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
763 tool!(
764 r,
765 ctx,
766 "job_create",
767 "Create a scheduled job",
768 "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. enabled: false creates it disabled. Answers the job as job_get does.",
769 obj(job_props(), &["name", "schedule", "target", "command"]),
770 ann.write,
771 |x: &Ctx, mut a: Value, _c: &Caller| -> Result<Value> {
772 let org = org_of(&a)?;
773 take(&mut a, &["org"]);
774 argv(&mut a)?;
775 let spec: JobSpec = args(a)?;
776 let j = x.jobs.create(&org, spec)?;
777 Ok(job_json(&x.jobs, &org, &j))
778 }
779 );
780 Ok(())
781}
782
783fn job_list_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
784 tool!(
785 r,
786 ctx,
787 "job_list",
788 "List jobs",
789 "An org's jobs, each as job_get answers one.",
790 obj(json!({}), &[]),
791 ann.ro,
792 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
793 let org = org_of(&a)?;
794 let jobs: Vec<Value> = x
795 .jobs
796 .list(&org)?
797 .iter()
798 .map(|j| job_json(&x.jobs, &org, j))
799 .collect();
800 Ok(json!({"jobs": jobs}))
801 }
802 );
803 Ok(())
804}
805
806fn job_get_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
807 tool!(
808 r,
809 ctx,
810 "job_get",
811 "Get a job",
812 "A job's settings (name, schedule, target, command, enabled, ...) with created_at, updated_at, next_run (RFC 3339; null when disabled) and last_run, all at the top level.",
813 obj(json!({"name": {"type": "string"}}), &["name"]),
814 ann.ro,
815 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
816 let org = org_of(&a)?;
817 let a: Named = args(a)?;
818 Ok(job_json(&x.jobs, &org, &x.jobs.get(&org, &a.name)?))
819 }
820 );
821 Ok(())
822}
823
824fn job_update_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
825 tool!(
826 r,
827 ctx,
828 "job_update",
829 "Update a job",
830 "Change a job's settings (a merge patch; the name is fixed). A changed schedule counts from now.",
831 obj(job_props(), &["name"]),
832 ann.write,
833 |x: &Ctx, mut a: Value, _c: &Caller| -> Result<Value> {
834 let org = org_of(&a)?;
835 let name = a
836 .get("name")
837 .and_then(Value::as_str)
838 .unwrap_or_default()
839 .to_string();
840 take(&mut a, &["org", "name"]);
841 argv(&mut a)?;
842 let j = x.jobs.update(&org, &name, &a)?;
843 Ok(job_json(&x.jobs, &org, &j))
844 }
845 );
846 Ok(())
847}
848
849fn job_delete_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
850 tool!(
851 r,
852 ctx,
853 "job_delete",
854 "Delete a job",
855 "Delete a job and its run records.",
856 obj(json!({"name": {"type": "string"}}), &["name"]),
857 ann.destructive,
858 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
859 let org = org_of(&a)?;
860 let a: Named = args(a)?;
861 x.jobs.delete(&org, &a.name)?;
862 Ok(json!({"ok": true}))
863 }
864 );
865 Ok(())
866}
867
868fn job_run_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
869 tool!(
870 r,
871 ctx,
872 "job_run",
873 "Run a job now",
874 "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).",
875 obj(
876 json!({"name": {"type": "string"}, "wait": {"type": "boolean"}, "timeout": {"type": "string"}}),
877 &["name"]
878 ),
879 ann.write,
880 |x: &Ctx, a: Value, c: &Caller| -> Result<Value> {
881 #[derive(Deserialize)]
882 #[serde(deny_unknown_fields)]
883 struct A {
884 name: String,
885 #[serde(default)]
886 wait: bool,
887 #[serde(default)]
888 timeout: Option<String>,
889 #[serde(default)]
890 #[allow(dead_code)]
891 org: Option<String>,
892 }
893 let org = org_of(&a)?;
894 let a: A = args(a)?;
895 let r = x.jobs.run_now(&org, &a.name, &caller_name(c))?;
896 let r = if a.wait {
897 wait_run(
898 &x.jobs.runs(&org, &a.name),
899 r.id,
900 timeout_arg(&a.timeout, Duration::from_secs(600))?,
901 )?
902 } else {
903 r
904 };
905 Ok(json!({"run": r}))
906 }
907 );
908 Ok(())
909}
910
911fn job_runs_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
912 tool!(
913 r,
914 ctx,
915 "job_runs",
916 "Job runs",
917 "A job's runs, newest first: trigger (schedule, missed, manual), status (running, succeeded, failed, skipped), exit code, duration, output size.",
918 obj(
919 json!({"name": {"type": "string"}, "limit": {"type": "integer", "minimum": 1, "maximum": 1000}}),
920 &["name"]
921 ),
922 ann.ro,
923 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
924 let org = org_of(&a)?;
925 let name = a
926 .get("name")
927 .and_then(Value::as_str)
928 .unwrap_or_default()
929 .to_string();
930 x.jobs.get(&org, &name)?;
931 let limit = a.get("limit").and_then(Value::as_u64).unwrap_or(50) as usize;
932 Ok(json!({"runs": x.jobs.runs(&org, &name).list(limit)}))
933 }
934 );
935 Ok(())
936}
937
938fn job_run_log_tool(r: &mut Registry, ctx: &Ctx, ann: &Ann) -> Result<()> {
939 tool!(
940 r,
941 ctx,
942 "job_run_log",
943 "Job run log",
944 "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.",
945 obj(
946 json!({"name": {"type": "string"}, "run": {"type": "integer", "minimum": 1}, "offset": {"type": "integer", "minimum": 0}}),
947 &["name", "run"]
948 ),
949 ann.ro,
950 |x: &Ctx, a: Value, _c: &Caller| -> Result<Value> {
951 let org = org_of(&a)?;
952 let name = a
953 .get("name")
954 .and_then(Value::as_str)
955 .unwrap_or_default()
956 .to_string();
957 x.jobs.get(&org, &name)?;
958 let run = a.get("run").and_then(Value::as_u64).unwrap_or(0);
959 let offset = a.get("offset").and_then(Value::as_u64).unwrap_or(0);
960 let store = x.jobs.runs(&org, &name);
961 let (text, offset, finished) = store.log(run, offset)?;
962 Ok(
963 json!({"text": text, "offset": offset, "finished": finished, "run": store.get(run)?}),
964 )
965 }
966 );
967 Ok(())
968}