Skip to main content

isb_daemon/daemon/
data.rs

1//! The data tools: databases (`database_*`, over apps with a database
2//! source), backups and restores (`backup_*`), and scheduled jobs
3//! (`job_*`). Every one acts in the org it names, so the authorizer's org
4//! check covers them; destinations on this host need a trusted caller.
5
6use 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/// What the data tools work on.
20#[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
31/// Remove the keys a tool handles itself before deserializing the rest.
32fn 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
44/// The caller may point a destination at this host: a superadmin (the
45/// local CLI included) or a platform admin.
46fn 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
57/// Wait (bounded) until run `id` in `store` finishes.
58pub(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
78/// Volume backups, like the rest of a volume's care, are for org admins.
79fn 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
89/// A database app with its connection details (password as a secret
90/// reference; with `password`, the value too).
91pub 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
99/// The properties backup_create and backup_update share.
100fn 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
115/// The properties job_create and job_update share.
116fn 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
135/// The MCP annotations the tools below share.
136struct Ann {
137    ro: Value,
138    destructive: Value,
139    write: Value,
140    // Destinations and backups reach outside isb (S3).
141    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
154/// A command line becomes argv.
155fn 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    // --- databases ---------------------------------------------------------
184
185    database_create_tool(r, &ctx, &ann)?;
186    database_list_tool(r, &ctx, &ann)?;
187    database_get_tool(r, &ctx, &ann)?;
188
189    // --- destinations ------------------------------------------------------
190
191    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    // --- backups -----------------------------------------------------------
197
198    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    // --- jobs --------------------------------------------------------------
208
209    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        // `reveal` hands out the password: a secret read (viewers, `read`
336        // tokens refused; always audited).
337        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}