id: mongodb-verified-collection-backup
namespace: company.team
inputs:
- id: database
type: STRING
displayName: Source database
defaults: ecommerce
- id: collection
type: STRING
displayName: Source collection
defaults: orders
- id: bucket
type: STRING
displayName: Backup bucket
description: S3-compatible bucket that receives the archives. It must already exist.
defaults: mongo-backups
- id: restore_database
type: STRING
displayName: Restore drill database
description: Scratch database where each archive is test-restored. Never point
this at a production database.
defaults: restore_check
- id: min_documents
type: INT
displayName: Minimum documents
description: An archive with fewer documents fails the drill. This catches a
wrong collection name or an emptied collection, which would otherwise pass
as 0 = 0.
defaults: 1
- id: notify_slack
type: BOOL
displayName: Notify Slack
description: Post the result to Slack. Needs the SLACK_WEBHOOK_URL secret.
defaults: false
variables:
archive_key: "archives/{{ inputs.database }}.{{ inputs.collection }}/{{
execution.startDate | date('yyyy-MM-dd-HH-mm-ss') }}.archive.gz"
latest_key: "verified/{{ inputs.database }}.{{ inputs.collection
}}/latest-good.archive.gz"
concurrency:
limit: 1
triggers:
- id: nightly
type: io.kestra.plugin.core.trigger.Schedule
description: Back up and restore-test every night at 02:00.
cron: "0 2 * * *"
tasks:
- id: source_indexes
type: io.kestra.plugin.mongodb.Aggregate
description: List the index names on the source collection. A restore without
its indexes is an outage, so the drill checks them too.
connection:
uri: "{{ secret('MONGODB_URI') }}"
database: "{{ inputs.database }}"
collection: "{{ inputs.collection }}"
pipeline:
- $indexStats: {}
- $project:
_id: 0
name: 1
- $sort:
name: 1
- id: dump
type: io.kestra.plugin.scripts.shell.Commands
description: Dump the collection with mongodump into a gzip archive. mongodump
keeps BSON types (ObjectId, Decimal128, dates) and the index definitions.
containerImage: mongo:8.0
env:
MONGODB_URI: "{{ secret('MONGODB_URI') }}"
outputFiles:
- dump.archive.gz
commands:
- mongodump --uri="$MONGODB_URI" --db="{{ inputs.database }}"
--collection="{{ inputs.collection }}" --archive=dump.archive.gz --gzip
2> dump.log || { cat dump.log; exit 1; }
- cat dump.log
- test -s dump.archive.gz && ! grep -q 'does not exist' dump.log || { echo
"Collection {{ inputs.database }}.{{ inputs.collection }} does not
exist, nothing to back up"; exit 1; }
- dumped=$(grep -oE '\(([0-9]+) documents?\)' dump.log | tail -1 | tr -dc
'0-9')
- echo "::{\"outputs\":{\"dumped\":${dumped:-0}}}::"
- id: upload_archive
type: io.kestra.plugin.minio.Upload
description: Write the archive to a timestamped key. It is not trusted until the
restore drill passes.
endpoint: "{{ secret('MINIO_ENDPOINT') }}"
accessKeyId: "{{ secret('MINIO_ACCESS_KEY_ID') }}"
secretKeyId: "{{ secret('MINIO_SECRET_KEY_ID') }}"
bucket: "{{ inputs.bucket }}"
key: "{{ render(vars.archive_key) }}"
from: "{{ outputs.dump.outputFiles['dump.archive.gz'] }}"
contentType: application/gzip
- id: download_archive
type: io.kestra.plugin.minio.Download
description: Read the archive back from the bucket. The drill restores what is
stored, not the local file.
endpoint: "{{ secret('MINIO_ENDPOINT') }}"
accessKeyId: "{{ secret('MINIO_ACCESS_KEY_ID') }}"
secretKeyId: "{{ secret('MINIO_SECRET_KEY_ID') }}"
bucket: "{{ inputs.bucket }}"
key: "{{ render(vars.archive_key) }}"
- id: restore
type: io.kestra.plugin.scripts.shell.Commands
description: Restore the downloaded archive into the scratch database with
mongorestore, replacing the previous drill.
containerImage: mongo:8.0
env:
MONGODB_URI: "{{ secret('MONGODB_URI') }}"
inputFiles:
restore.archive.gz: "{{ outputs.download_archive.uri }}"
commands:
- mongorestore --uri="$MONGODB_URI" --archive=restore.archive.gz --gzip
--drop --nsFrom="{{ inputs.database }}.{{ inputs.collection }}"
--nsTo="{{ inputs.restore_database }}.{{ inputs.collection }}" 2>
restore.log || { cat restore.log; exit 1; }
- cat restore.log
- failed=$(grep -oE '([0-9]+) document\(s\) failed' restore.log | tail -1
| tr -dc '0-9')
- echo "::{\"outputs\":{\"failed\":${failed:-0}}}::"
- id: restored_count
type: io.kestra.plugin.mongodb.Aggregate
description: Count the documents in the restored copy.
connection:
uri: "{{ secret('MONGODB_URI') }}"
database: "{{ inputs.restore_database }}"
collection: "{{ inputs.collection }}"
pipeline:
- $count: restored
- id: restored_indexes
type: io.kestra.plugin.mongodb.Aggregate
description: List the index names on the restored copy.
connection:
uri: "{{ secret('MONGODB_URI') }}"
database: "{{ inputs.restore_database }}"
collection: "{{ inputs.collection }}"
pipeline:
- $indexStats: {}
- $project:
_id: 0
name: 1
- $sort:
name: 1
- id: drill
type: io.kestra.plugin.core.output.OutputValues
description: The numbers the gate uses, in one place.
values:
dumped: "{{ outputs.dump.vars.dumped }}"
restored: "{{ outputs.restored_count.rows | length > 0 ?
outputs.restored_count.rows[0].restored : 0 }}"
failed: "{{ outputs.restore.vars.failed }}"
source_indexes: "{{ outputs.source_indexes.rows | jq('map(.name) | join(\",\")')
| first }}"
restored_indexes: "{{ outputs.restored_indexes.rows | jq('map(.name) |
join(\",\")') | first }}"
- id: restore_gate
type: io.kestra.plugin.core.flow.If
description: >
Promote only when every dumped document came back, none failed, the
indexes match and the archive is not suspiciously small.
condition: >-
{{ (outputs.drill.values.restored | number) ==
(outputs.drill.values.dumped | number) and (outputs.drill.values.failed |
number) == 0 and outputs.drill.values.restored_indexes ==
outputs.drill.values.source_indexes and (outputs.drill.values.dumped |
number) >= inputs.min_documents }}
then:
- id: promote_archive
type: io.kestra.plugin.minio.Copy
description: Copy the tested archive onto the stable key that restore tooling
reads. The timestamped archive stays as history.
endpoint: "{{ secret('MINIO_ENDPOINT') }}"
accessKeyId: "{{ secret('MINIO_ACCESS_KEY_ID') }}"
secretKeyId: "{{ secret('MINIO_SECRET_KEY_ID') }}"
from:
bucket: "{{ inputs.bucket }}"
key: "{{ render(vars.archive_key) }}"
to:
bucket: "{{ inputs.bucket }}"
key: "{{ render(vars.latest_key) }}"
- id: log_verified
type: io.kestra.plugin.core.log.Log
message: >-
Restore drill passed for {{ inputs.database }}.{{ inputs.collection
}}: {{ outputs.drill.values.restored }} of {{
outputs.drill.values.dumped }} documents restored, indexes [{{
outputs.drill.values.restored_indexes }}] match. Promoted to s3://{{
inputs.bucket }}/{{ render(vars.latest_key) }}.
- id: notify_verified
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.notify_slack }}"
then:
- id: slack_verified
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "MongoDB backup of {{ inputs.database }}.{{ inputs.collection }} passed its restore drill: {{ outputs.drill.values.restored }} documents and indexes [{{ outputs.drill.values.restored_indexes }}] restored, promoted to latest-good. Execution {{ execution.id }}."
}
else:
- id: discard_archive
type: io.kestra.plugin.minio.Delete
description: Remove the archive that failed the drill. The previous latest-good
archive is left as it is.
endpoint: "{{ secret('MINIO_ENDPOINT') }}"
accessKeyId: "{{ secret('MINIO_ACCESS_KEY_ID') }}"
secretKeyId: "{{ secret('MINIO_SECRET_KEY_ID') }}"
bucket: "{{ inputs.bucket }}"
key: "{{ render(vars.archive_key) }}"
- id: notify_failed
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.notify_slack }}"
then:
- id: slack_failed
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "MongoDB backup of {{ inputs.database }}.{{ inputs.collection }} FAILED its restore drill: {{ outputs.drill.values.restored }} restored of {{ outputs.drill.values.dumped }} dumped, {{ outputs.drill.values.failed }} failed, indexes [{{ outputs.drill.values.restored_indexes }}] vs source [{{ outputs.drill.values.source_indexes }}], minimum {{ inputs.min_documents }}. The last good backup was kept. Execution {{ execution.id }}."
}
- id: fail_run
type: io.kestra.plugin.core.execution.Fail
description: End the run as FAILED so a broken backup is visible in Kestra.
errorMessage: "Restore drill failed: {{ outputs.drill.values.restored }}
restored of {{ outputs.drill.values.dumped }} dumped, indexes [{{
outputs.drill.values.restored_indexes }}] vs [{{
outputs.drill.values.source_indexes }}]."
errors:
- id: notify_error
type: io.kestra.plugin.core.flow.If
description: >
A task failed before the drill could decide, for example MongoDB or the
bucket was unreachable. A failed drill already posted its own alert, so
this one is skipped then.
condition: "{{ inputs.notify_slack and (outputs.drill ?? null) is null }}"
then:
- id: slack_error
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "MongoDB backup of {{ inputs.database }}.{{ inputs.collection }} ERRORED before the restore drill finished. No new backup was promoted. Check execution {{ execution.id }}."
}
finally:
- id: clear_scratch
type: io.kestra.plugin.mongodb.Delete
description: Empty the scratch copy whatever the outcome. The next drill drops
and recreates it.
connection:
uri: "{{ secret('MONGODB_URI') }}"
database: "{{ inputs.restore_database }}"
collection: "{{ inputs.collection }}"
operation: DELETE_MANY
filter: {}