Schedule icon
Query icon
If icon
Switch icon
Log icon
SlackIncomingWebhook icon

Couchbase TTL expiry triage sentinel

Detect Couchbase documents nearing TTL expiry via N1QL and either dry-run log or quarantine them, with Slack alerts.

Categories
DataInfrastructureinfrastructure

Diagram unavailable

We could not build the topology for this blueprint. The flow itself is valid, use the YAML on the left to run it.

Couchbase automatically ejects a document once its TTL (Time-To-Live) expiration passes, with no built-in way to intervene first. For session stores, token registries, or compliance logs, that is often too late — someone needed to archive, audit, or flag the record before it was gone. This blueprint runs a N1QL query against META().expiration to find documents expiring within a threshold window, then either logs what it found (dryRun: true, the safe default) or actively quarantines the matched documents (dryRun: false) by tagging them, before they disappear.

How it works

  1. schedule-couchbase-ttl (io.kestra.plugin.core.trigger.Schedule) runs every 6 hours. Shipped disabled so you can validate a manual run first.
  2. query_expiring_docs (io.kestra.plugin.couchbase.Query) runs a N1QL SELECT against `bucketName`.`scopeName`.`collectionName`, filtering META(d).expiration > 0 (0 means no TTL set) and <= NOW_MILLIS() / 1000 + thresholdSeconds. fetchType: FETCH populates rows with every matching document's ID and expiration, and size with the count.
  3. evaluate_expiring_docs (io.kestra.plugin.core.flow.If) only proceeds into triage when outputs.query_expiring_docs.size > 0.
  4. switch_dry_run (io.kestra.plugin.core.flow.Switch) branches on inputs.dryRun:
    • "true": log_dry_run_triage logs the matched count and document IDs. Nothing in Couchbase changes.
    • "false": quarantine_expiring_docs (io.kestra.plugin.couchbase.Query) runs a N1QL UPDATE tagging every matched document with ttlTriageStatus = "quarantined" and ttlTriagedAt = NOW_STR(), using fetchType: NONE since the mutation has nothing to fetch. log_active_triage then logs that it ran.
  5. When no documents match, log_no_expiring_docs records a clean run instead.
  6. notify_slack posts one summary message every run: the flagged count, the dryRun value, and whether the quarantine task actually executed.

What you get

  • A scheduled early-warning sweep for documents about to be permanently ejected by Couchbase's TTL mechanism.
  • A safe-by-default dry-run mode that only logs, with an explicit flag required to actually tag documents.
  • A quarantine marker (ttlTriageStatus/ttlTriagedAt) added without touching the document's existing expiration clock.
  • One Slack message per run covering all three outcomes: clean, dry-run audit, and active quarantine.

Who it's for

  • Platform and data engineering teams running Couchbase for session stores, token registries, or compliance logs where losing a record silently is a problem.
  • SREs who want a heads-up before a TTL-driven ejection, not a post-mortem after it.
  • Anyone who has had Couchbase eject a document that should have been archived first.

Why orchestrate this with Kestra

Couchbase's TTL mechanism is a one-way clock: there is no native hook to intervene before ejection. Kestra adds the missing piece — a schedule, a branching decision between audit-only and active triage, and an execution history showing exactly which documents were flagged on every run. The same flow can later fan out to an archival task or a ticketing system by adding one more branch.

Prerequisites

  • A reachable Couchbase cluster with the target bucket/scope/collection and TTLs set on at least some documents to see non-empty results. Locally: docker run -d --name couchbase -p 8091-8096:8091-8096 -p 11210-11211:11210-11211 couchbase/server:community, then if the Kestra worker runs in a separate Docker network: docker network connect YOUR_KESTRA_NETWORK couchbase.
  • A Couchbase user with read access to run the N1QL SELECT, and write access if dryRun: false is used (the UPDATE needs document-mutate permission on the collection).
  • A Slack incoming webhook for the summary notification.

Secrets

  • COUCHBASE_PASSWORD: password for username (Administrator by default), used by both Query tasks.
  • SLACK_WEBHOOK_URL: Slack incoming webhook used by notify_slack.
  • In Kestra OSS (no Enterprise secrets backend), secrets are supplied as environment variables prefixed SECRET_, base64-encoded, and read back in flows with {{ secret('NAME') }} — for example SECRET_COUCHBASE_PASSWORD=$(echo -n 'password' | base64). This keeps credentials out of the flow YAML but, per Kestra's own documentation, offers no encryption at rest or access control beyond the host environment; use the Enterprise secrets backend for stronger guarantees.

Inputs

  • connectionString (STRING, default couchbase://127.0.0.1): cluster connection string.
  • username (STRING, default Administrator): cluster username.
  • bucketName (STRING, default telemetry): bucket to scan.
  • scopeName (STRING, default _default): scope within the bucket.
  • collectionName (STRING, default _default): collection within the scope.
  • thresholdSeconds (INT, default 86400): expiry window in seconds (default 24 hours).
  • dryRun (BOOLEAN, default true): must be explicitly false to run the active quarantine UPDATE.

Outputs

  • outputs.query_expiring_docs.rows / .size: the matched documents (ID and expiration) and how many were found.
  • outputs.quarantine_expiring_docs: present only when dryRun: false and documents matched; fetchType: NONE means it carries no row/rows/uri, only confirms the UPDATE executed.

Quick start

  1. Start Couchbase locally with the command above, create the telemetry bucket (or point at an existing one), and set a short TTL on a few test documents so the query has something to find.
  2. Add the COUCHBASE_PASSWORD and SLACK_WEBHOOK_URL secrets.
  3. Run the flow manually with the defaults (dryRun: true) and confirm log_dry_run_triage lists the expected documents and nothing in Couchbase changes.
  4. Run once with dryRun: false against a test collection and confirm ttlTriageStatus was set on the matched documents.
  5. Enable schedule-couchbase-ttl once you trust the dry-run step; keep dryRun: false as a manual, reviewed flag rather than scheduling it.

How to extend

  • Replace the quarantine UPDATE with an archival step (e.g. write matched rows to S3 or another bucket) before deletion, instead of only tagging them in place.
  • Add a second threshold for a "critical" window (e.g. 1 hour) with its own Slack alert, alongside the existing thresholdSeconds warning window.
  • Swap the Schedule trigger for io.kestra.plugin.couchbase.Trigger to react the moment matching documents appear instead of polling on a cron.
  • Loop over several bucket/scope/collection combinations with io.kestra.plugin.core.flow.Loop if more than one needs TTL triage.

Pitfalls

  • META().expiration of 0 means no TTL, not "already expired." The query explicitly filters expiration > 0 to exclude documents with no TTL set at all; omitting that filter would wrongly match every untimed document as "expiring within the threshold."
  • There is no NOW_NUM() function in Couchbase N1QL. The verified clock function is NOW_MILLIS() (epoch milliseconds); this blueprint divides by 1000 to compare against META().expiration, which is stored in epoch seconds.
  • io.kestra.plugin.couchbase.Query has no bucketName/scopeName/collectionName property. Those three inputs exist only in this flow's own YAML; they are interpolated into the N1QL keyspace path (`bucket`.`scope`.`collection`) inside query, not bound to a task property.
  • The quarantine UPDATE does not touch the document's TTL. Tagging a document with ttlTriageStatus/ttlTriagedAt does not reset or extend Couchbase's own expiration clock; a tagged document can still be ejected on schedule unless you separately clear or extend its TTL.
  • Switch case keys "true"/"false" are quoted strings. inputs.dryRun renders to the literal string "true" or "false"; unquoted true:/false: YAML map keys would parse as booleans instead and would not match.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.