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", "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
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
164/// A job as the tools answer: its settings, plus `created_at`, `updated_at`,
165/// `next_run` and `last_run`, all at the top level.
166fn 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    // --- databases ---------------------------------------------------------
185
186    database_create_tool(r, &ctx, &ann)?;
187    database_list_tool(r, &ctx, &ann)?;
188    database_get_tool(r, &ctx, &ann)?;
189
190    // --- destinations ------------------------------------------------------
191
192    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    // --- backups -----------------------------------------------------------
198
199    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    // --- jobs --------------------------------------------------------------
209
210    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        // `reveal` hands out the password: a secret read (viewers, `read`
339        // tokens refused; always audited).
340        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}