New to Kestra?
Use blueprints to kickstart your first workflows.
Check Kestra flows in a GitHub repo for 2.0 readiness with kestra-migrate, then open a verified, approved, duplicate-free migration pull request.
id: kestra-2-migration-readiness-gate
namespace: company.team
description: |
Kestra 2.0 Migration Readiness Gate. Scans the Kestra flows of a public GitHub repository at an exact commit with
the pinned kestra-migrate 2.6.1 tool and classifies them CLEAN, AUTO_MIGRATABLE, ADVISORY, BLOCKING or TOOL_ERROR.
Only an AUTO_MIGRATABLE scope can be remediated, and only after a human approves: the gate then rewrites the flows
in a fresh clone, verifies the result, and delivers it as one branch, one issue and one pull request. Running it
again for the same repository, commit and path reuses them instead of opening duplicates.
inputs:
- id: repo_url
type: STRING
displayName: GitHub repository
description: Public GitHub repository, https://github.com/<owner>/<repo>. No
other host, scheme, port or credentials.
defaults: https://github.com/kestra-io/kestra2-flow-migration
validator: '^https://github\.com/[A-Za-z0-9](?:[A-Za-z0-9-]{0,37}[A-Za-z0-9])?/[A-Za-z0-9_][A-Za-z0-9._-]{0,99}$'
- id: commit
type: STRING
displayName: Commit SHA
description: Full 40-character commit SHA to scan and, after approval,
remediate. Branch names are not accepted.
defaults: 41e0fe57268711013841f633bd189cb04d83b198
validator: '^[0-9a-f]{40}$'
- id: flows_path
type: STRING
displayName: Flows path
description: Folder or flow file inside the repository, or "." for the whole
repository. No absolute paths, no "..".
defaults: input-flows/additional-test-cases/ion-read.yaml
validator: '^(?:\.|(?!.*(?:^|/)\.\.?(?:/|$))[A-Za-z0-9._][A-Za-z0-9._-]{0,99}(?:/[A-Za-z0-9._][A-Za-z0-9._-]{0,99}){0,19})$'
tasks:
- id: scan_workdir
type: io.kestra.plugin.core.flow.WorkingDirectory
description: Read-only readiness scan of the requested commit. Nothing in this
checkout is ever modified.
tasks:
- id: clone_a
type: io.kestra.plugin.git.Clone
description: Clones the repository and checks out the exact commit; submodules
are never fetched.
url: "{{ inputs.repo_url }}"
commit: "{{ inputs.commit }}"
cloneSubmodules: false
directory: repo
- id: verify_a
type: io.kestra.plugin.scripts.shell.Commands
description: Runs on the Kestra worker before any container starts. Checks that
the checkout is exactly the requested commit, refuses repositories
containing symbolic links, and keeps flows_path inside the clone.
taskRunner:
type: io.kestra.plugin.core.runner.Process
env:
REQUESTED_COMMIT: "{{ inputs.commit }}"
FLOWS_PATH: "{{ inputs.flows_path }}"
commands:
- |
set -eu
HEAD_SHA=$(cat repo/.git/HEAD)
if [ "$HEAD_SHA" != "$REQUESTED_COMMIT" ]; then echo "Checked-out HEAD '$HEAD_SHA' does not match the requested commit $REQUESTED_COMMIT."; exit 1; fi
LINKS=$(find repo -type l -printf '%P -> %l\n')
if [ -n "$LINKS" ]; then printf 'The repository contains symbolic links, which this gate refuses to process:\n%s\n' "$LINKS"; exit 1; fi
ROOT=$(realpath -e repo)
TARGET=$(realpath -m -- "repo/$FLOWS_PATH")
case "$TARGET/" in "$ROOT"/*) ;; *) echo "flows_path resolves outside the repository."; exit 1 ;; esac
echo "Verified checkout $HEAD_SHA: no symbolic links, flows_path stays inside the repository."
echo "::{\"outputs\":{\"headSha\":\"$HEAD_SHA\"}}::"
- id: install_a
type: io.kestra.plugin.scripts.shell.Commands
description: Downloads kestra-migrate 2.6.1 from its GitHub release and refuses
it unless the SHA-256 matches the pinned value.
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: buildpack-deps:trixie-scm@sha256:ed6425e2963332c914dd78c0768b4b0f86680ac5f3d630f39f91ec8aa82384a1
commands:
- |
set -eu
VERSION=2.6.1
case "$(uname -m)" in
x86_64) ARCH=amd64; WANT=84e47ab533e7ae159cc6e9b78a5962dc80554742faf4a5b3cbd4c957d8ee57f2 ;;
aarch64) ARCH=arm64; WANT=46cbb256f7ab862a5b1dc85038bf00eba60d6f777446f8c82a352769458b1ba9 ;;
*) echo "Unsupported architecture $(uname -m)."; exit 1 ;;
esac
ASSET="kestra-migrate_${VERSION}_linux_${ARCH}"
mkdir -p bin
wget -q --tries=3 -O "bin/$ASSET" "https://github.com/kestra-io/kestra2-flow-migration/releases/download/${VERSION}/${ASSET}"
GOT=$(sha256sum "bin/$ASSET" | cut -d' ' -f1)
if [ "$GOT" != "$WANT" ]; then echo "kestra-migrate checksum mismatch: expected $WANT, got $GOT."; exit 1; fi
chmod +x "bin/$ASSET"; mv "bin/$ASSET" bin/kestra-migrate
echo "kestra-migrate $VERSION ($ARCH) verified, sha256 $GOT"
- id: scan
type: io.kestra.plugin.scripts.shell.Commands
description: Runs kestra-migrate --check --summary and parses its output with a
strict parser; any unexpected line, count mismatch or missing summary
makes the result TOOL_ERROR.
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: buildpack-deps:trixie-scm@sha256:ed6425e2963332c914dd78c0768b4b0f86680ac5f3d630f39f91ec8aa82384a1
env:
FLOWS_PATH: "{{ inputs.flows_path }}"
inputFiles:
parse.awk: |
# Parses ANSI-stripped `kestra-migrate --check --summary` stdout (2.6.1) into a Kestra outputs line.
# Contract kestra-migrate-scan/1.1: every non-blank line must match a documented output shape, else TOOL_ERROR.
# Portable to BusyBox awk and mawk: no 3-arg match(), no gensub(), no arrays of arrays.
function esc(s) { gsub(/\\/, "\\\\", s); gsub(/"/, "\\\"", s); return s }
function bad(where) { unknown++; print "UNRECOGNISED(" where "): " $0 > "/dev/stderr" }
function flush() {
if (cur == "") return
if (perr) cls = "BLOCKING"
else if (nb > 0) cls = "BLOCKING"
else if (na > 0) cls = "ADVISORY"
else if (rw) cls = "AUTO_MIGRATABLE"
else cls = "CLEAN"
count[cls]++
warnBlk += nb; warnAdv += na
if (nb + na > 0) affected++
if (!perr) summarised++
list[cls] = list[cls] (list[cls] == "" ? "" : ",") "\"" esc(cur) "\""
if (rw) rewrittenList = rewrittenList (rewrittenList == "" ? "" : ",") "\"" esc(cur) "\""
flows = flows (flows == "" ? "" : ",") sprintf("{\"file\":\"%s\",\"class\":\"%s\",\"rewritten\":%s,\"blocking\":%d,\"advisory\":%d,\"parseError\":%s}", esc(cur), cls, rw ? "true" : "false", nb, na, perr ? "true" : "false")
nflows++
cur = ""
}
function startFlow(name, rewritten, parseError) { flush(); cur = name; rw = rewritten; perr = parseError; nb = 0; na = 0; inDiff = 0 }
BEGIN { tallyTotal = -1; tallyNeed = -1 }
# Nothing but blank lines may follow the final tally.
afterTally { if ($0 != "") bad("after tally"); next }
# 1. Tally (exact sentences printed by main.go runCheck). It ends the output and closes the summary.
/^⚠ [0-9]+\/[0-9]+ flows need migration$/ { flush(); split($2, p, "/"); tallyNeed = p[1] + 0; tallyTotal = p[2] + 0; afterTally = 1; inSummary = 0; next }
/^✔ All [0-9]+ flows are v2-compatible$/ { flush(); tallyNeed = 0; tallyTotal = $3 + 0; afterTally = 1; inSummary = 0; next }
# Summary block (report.Summarize): only its documented line shapes are accepted; its totals are
# reconciled with the per-flow warnings in END, so a truncated or altered summary cannot pass.
/^Summary: [0-9]+ warnings? across [0-9]+ flows? — [0-9]+ blocking, [0-9]+ advisory$/ && !inSummary && !summarySeen {
flush(); inSummary = 1; summarySeen = 1; section = ""; afterRow = 0
hdrWarn = $2 + 0; hdrFlows = $5 + 0; hdrBlk = $8 + 0; hdrAdv = $10 + 0; next
}
inSummary {
if ($0 == "") afterRow = 0
else if ($0 == " BLOCKING (Kestra 2.0 rejects the flow)" && section == "") { section = "B"; afterRow = 0 }
else if ($0 == " ADVISORY (deploys, breaks at run time)" && section != "A") { section = "A"; afterRow = 0 }
else if (section == "B" && $0 ~ /^ [ 0-9]*[0-9]× ✗ .*→ [^ ]+$/) { rowBlk += $1 + 0; afterRow = 1 }
else if (section == "A" && $0 ~ /^ [ 0-9]*[0-9]× ⚠ .*→ [^ ]+$/) { rowAdv += $1 + 0; afterRow = 1 }
else if (afterRow && $0 ~ /^ ([0-9]+ [^ ]|in [0-9]+ flows$)/) { }
else bad("summary")
next
}
# Unified diff of a rewritten flow: consume exactly the line counts announced by each hunk header.
/^--- original$/ && cur != "" { inDiff = 1; next }
/^\+\+\+ migrated$/ && inDiff { next }
/^@@ -[0-9,]+ \+[0-9,]+ @@/ && inDiff {
split($2, o, ","); split($3, n, ",")
oldLeft = (o[2] == "") ? 1 : o[2] + 0; newLeft = (n[2] == "") ? 1 : n[2] + 0; next
}
inDiff && (oldLeft > 0 || newLeft > 0) {
c = substr($0, 1, 1)
if (c == " ") { oldLeft--; newLeft-- }
else if (c == "-") { oldLeft-- }
else if (c == "+") { newLeft-- }
else if ($0 == "") { oldLeft--; newLeft-- }
else if (c == "\\") { }
else { unknown++; print "UNRECOGNISED(diff): " $0 > "/dev/stderr" }
next
}
# 2. Clean flow
/^✔ [^ ]/ { startFlow(substr($0, length("✔ ") + 1), 0, 0); flush(); next }
# 3. Flow needing work: rewritten (✎) or warning-only (⚠)
/^✎ [^ ]/ { startFlow(substr($0, length("✎ ") + 1), 1, 0); next }
/^⚠ [^ ]/ { startFlow(substr($0, length("⚠ ") + 1), 0, 0); next }
# 4. Per-flow processing error (YAML parse error): "✗ <file>: <error>"
/^✗ [^ ]/ { s = substr($0, length("✗ ") + 1); i = index(s, ": "); startFlow(i ? substr(s, 1, i - 1) : s, 0, 1); perrMsg[cur] = i ? substr(s, i + 2) : ""; flush(); next }
# 5. Warnings with severity, and their docs links
/^ ✗ / && cur != "" { nb++; inDiff = 0; next }
/^ ⚠ / && cur != "" { na++; inDiff = 0; next }
/^ ↳ docs: https?:\/\// { next }
/^$/ { next }
{ unknown++; print "UNRECOGNISED: " $0 > "/dev/stderr" }
# 6. Structured result
END {
flush()
reason = ""; code = ""
# runCheck only summarises flows that parsed, and only prints a summary for 2+ of them with warnings.
summaryExpected = (summarised >= 2 && warnBlk + warnAdv > 0)
if (tallyTotal < 0) { code = "NO_TALLY"; reason = "no tally line: kestra-migrate did not complete" }
else if (tallyTotal == 0) { code = "ZERO_FLOWS"; reason = "no flows found under the scanned path" }
else if (nflows != tallyTotal) { code = "COUNT_MISMATCH"; reason = sprintf("parsed %d flows but tally reports %d", nflows, tallyTotal) }
else if (nflows - count["CLEAN"] != tallyNeed) { code = "COUNT_MISMATCH"; reason = sprintf("parsed %d flows needing migration but tally reports %d", nflows - count["CLEAN"], tallyNeed) }
else if (unknown > 0) { code = "UNRECOGNISED_OUTPUT"; reason = sprintf("%d unrecognised output lines", unknown) }
else if (summarySeen != summaryExpected) { code = "SUMMARY_MISMATCH"; reason = summarySeen ? "summary printed where none was expected" : "summary missing: run kestra-migrate with --check --summary" }
else if (summarySeen && (hdrBlk != warnBlk || rowBlk != warnBlk || hdrAdv != warnAdv || rowAdv != warnAdv || hdrWarn != warnBlk + warnAdv || hdrFlows != affected)) {
code = "SUMMARY_MISMATCH"
reason = sprintf("summary reports %d blocking/%d advisory in %d flows (rows %d/%d) but flows carry %d/%d in %d", hdrBlk, hdrAdv, hdrFlows, rowBlk, rowAdv, warnBlk, warnAdv, affected)
}
if (reason != "") status = "TOOL_ERROR"
else if (count["BLOCKING"] > 0) status = "BLOCKING"
else if (count["ADVISORY"] > 0) status = "ADVISORY"
else if (count["AUTO_MIGRATABLE"] > 0) status = "AUTO_MIGRATABLE"
else status = "CLEAN"
printf "::{\"outputs\":{\"contract\":\"kestra-migrate-scan/1.1\",\"status\":\"%s\",\"reasonCode\":\"%s\",\"reason\":\"%s\",\"total\":%d,\"tallyTotal\":%d,\"tallyNeedMigration\":%d,\"unrecognisedLines\":%d,", status, code, esc(reason), nflows, tallyTotal, tallyNeed, unknown
printf "\"counts\":{\"CLEAN\":%d,\"AUTO_MIGRATABLE\":%d,\"ADVISORY\":%d,\"BLOCKING\":%d},", count["CLEAN"], count["AUTO_MIGRATABLE"], count["ADVISORY"], count["BLOCKING"]
printf "\"clean\":[%s],\"autoMigratable\":[%s],\"advisory\":[%s],\"blocking\":[%s],\"rewritten\":[%s],\"flows\":[%s]}}::\n", list["CLEAN"], list["AUTO_MIGRATABLE"], list["ADVISORY"], list["BLOCKING"], rewrittenList, flows
}
outputFiles:
- scan-report.txt
- scan-result.line
commands:
- |
set -eu
GOT=$(sha256sum parse.awk | cut -d' ' -f1)
if [ "$GOT" != 1b38b4c64e3473a0704005dc0a562cc70d2d6ff6bf12202fde03f6d412074593 ]; then echo "Parser integrity check failed: got $GOT."; exit 1; fi
rc=0
KESTRA_MIGRATE_NO_UPDATE_CHECK=1 ./bin/kestra-migrate --check --summary -- "repo/$FLOWS_PATH" > scan-raw.txt 2> scan-stderr.txt || rc=$?
echo "kestra-migrate exit code $rc (diagnostic only, not used for classification)"
if [ -s scan-stderr.txt ]; then sed 's/^/kestra-migrate stderr: /' scan-stderr.txt; fi
sed "s/$(printf '\033')\[[0-9;]*m//g" scan-raw.txt > scan-report.txt
awk -f parse.awk scan-report.txt > scan-result.line
echo "::{\"outputs\":{\"toolExitCode\":$rc,\"parserSha256\":\"$GOT\"}}::"
cat scan-result.line
- id: contract_guard
type: io.kestra.plugin.core.flow.If
description: Stops if the scan result does not follow the expected contract.
condition: "{{ outputs.scan.vars.contract != 'kestra-migrate-scan/1.1' }}"
then:
- id: contract_mismatch
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Unexpected scan contract '{{ outputs.scan.vars.contract }}'."
- id: scan_summary
type: io.kestra.plugin.core.log.Log
message: >-
SCAN repo={{ inputs.repo_url }} commit={{ outputs.verify_a.vars.headSha }}
flowsPath={{ inputs.flows_path }} status={{ outputs.scan.vars.status }}
reasonCode={{ outputs.scan.vars.reasonCode }} total={{
outputs.scan.vars.total }} clean={{ outputs.scan.vars.counts.CLEAN }}
autoMigratable={{ outputs.scan.vars.counts.AUTO_MIGRATABLE }} advisory={{
outputs.scan.vars.counts.ADVISORY }} blocking={{
outputs.scan.vars.counts.BLOCKING }} rewritten={{
outputs.scan.vars.rewritten | length }}
- id: policy
type: io.kestra.plugin.core.flow.Switch
description: Remediation policy. Only AUTO_MIGRATABLE goes to the approval gate;
CLEAN, ADVISORY and BLOCKING end with a report, TOOL_ERROR fails the
execution.
value: "{{ outputs.scan.vars.status }}"
cases:
CLEAN:
- id: outcome_clean
type: io.kestra.plugin.core.debug.Return
format: '{"outcome":"NO_REMEDIATION_REQUIRED","scanStatus":"CLEAN"}'
ADVISORY:
- id: outcome_advisory
type: io.kestra.plugin.core.debug.Return
format: '{"outcome":"NOT_ELIGIBLE","scanStatus":"ADVISORY","why":"kestra-migrate
cannot rewrite advisory findings safely; they need a human fix."}'
BLOCKING:
- id: outcome_blocking
type: io.kestra.plugin.core.debug.Return
format: '{"outcome":"NOT_ELIGIBLE","scanStatus":"BLOCKING","why":"Kestra 2.0
rejects these flows; kestra-migrate -o cannot fix them."}'
TOOL_ERROR:
- id: scan_failed
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Migration scan cannot be trusted: {{ outputs.scan.vars.reasonCode
}}: {{ outputs.scan.vars.reason }}"
AUTO_MIGRATABLE:
- id: show_proposal
type: io.kestra.plugin.core.log.Log
description: Logs what an approval would change, right before the execution
pauses.
message: |
REMEDIATION PROPOSAL (approval required)
repository: {{ inputs.repo_url }}
commit: {{ outputs.verify_a.vars.headSha }}
flows_path: {{ inputs.flows_path }}
classification: {{ outputs.scan.vars.status }}
total={{ outputs.scan.vars.total }} clean={{ outputs.scan.vars.counts.CLEAN }} autoMigratable={{ outputs.scan.vars.counts.AUTO_MIGRATABLE }} advisory={{ outputs.scan.vars.counts.ADVISORY }} blocking={{ outputs.scan.vars.counts.BLOCKING }}
blocking warnings={{ outputs.scan.vars.flows | jq('[.[].blocking] | add // 0') | first }} advisory warnings={{ outputs.scan.vars.flows | jq('[.[].advisory] | add // 0') | first }}
files that will be rewritten ({{ outputs.scan.vars.rewritten | length }}): {{ outputs.scan.vars.rewritten | join(', ') }}
migration tool: kestra-migrate 2.6.1 (-o, in a fresh clone of the same commit); modifies files: {{ (outputs.scan.vars.rewritten | length) > 0 }}
on APPROVE: the rewrite is verified, then pushed to a kestra-migration/ branch of {{ inputs.repo_url }} with one issue and one pull request (reused if they already exist)
on REJECT or no answer within 1 day: nothing is changed
proposed changes (kestra-migrate --check output):
{{ read(outputs.scan.outputFiles['scan-report.txt']) | abbreviate(8000) }}
- id: approve
type: io.kestra.plugin.core.flow.Pause
description: Human approval gate. Resume with decision APPROVE to remediate;
REJECT (the default) or no answer within 1 day changes nothing.
pauseDuration: P1D
behavior: CANCEL
onResume:
- id: decision
type: SELECT
values:
- REJECT
- APPROVE
defaults: REJECT
description: APPROVE rewrites the listed files in a fresh clone, verifies them,
then pushes a branch and opens one issue and one pull request in
the repository. REJECT changes nothing.
- id: reviewer
type: STRING
description: Name of the person taking the decision.
validator: '^[A-Za-z0-9][A-Za-z0-9 ._@-]{1,63}$'
- id: decision_gate
type: io.kestra.plugin.core.flow.If
description: Remediates and delivers only on an explicit APPROVE.
condition: "{{ outputs.approve.onResume.decision == 'APPROVE' }}"
else:
- id: outcome_rejected
type: io.kestra.plugin.core.debug.Return
format: '{"outcome":"REJECTED","scanStatus":"AUTO_MIGRATABLE","reviewer":"{{
outputs.approve.onResume.reviewer }}"}'
then:
- id: remediation_workdir
type: io.kestra.plugin.core.flow.WorkingDirectory
description: Remediation in a fresh clone of the approved commit.
tasks:
- id: clone_b
type: io.kestra.plugin.git.Clone
description: Fresh clone of the same commit; remediation never reuses the
scanned checkout.
url: "{{ inputs.repo_url }}"
commit: "{{ inputs.commit }}"
cloneSubmodules: false
directory: repo
- id: verify_b
type: io.kestra.plugin.scripts.shell.Commands
description: Runs on the Kestra worker before any container starts. Checks that
the checkout is exactly the requested commit, refuses
repositories containing symbolic links, and keeps flows_path
inside the clone.
taskRunner:
type: io.kestra.plugin.core.runner.Process
env:
REQUESTED_COMMIT: "{{ inputs.commit }}"
FLOWS_PATH: "{{ inputs.flows_path }}"
commands:
- |
set -eu
HEAD_SHA=$(cat repo/.git/HEAD)
if [ "$HEAD_SHA" != "$REQUESTED_COMMIT" ]; then echo "Checked-out HEAD '$HEAD_SHA' does not match the requested commit $REQUESTED_COMMIT."; exit 1; fi
LINKS=$(find repo -type l -printf '%P -> %l\n')
if [ -n "$LINKS" ]; then printf 'The repository contains symbolic links, which this gate refuses to process:\n%s\n' "$LINKS"; exit 1; fi
ROOT=$(realpath -e repo)
TARGET=$(realpath -m -- "repo/$FLOWS_PATH")
case "$TARGET/" in "$ROOT"/*) ;; *) echo "flows_path resolves outside the repository."; exit 1 ;; esac
echo "Verified checkout $HEAD_SHA: no symbolic links, flows_path stays inside the repository."
echo "::{\"outputs\":{\"headSha\":\"$HEAD_SHA\"}}::"
- id: install_b
type: io.kestra.plugin.scripts.shell.Commands
description: Downloads kestra-migrate 2.6.1 from its GitHub release and refuses
it unless the SHA-256 matches the pinned value.
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: buildpack-deps:trixie-scm@sha256:ed6425e2963332c914dd78c0768b4b0f86680ac5f3d630f39f91ec8aa82384a1
commands:
- |
set -eu
VERSION=2.6.1
case "$(uname -m)" in
x86_64) ARCH=amd64; WANT=84e47ab533e7ae159cc6e9b78a5962dc80554742faf4a5b3cbd4c957d8ee57f2 ;;
aarch64) ARCH=arm64; WANT=46cbb256f7ab862a5b1dc85038bf00eba60d6f777446f8c82a352769458b1ba9 ;;
*) echo "Unsupported architecture $(uname -m)."; exit 1 ;;
esac
ASSET="kestra-migrate_${VERSION}_linux_${ARCH}"
mkdir -p bin
wget -q --tries=3 -O "bin/$ASSET" "https://github.com/kestra-io/kestra2-flow-migration/releases/download/${VERSION}/${ASSET}"
GOT=$(sha256sum "bin/$ASSET" | cut -d' ' -f1)
if [ "$GOT" != "$WANT" ]; then echo "kestra-migrate checksum mismatch: expected $WANT, got $GOT."; exit 1; fi
chmod +x "bin/$ASSET"; mv "bin/$ASSET" bin/kestra-migrate
echo "kestra-migrate $VERSION ($ARCH) verified, sha256 $GOT"
- id: remediate
type: io.kestra.plugin.scripts.shell.Commands
description: Re-scans and requires the exact approved result, runs
kestra-migrate -o, then checks that only the expected files
changed (no new, deleted, re-moded or symlinked files,
nothing outside flows_path), that a post-scan is CLEAN, and
that the binary patch re-applies to the source commit and
reproduces the same tree.
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: buildpack-deps:trixie-scm@sha256:ed6425e2963332c914dd78c0768b4b0f86680ac5f3d630f39f91ec8aa82384a1
env:
FLOWS_PATH: "{{ inputs.flows_path }}"
REQUESTED_COMMIT: "{{ inputs.commit }}"
REPO_URL: "{{ inputs.repo_url }}"
inputFiles:
phase-a-scan.line: "{{ outputs.scan.outputFiles['scan-result.line'] }}"
parse.awk: |
# Parses ANSI-stripped `kestra-migrate --check --summary` stdout (2.6.1) into a Kestra outputs line.
# Contract kestra-migrate-scan/1.1: every non-blank line must match a documented output shape, else TOOL_ERROR.
# Portable to BusyBox awk and mawk: no 3-arg match(), no gensub(), no arrays of arrays.
function esc(s) { gsub(/\\/, "\\\\", s); gsub(/"/, "\\\"", s); return s }
function bad(where) { unknown++; print "UNRECOGNISED(" where "): " $0 > "/dev/stderr" }
function flush() {
if (cur == "") return
if (perr) cls = "BLOCKING"
else if (nb > 0) cls = "BLOCKING"
else if (na > 0) cls = "ADVISORY"
else if (rw) cls = "AUTO_MIGRATABLE"
else cls = "CLEAN"
count[cls]++
warnBlk += nb; warnAdv += na
if (nb + na > 0) affected++
if (!perr) summarised++
list[cls] = list[cls] (list[cls] == "" ? "" : ",") "\"" esc(cur) "\""
if (rw) rewrittenList = rewrittenList (rewrittenList == "" ? "" : ",") "\"" esc(cur) "\""
flows = flows (flows == "" ? "" : ",") sprintf("{\"file\":\"%s\",\"class\":\"%s\",\"rewritten\":%s,\"blocking\":%d,\"advisory\":%d,\"parseError\":%s}", esc(cur), cls, rw ? "true" : "false", nb, na, perr ? "true" : "false")
nflows++
cur = ""
}
function startFlow(name, rewritten, parseError) { flush(); cur = name; rw = rewritten; perr = parseError; nb = 0; na = 0; inDiff = 0 }
BEGIN { tallyTotal = -1; tallyNeed = -1 }
# Nothing but blank lines may follow the final tally.
afterTally { if ($0 != "") bad("after tally"); next }
# 1. Tally (exact sentences printed by main.go runCheck). It ends the output and closes the summary.
/^⚠ [0-9]+\/[0-9]+ flows need migration$/ { flush(); split($2, p, "/"); tallyNeed = p[1] + 0; tallyTotal = p[2] + 0; afterTally = 1; inSummary = 0; next }
/^✔ All [0-9]+ flows are v2-compatible$/ { flush(); tallyNeed = 0; tallyTotal = $3 + 0; afterTally = 1; inSummary = 0; next }
# Summary block (report.Summarize): only its documented line shapes are accepted; its totals are
# reconciled with the per-flow warnings in END, so a truncated or altered summary cannot pass.
/^Summary: [0-9]+ warnings? across [0-9]+ flows? — [0-9]+ blocking, [0-9]+ advisory$/ && !inSummary && !summarySeen {
flush(); inSummary = 1; summarySeen = 1; section = ""; afterRow = 0
hdrWarn = $2 + 0; hdrFlows = $5 + 0; hdrBlk = $8 + 0; hdrAdv = $10 + 0; next
}
inSummary {
if ($0 == "") afterRow = 0
else if ($0 == " BLOCKING (Kestra 2.0 rejects the flow)" && section == "") { section = "B"; afterRow = 0 }
else if ($0 == " ADVISORY (deploys, breaks at run time)" && section != "A") { section = "A"; afterRow = 0 }
else if (section == "B" && $0 ~ /^ [ 0-9]*[0-9]× ✗ .*→ [^ ]+$/) { rowBlk += $1 + 0; afterRow = 1 }
else if (section == "A" && $0 ~ /^ [ 0-9]*[0-9]× ⚠ .*→ [^ ]+$/) { rowAdv += $1 + 0; afterRow = 1 }
else if (afterRow && $0 ~ /^ ([0-9]+ [^ ]|in [0-9]+ flows$)/) { }
else bad("summary")
next
}
# Unified diff of a rewritten flow: consume exactly the line counts announced by each hunk header.
/^--- original$/ && cur != "" { inDiff = 1; next }
/^\+\+\+ migrated$/ && inDiff { next }
/^@@ -[0-9,]+ \+[0-9,]+ @@/ && inDiff {
split($2, o, ","); split($3, n, ",")
oldLeft = (o[2] == "") ? 1 : o[2] + 0; newLeft = (n[2] == "") ? 1 : n[2] + 0; next
}
inDiff && (oldLeft > 0 || newLeft > 0) {
c = substr($0, 1, 1)
if (c == " ") { oldLeft--; newLeft-- }
else if (c == "-") { oldLeft-- }
else if (c == "+") { newLeft-- }
else if ($0 == "") { oldLeft--; newLeft-- }
else if (c == "\\") { }
else { unknown++; print "UNRECOGNISED(diff): " $0 > "/dev/stderr" }
next
}
# 2. Clean flow
/^✔ [^ ]/ { startFlow(substr($0, length("✔ ") + 1), 0, 0); flush(); next }
# 3. Flow needing work: rewritten (✎) or warning-only (⚠)
/^✎ [^ ]/ { startFlow(substr($0, length("✎ ") + 1), 1, 0); next }
/^⚠ [^ ]/ { startFlow(substr($0, length("⚠ ") + 1), 0, 0); next }
# 4. Per-flow processing error (YAML parse error): "✗ <file>: <error>"
/^✗ [^ ]/ { s = substr($0, length("✗ ") + 1); i = index(s, ": "); startFlow(i ? substr(s, 1, i - 1) : s, 0, 1); perrMsg[cur] = i ? substr(s, i + 2) : ""; flush(); next }
# 5. Warnings with severity, and their docs links
/^ ✗ / && cur != "" { nb++; inDiff = 0; next }
/^ ⚠ / && cur != "" { na++; inDiff = 0; next }
/^ ↳ docs: https?:\/\// { next }
/^$/ { next }
{ unknown++; print "UNRECOGNISED: " $0 > "/dev/stderr" }
# 6. Structured result
END {
flush()
reason = ""; code = ""
# runCheck only summarises flows that parsed, and only prints a summary for 2+ of them with warnings.
summaryExpected = (summarised >= 2 && warnBlk + warnAdv > 0)
if (tallyTotal < 0) { code = "NO_TALLY"; reason = "no tally line: kestra-migrate did not complete" }
else if (tallyTotal == 0) { code = "ZERO_FLOWS"; reason = "no flows found under the scanned path" }
else if (nflows != tallyTotal) { code = "COUNT_MISMATCH"; reason = sprintf("parsed %d flows but tally reports %d", nflows, tallyTotal) }
else if (nflows - count["CLEAN"] != tallyNeed) { code = "COUNT_MISMATCH"; reason = sprintf("parsed %d flows needing migration but tally reports %d", nflows - count["CLEAN"], tallyNeed) }
else if (unknown > 0) { code = "UNRECOGNISED_OUTPUT"; reason = sprintf("%d unrecognised output lines", unknown) }
else if (summarySeen != summaryExpected) { code = "SUMMARY_MISMATCH"; reason = summarySeen ? "summary printed where none was expected" : "summary missing: run kestra-migrate with --check --summary" }
else if (summarySeen && (hdrBlk != warnBlk || rowBlk != warnBlk || hdrAdv != warnAdv || rowAdv != warnAdv || hdrWarn != warnBlk + warnAdv || hdrFlows != affected)) {
code = "SUMMARY_MISMATCH"
reason = sprintf("summary reports %d blocking/%d advisory in %d flows (rows %d/%d) but flows carry %d/%d in %d", hdrBlk, hdrAdv, hdrFlows, rowBlk, rowAdv, warnBlk, warnAdv, affected)
}
if (reason != "") status = "TOOL_ERROR"
else if (count["BLOCKING"] > 0) status = "BLOCKING"
else if (count["ADVISORY"] > 0) status = "ADVISORY"
else if (count["AUTO_MIGRATABLE"] > 0) status = "AUTO_MIGRATABLE"
else status = "CLEAN"
printf "::{\"outputs\":{\"contract\":\"kestra-migrate-scan/1.1\",\"status\":\"%s\",\"reasonCode\":\"%s\",\"reason\":\"%s\",\"total\":%d,\"tallyTotal\":%d,\"tallyNeedMigration\":%d,\"unrecognisedLines\":%d,", status, code, esc(reason), nflows, tallyTotal, tallyNeed, unknown
printf "\"counts\":{\"CLEAN\":%d,\"AUTO_MIGRATABLE\":%d,\"ADVISORY\":%d,\"BLOCKING\":%d},", count["CLEAN"], count["AUTO_MIGRATABLE"], count["ADVISORY"], count["BLOCKING"]
printf "\"clean\":[%s],\"autoMigratable\":[%s],\"advisory\":[%s],\"blocking\":[%s],\"rewritten\":[%s],\"flows\":[%s]}}::\n", list["CLEAN"], list["AUTO_MIGRATABLE"], list["ADVISORY"], list["BLOCKING"], rewrittenList, flows
}
outputFiles:
- remediation.patch
- remediation-result.json
- remediation/*.txt
commands:
- |
set -eu
umask 022
export KESTRA_MIGRATE_NO_UPDATE_CHECK=1 GIT_CONFIG_NOSYSTEM=1 GIT_CONFIG_GLOBAL=/dev/null HOME=/nonexistent
GOT=$(sha256sum parse.awk | cut -d' ' -f1)
if [ "$GOT" != 1b38b4c64e3473a0704005dc0a562cc70d2d6ff6bf12202fde03f6d412074593 ]; then echo "Parser integrity check failed: got $GOT."; exit 1; fi
gitr() { git -C repo -c safe.directory='*' -c core.quotePath=false -c core.fsmonitor=false -c core.hooksPath=/dev/null "$@"; }
ESC=$(printf '\033')
strip() { sed "s/${ESC}\[[0-9;]*m//g"; }
esc() { sed 's/\\/\\\\/g; s/"/\\"/g' | tr '\n\r\t' ' '; }
jlist() { awk 'BEGIN{printf "["} {gsub(/\\/,"\\\\"); gsub(/"/,"\\\""); printf "%s\"%s\"", (NR>1?",":""), $0} END{printf "]"}' "$1"; }
W=remediation; mkdir -p "$W"
TARGET="repo/$FLOWS_PATH"
if [ "$FLOWS_PATH" = "." ]; then PREFIX=""; OUTDIR="$TARGET"
elif [ -d "$TARGET" ]; then PREFIX="$FLOWS_PATH/"; OUTDIR="$TARGET"
else PREFIX="$(dirname -- "$FLOWS_PATH")/"; [ "$PREFIX" = "./" ] && PREFIX=""; OUTDIR=$(dirname -- "$TARGET"); fi
STATUS=REMEDIATED; CODE=""; REASON=""
fail_with() { if [ "$STATUS" = REMEDIATED ]; then STATUS=FAILED; CODE=$1; REASON=$2; echo "Remediation check failed: $1: $2"; fi; }
# 1. The fresh clone must reproduce the approved Phase A scan byte for byte.
rc_pre=0
./bin/kestra-migrate --check --summary -- "$TARGET" > "$W/pre-raw.txt" 2> "$W/pre-stderr.txt" || rc_pre=$?
strip < "$W/pre-raw.txt" > "$W/pre-report.txt"
awk -f parse.awk "$W/pre-report.txt" > "$W/pre.line" 2> /dev/null
cmp -s "$W/pre.line" phase-a-scan.line || fail_with PRECHECK_MISMATCH "The fresh clone does not reproduce the approved Phase A scan."
grep -q '"status":"AUTO_MIGRATABLE"' "$W/pre.line" || fail_with NOT_ELIGIBLE "Only an AUTO_MIGRATABLE scope may be remediated."
sed -n 's/^✎ //p' "$W/pre-report.txt" | awk -v p="$PREFIX" '{print p $0}' | LC_ALL=C sort > "$W/expected.txt"
# 2. Rewrite in place with the pinned tool. Exit code alone is not trusted.
rc_mig=-1
if [ "$STATUS" = REMEDIATED ]; then
rc_mig=0
./bin/kestra-migrate -o "$OUTDIR" -- "$TARGET" > "$W/migrate-stdout.txt" 2> "$W/migrate-stderr.raw" || rc_mig=$?
strip < "$W/migrate-stderr.raw" > "$W/migrate-stderr.txt"
if [ "$rc_mig" -ne 0 ] || grep -q '^error:' "$W/migrate-stderr.txt"; then
fail_with MIGRATION_TOOL_ERROR "kestra-migrate -o exited $rc_mig: $(head -c 300 "$W/migrate-stderr.txt" | esc)"
elif [ -s "$W/migrate-stderr.txt" ]; then
fail_with MIGRATION_WARNINGS "kestra-migrate -o reported warnings for an AUTO_MIGRATABLE scope: $(head -c 300 "$W/migrate-stderr.txt" | esc)"
fi
fi
# 3. Repository state: only the expected files may be modified, nothing else may appear or change mode.
gitr diff --name-only | LC_ALL=C sort > "$W/changed.txt"
gitr ls-files --others | LC_ALL=C sort > "$W/untracked.txt"
gitr diff --name-only --diff-filter=DT > "$W/deleted-or-typechanged.txt"
gitr diff --summary | grep 'mode change' > "$W/mode-changes.txt" || true
find repo -path repo/.git -prune -o -type l -print > "$W/symlinks.txt"
awk -v p="$PREFIX" 'index($0, p) != 1' "$W/changed.txt" > "$W/outside-scope.txt"
if [ "$STATUS" = REMEDIATED ]; then
[ -s "$W/untracked.txt" ] && fail_with UNTRACKED_FILES "Unexpected untracked files: $(head -c 300 "$W/untracked.txt" | esc)"
[ -s "$W/deleted-or-typechanged.txt" ] && fail_with DELETED_OR_TYPECHANGED "Files were deleted or changed type: $(head -c 300 "$W/deleted-or-typechanged.txt" | esc)"
[ -s "$W/mode-changes.txt" ] && fail_with MODE_CHANGE "File modes changed: $(head -c 300 "$W/mode-changes.txt" | esc)"
[ -s "$W/symlinks.txt" ] && fail_with SYMLINK_CREATED "Symbolic links appeared: $(head -c 300 "$W/symlinks.txt" | esc)"
[ -s "$W/outside-scope.txt" ] && fail_with OUTSIDE_SCOPE "Files outside flows_path changed: $(head -c 300 "$W/outside-scope.txt" | esc)"
[ -s "$W/changed.txt" ] || fail_with NO_CHANGES "kestra-migrate -o changed no file."
cmp -s "$W/expected.txt" "$W/changed.txt" || fail_with CHANGESET_MISMATCH "Changed files differ from the files the scan marked as rewritten."
fi
# 4. Post-check with the same tool and parser: the scope must now be CLEAN.
rc_post=0
./bin/kestra-migrate --check --summary -- "$TARGET" > "$W/post-raw.txt" 2> "$W/post-stderr.txt" || rc_post=$?
strip < "$W/post-raw.txt" > "$W/post-report.txt"
awk -f parse.awk "$W/post-report.txt" > "$W/post.line" 2> /dev/null
POST_STATUS=$(sed -n 's/.*"status":"\([A-Z_]*\)".*/\1/p' "$W/post.line")
POST_TOTAL=$(sed -n 's/.*"total":\([0-9]*\).*/\1/p' "$W/post.line")
PRE_TOTAL=$(sed -n 's/.*"total":\([0-9]*\).*/\1/p' "$W/pre.line")
if [ "$POST_STATUS" = TOOL_ERROR ]; then fail_with POST_CHECK_TOOL_ERROR "The post-remediation scan could not be trusted."
elif [ "$POST_STATUS" != CLEAN ]; then fail_with POST_CHECK_NOT_CLEAN "The post-remediation scan is $POST_STATUS, not CLEAN."
elif [ "$POST_TOTAL" != "$PRE_TOTAL" ]; then fail_with POST_CHECK_NOT_CLEAN "The post-remediation scan covers $POST_TOTAL flows instead of $PRE_TOTAL."; fi
# 5. Patch: binary-safe diff against the source commit, re-applied to a pristine copy to prove it.
PATCH_SHA=""; PATCH_BYTES=0; RESULT_TREE=""
if [ "$STATUS" = REMEDIATED ]; then
gitr diff --binary --full-index --no-ext-diff --no-textconv --no-color HEAD -- > remediation.patch
PATCH_SHA=$(sha256sum remediation.patch | cut -d' ' -f1); PATCH_BYTES=$(wc -c < remediation.patch)
GIT_INDEX_FILE="$PWD/$W/result.index" gitr read-tree HEAD
GIT_INDEX_FILE="$PWD/$W/result.index" gitr add -A
RESULT_TREE=$(GIT_INDEX_FILE="$PWD/$W/result.index" gitr write-tree)
rm -rf /tmp/pristine
git -c safe.directory='*' -c init.defaultBranch=main -c advice.defaultBranchName=false clone -q --no-checkout --no-hardlinks repo /tmp/pristine
git -C /tmp/pristine -c advice.detachedHead=false checkout -q --detach "$REQUESTED_COMMIT"
if git -C /tmp/pristine apply --check --binary "$PWD/remediation.patch" && git -C /tmp/pristine apply --binary "$PWD/remediation.patch"; then
git -C /tmp/pristine add -A
[ "$(git -C /tmp/pristine write-tree)" = "$RESULT_TREE" ] || fail_with PATCH_VERIFY_FAILED "Applying the patch to the source commit does not reproduce the remediated tree."
else
fail_with PATCH_VERIFY_FAILED "The patch does not apply cleanly to the source commit."
fi
if [ "$STATUS" != REMEDIATED ]; then rm -f remediation.patch; PATCH_SHA=""; PATCH_BYTES=0; RESULT_TREE=""; fi
fi
printf '{"contract":"kestra-migrate-remediation/1.0","status":"%s","reasonCode":"%s","reason":"%s","sourceRepo":"%s","sourceCommit":"%s","flowsPath":"%s","migratorVersion":"2.6.1","preCheckExitCode":%s,"migrateExitCode":%s,"postCheckExitCode":%s,"expectedFiles":%s,"changedFiles":%s,"untrackedFiles":%s,"changedCount":%s,"postCheckStatus":"%s","resultTree":"%s","patchSha256":"%s","patchBytes":%s}\n' \
"$STATUS" "$CODE" "$(printf '%s' "$REASON" | esc)" "$REPO_URL" "$REQUESTED_COMMIT" "$FLOWS_PATH" "$rc_pre" "$rc_mig" "$rc_post" \
"$(jlist "$W/expected.txt")" "$(jlist "$W/changed.txt")" "$(jlist "$W/untracked.txt")" "$(wc -l < "$W/changed.txt")" \
"$POST_STATUS" "$RESULT_TREE" "$PATCH_SHA" "$PATCH_BYTES" > remediation-result.json
echo "::{\"outputs\":{\"remediation\":$(cat remediation-result.json)}}::"
echo "Remediation $STATUS${CODE:+ ($CODE)}; changed $(wc -l < "$W/changed.txt") files; patch ${PATCH_SHA:-none}"
- id: remediation_guard
type: io.kestra.plugin.core.flow.If
description: Nothing is delivered unless every remediation check passed.
condition: "{{ outputs.remediate.vars.remediation.status != 'REMEDIATED' }}"
then:
- id: remediation_failed
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Remediation FAILED: {{
outputs.remediate.vars.remediation.reasonCode }}: {{
outputs.remediate.vars.remediation.reason }}"
- id: delivery_workdir
type: io.kestra.plugin.core.flow.WorkingDirectory
description: Prepares the delivery and pushes the remediation branch.
tasks:
- id: clone_c
type: io.kestra.plugin.git.Clone
description: Fresh clone of the same commit, used to rebuild the approved patch
as a commit.
url: "{{ inputs.repo_url }}"
commit: "{{ inputs.commit }}"
cloneSubmodules: false
directory: repo
- id: verify_c
type: io.kestra.plugin.scripts.shell.Commands
description: Runs on the Kestra worker before any container starts. Checks that
the checkout is exactly the requested commit, refuses
repositories containing symbolic links, and keeps flows_path
inside the clone.
taskRunner:
type: io.kestra.plugin.core.runner.Process
env:
REQUESTED_COMMIT: "{{ inputs.commit }}"
FLOWS_PATH: "{{ inputs.flows_path }}"
commands:
- |
set -eu
HEAD_SHA=$(cat repo/.git/HEAD)
if [ "$HEAD_SHA" != "$REQUESTED_COMMIT" ]; then echo "Checked-out HEAD '$HEAD_SHA' does not match the requested commit $REQUESTED_COMMIT."; exit 1; fi
LINKS=$(find repo -type l -printf '%P -> %l\n')
if [ -n "$LINKS" ]; then printf 'The repository contains symbolic links, which this gate refuses to process:\n%s\n' "$LINKS"; exit 1; fi
ROOT=$(realpath -e repo)
TARGET=$(realpath -m -- "repo/$FLOWS_PATH")
case "$TARGET/" in "$ROOT"/*) ;; *) echo "flows_path resolves outside the repository."; exit 1 ;; esac
echo "Verified checkout $HEAD_SHA: no symbolic links, flows_path stays inside the repository."
echo "::{\"outputs\":{\"headSha\":\"$HEAD_SHA\"}}::"
- id: delivery_plan
type: io.kestra.plugin.scripts.shell.Commands
description: Read-only. Re-checks the approved patch, builds the deterministic
branch commit, and decides CREATE or REUSE for the branch,
issue and pull request from the remediation key. Any
conflicting branch, issue or pull request stops the
delivery.
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: buildpack-deps:trixie-scm@sha256:ed6425e2963332c914dd78c0768b4b0f86680ac5f3d630f39f91ec8aa82384a1
env:
REPO_URL: "{{ inputs.repo_url }}"
REQUESTED_COMMIT: "{{ inputs.commit }}"
FLOWS_PATH: "{{ inputs.flows_path }}"
GITHUB_TOKEN: "{{ secret('GITHUB_TOKEN') }}"
inputFiles:
remediation.patch: "{{ outputs.remediate.outputFiles['remediation.patch'] }}"
remediation-result.json: "{{
outputs.remediate.outputFiles['remediation-result.json']
}}"
plan.py: |
"""Delivery plan: verifies the approved remediation, builds the deterministic branch commit, and decides
CREATE / REUSE for branch, issue and PR from current GitHub state. Performs no GitHub mutation."""
import hashlib, json, os, re, subprocess, sys, time, urllib.error, urllib.request
VERSION = '2.6.1'
CONTRACT = 'kestra-migrate-remediation/1.0'
REPO_URL, COMMIT, FLOWS = os.environ['REPO_URL'], os.environ['REQUESTED_COMMIT'], os.environ['FLOWS_PATH']
GIT = ['git', '-C', 'repo', '-c', 'safe.directory=*', '-c', 'core.quotePath=false', '-c', 'core.fsmonitor=false', '-c', 'core.hooksPath=/dev/null']
plan = {'status': 'OK', 'reasonCode': '', 'reason': ''}
def emit():
with open('plan.json', 'w') as f:
json.dump(plan, f)
print('::' + json.dumps({'outputs': {'plan': plan}}) + '::')
def stop(code, reason):
plan.update(status='REFUSED', reasonCode=code, reason=reason)
print(f'Delivery refused: {code}: {reason}')
emit()
sys.exit(0)
def git(*args, env=None, check=True):
# Read-only network commands are retried; everything else runs once.
attempts = 3 if args[0] in ('ls-remote', 'fetch') else 1
for attempt in range(attempts):
r = subprocess.run(GIT + list(args), capture_output=True, text=True, env=env)
if r.returncode == 0 or attempt == attempts - 1:
break
time.sleep(3 * (attempt + 1))
if check and r.returncode:
stop('GIT_ERROR', f'git {args[0]} failed: {r.stderr.strip()[:300]}')
return r.stdout.strip()
def api(path):
req = urllib.request.Request('https://api.github.com' + path, headers={
'Authorization': 'Bearer ' + os.environ['GITHUB_TOKEN'], 'Accept': 'application/vnd.github+json',
'X-GitHub-Api-Version': '2022-11-28', 'User-Agent': 'kestra-migration-gate'})
for attempt in range(3): # GETs are idempotent; transient network errors and 5xx are retried
try:
with urllib.request.urlopen(req, timeout=30) as r:
return json.loads(r.read())
except urllib.error.HTTPError as e:
if e.code < 500 or attempt == 2:
stop('GITHUB_API_ERROR', f'GET {path.split("?")[0]} returned HTTP {e.code}')
except (urllib.error.URLError, OSError) as e:
if attempt == 2:
stop('GITHUB_API_ERROR', f'GET {path.split("?")[0]} failed: {type(e).__name__}')
time.sleep(3 * (attempt + 1))
# 1. The remediation being delivered must be the approved, validated one.
res = json.load(open('remediation-result.json'))
patch_sha = hashlib.sha256(open('remediation.patch', 'rb').read()).hexdigest()
if res.get('contract') != CONTRACT or res.get('status') != 'REMEDIATED' or res.get('postCheckStatus') != 'CLEAN':
stop('REMEDIATION_NOT_DELIVERABLE', f"remediation is {res.get('status')} with post-check {res.get('postCheckStatus')}")
if res.get('sourceRepo') != REPO_URL or res.get('flowsPath') != FLOWS:
stop('INPUT_MISMATCH', 'the remediation result does not belong to this repository and flows_path')
head = open('repo/.git/HEAD').read().strip()
if res.get('sourceCommit') != COMMIT or head != COMMIT:
stop('SOURCE_COMMIT_MISMATCH', f"remediation source {res.get('sourceCommit')}, requested {COMMIT}, checked out {head}")
if not res.get('patchSha256') or patch_sha != res['patchSha256']:
stop('PATCH_SHA_MISMATCH', f"stored patch sha256 {patch_sha} differs from the validated {res.get('patchSha256')}")
# 2. Remediation identity and the names derived from it.
m = re.fullmatch(r'https://github\.com/([A-Za-z0-9-]+)/([A-Za-z0-9._-]+?)(?:\.git)?', REPO_URL)
owner, name = m.group(1), m.group(2)
canonical = (f'{CONTRACT}\nrepo=github.com/{owner.lower()}/{name.lower()}\ncommit={COMMIT}\n'
f'flowsPath={FLOWS}\nmigrator=kestra-migrate {VERSION}\n')
key = hashlib.sha256(canonical.encode()).hexdigest()
branch = f'kestra-migration/v{VERSION}-{COMMIT[:12]}-{key[:12]}'
marker = f'<!-- kestra-remediation-key: {key} -->'
plan.update(remediationKey=key, branch=branch, marker=marker, repository=f'{owner}/{name}', patchSha256=patch_sha,
changedFiles=res['changedFiles'])
# 3. Rebuild the validated tree from the source commit + stored patch; nothing else may change.
git('apply', '--check', '--binary', '../remediation.patch')
git('apply', '--binary', '../remediation.patch')
changed = sorted(git('diff', '--name-only').splitlines())
if changed != sorted(res['changedFiles']):
stop('CHANGESET_MISMATCH', f'applying the patch changed {changed}, expected {res["changedFiles"]}')
if git('ls-files', '--others'):
stop('UNTRACKED_FILES', 'untracked files are present after applying the patch')
if git('diff', '--summary').count('mode change'):
stop('MODE_CHANGE', 'the patch changes file modes')
idx = dict(os.environ, GIT_INDEX_FILE=os.path.abspath('delivery.index'))
git('read-tree', 'HEAD', env=idx)
git('add', '-A', env=idx)
tree = git('write-tree', env=idx)
if tree != res['resultTree']:
stop('TREE_MISMATCH', f'rebuilt tree {tree} differs from the validated {res["resultTree"]}')
# 4. Deterministic commit: same remediation, same commit SHA.
when = git('show', '-s', '--format=%cI', COMMIT)
ident = dict(os.environ, GIT_AUTHOR_NAME='Kestra Migration Gate', GIT_AUTHOR_EMAIL='kestra-migration-gate@users.noreply.github.com',
GIT_COMMITTER_NAME='Kestra Migration Gate', GIT_COMMITTER_EMAIL='kestra-migration-gate@users.noreply.github.com',
GIT_AUTHOR_DATE=when, GIT_COMMITTER_DATE=when)
message = (f'Migrate {len(changed)} Kestra flow(s) to 2.0 with kestra-migrate {VERSION}\n\n'
f'Source commit: {COMMIT}\nScope: {FLOWS}\nPatch SHA256: {patch_sha}\nRemediation key: {key}\n')
with open('commit-message.txt', 'w') as f:
f.write(message)
expected_head = git('commit-tree', tree, '-p', COMMIT, '-F', '../commit-message.txt', env=ident)
plan.update(expectedHead=expected_head, resultTree=tree)
# 5. The PR base is the repository's real default branch, and the source commit must be on it.
symref = git('ls-remote', '--symref', 'origin', 'HEAD')
mm = re.search(r'^ref: refs/heads/(\S+)\tHEAD$', symref, re.M)
if not mm:
stop('DEFAULT_BRANCH_UNKNOWN', 'could not resolve the remote default branch')
default = mm.group(1)
if branch == default or not branch.startswith('kestra-migration/'):
stop('UNSAFE_BRANCH', f'refusing to target branch {branch}')
git('fetch', '-q', 'origin', f'+refs/heads/{default}:refs/remotes/kmg/default')
if subprocess.run(GIT + ['merge-base', '--is-ancestor', COMMIT, 'refs/remotes/kmg/default']).returncode:
stop('SOURCE_NOT_ON_DEFAULT_BRANCH', f'{COMMIT} is not an ancestor of {default}')
plan['defaultBranch'] = default
# 6. Existing branch: reuse only if it is exactly this remediation on top of the source commit.
remote = git('ls-remote', 'origin', f'refs/heads/{branch}')
if remote:
rsha = remote.split()[0]
git('fetch', '-q', 'origin', f'+refs/heads/{branch}:refs/remotes/kmg/branch')
parents = git('rev-list', '--parents', '-n', '1', rsha).split()[1:]
rtree = git('rev-parse', rsha + '^' + 'tree'.join('{}'))
if rsha == expected_head or (parents == [COMMIT] and rtree == tree):
plan.update(branchAction='REUSE', branchHead=rsha)
else:
stop('BRANCH_CONFLICT', f'branch {branch} exists at {rsha} with unexpected content')
else:
plan.update(branchAction='CREATE', branchHead=expected_head)
# 7. Existing issue / PR, found by the hidden marker through strongly consistent REST list endpoints.
login = api('/user')['login']
items = api(f'/repos/{owner}/{name}/issues?state=all&creator={login}&per_page=100')
if len(items) >= 100:
stop('TOO_MANY_ITEMS', 'more than 99 issues/PRs by this account; refusing to guess')
marked = [i for i in items if marker in (i.get('body') or '')]
issues = [i for i in marked if 'pull_request' not in i]
if len(issues) > 1:
stop('ISSUE_CONFLICT', f'{len(issues)} issues carry this remediation key')
for p in [i for i in marked if 'pull_request' in i]:
pd = api(f'/repos/{owner}/{name}/pulls/{p["number"]}')
if pd['head']['ref'] != branch or (pd['head'].get('repo') or {}).get('full_name') != f'{owner}/{name}':
stop('PR_CONFLICT', f'PR #{p["number"]} carries this remediation key but comes from {pd["head"]["ref"]}')
by_head = api(f'/repos/{owner}/{name}/pulls?state=all&head={owner}:{branch}&per_page=100')
for p in by_head:
if marker not in (p.get('body') or ''):
stop('PR_CONFLICT', f'PR #{p["number"]} from {branch} does not carry this remediation key')
if p['base']['ref'] != default:
stop('PR_CONFLICT', f'PR #{p["number"]} targets {p["base"]["ref"]}, not {default}')
if len(by_head) > 1:
stop('PR_CONFLICT', f'{len(by_head)} PRs exist from {branch}')
plan.update(issueAction='REUSE' if issues else 'CREATE', issueNumber=issues[0]['number'] if issues else 0,
prAction='REUSE' if by_head else 'CREATE', prNumber=by_head[0]['number'] if by_head else 0)
print(f'Delivery plan: branch {plan["branchAction"]} {branch}, issue {plan["issueAction"]}, PR {plan["prAction"]}, base {default}')
emit()
commands:
- |
set -eu
[ "$(sha256sum plan.py | cut -d' ' -f1)" = 95330c056f02623e94e6da0c99983668de9880ef74edc3f31cf0f141e77d6441 ] || { echo "plan.py integrity check failed."; exit 1; }
export GIT_CONFIG_NOSYSTEM=1 GIT_CONFIG_GLOBAL=/dev/null GIT_TERMINAL_PROMPT=0 HOME=/nonexistent
python3 plan.py
- id: push_branch
type: io.kestra.plugin.scripts.shell.Commands
description: Create-only push of the remediation branch. The token reaches git
through a temporary askpass helper, never through a URL,
argument or git config.
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: buildpack-deps:trixie-scm@sha256:ed6425e2963332c914dd78c0768b4b0f86680ac5f3d630f39f91ec8aa82384a1
env:
GITHUB_TOKEN: "{{ secret('GITHUB_TOKEN') }}"
commands:
- |
set -eu
export GIT_CONFIG_NOSYSTEM=1 GIT_CONFIG_GLOBAL=/dev/null GIT_TERMINAL_PROMPT=0 HOME=/nonexistent
field() { python3 -c 'import json,sys; print(json.load(open("plan.json"))[sys.argv[1]])' "$1"; }
if [ "$(field status)" != OK ]; then echo "Plan is $(field status) ($(field reasonCode)); refusing to push."; exit 1; fi
BRANCH=$(field branch); HEAD_SHA=$(field expectedHead); DEFAULT=$(field defaultBranch)
echo "$BRANCH" | grep -Eq '^kestra-migration/v[0-9.]+-[0-9a-f]{12}-[0-9a-f]{12}$' || { echo "Unexpected branch name."; exit 1; }
[ "$BRANCH" != "$DEFAULT" ] || { echo "Refusing to push to the default branch."; exit 1; }
if [ "$(field branchAction)" = REUSE ]; then echo "Branch $BRANCH already holds this remediation; nothing to push."; exit 0; fi
printf '#!/bin/sh\ncase "$1" in Username*) echo x-access-token ;; *) printf "%%s\\n" "$GITHUB_TOKEN" ;; esac\n' > askpass.sh
chmod 700 askpass.sh
gitp() { git -C repo -c 'safe.directory=*' -c core.hooksPath=/dev/null -c credential.helper= "$@"; }
export GIT_ASKPASS="$PWD/askpass.sh"
# Create-only push. After a failed attempt the remote is re-read: if it already holds our
# deterministic commit, the push landed and only the response was lost.
for attempt in 1 2 3; do
if gitp push --porcelain --force-with-lease="refs/heads/$BRANCH:" origin "$HEAD_SHA:refs/heads/$BRANCH"; then break; fi
[ "$(gitp ls-remote origin "refs/heads/$BRANCH" | cut -f1)" = "$HEAD_SHA" ] && { echo "Branch already at $HEAD_SHA after attempt $attempt."; break; }
[ "$attempt" = 3 ] && { unset GIT_ASKPASS; rm -f askpass.sh; echo "Push failed after 3 attempts."; exit 1; }
sleep $((attempt * 3))
done
unset GIT_ASKPASS; rm -f askpass.sh
REMOTE=$(gitp ls-remote origin "refs/heads/$BRANCH" | cut -f1)
[ "$REMOTE" = "$HEAD_SHA" ] || { echo "Remote branch is at '$REMOTE', expected $HEAD_SHA."; exit 1; }
echo "Pushed $BRANCH at $HEAD_SHA"
- id: delivery_guard
type: io.kestra.plugin.core.flow.If
description: Stops before any issue or pull request is created if the delivery
plan was refused.
condition: "{{ outputs.delivery_plan.vars.plan.status != 'OK' }}"
then:
- id: delivery_refused
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Delivery refused: {{ outputs.delivery_plan.vars.plan.reasonCode
}}: {{ outputs.delivery_plan.vars.plan.reason }}"
- id: issue_gate
type: io.kestra.plugin.core.flow.If
description: Creates the issue only if no issue carries this remediation key
yet.
condition: "{{ outputs.delivery_plan.vars.plan.issueAction == 'CREATE' }}"
then:
- id: issue_create
type: io.kestra.plugin.github.issues.Create
oauthToken: "{{ secret('GITHUB_TOKEN') }}"
repository: "{{ outputs.delivery_plan.vars.plan.repository }}"
title: "Kestra 2.0 migration: {{ outputs.remediate.vars.remediation.changedCount
}} flow(s) in {{ inputs.flows_path }} can be auto-migrated"
body: |
{{ outputs.delivery_plan.vars.plan.marker }}
The Kestra 2.0 Migration Readiness Gate found flows that kestra-migrate {{ outputs.remediate.vars.remediation.migratorVersion }} rewrites automatically, and a reviewer approved the rewrite.
| | |
|---|---|
| Repository | {{ inputs.repo_url }} |
| Source commit | `{{ inputs.commit }}` |
| Scope (`flows_path`) | `{{ inputs.flows_path }}` |
| Classification | AUTO_MIGRATABLE |
| Files rewritten | {{ outputs.remediate.vars.remediation.changedCount }} |
| Validation | post-rewrite `kestra-migrate --check` is CLEAN; changed files match the scan exactly |
| Remediation key | `{{ outputs.delivery_plan.vars.plan.remediationKey }}` |
| Approval | Kestra execution `{{ execution.id }}`; reviewer name entered at approval: {{ outputs.approve.onResume.reviewer }} (not an authenticated identity) |
The pull request from branch `{{ outputs.delivery_plan.vars.plan.branch }}` closes this issue.
- id: pr_gate
type: io.kestra.plugin.core.flow.If
description: Creates the pull request only if none exists from the remediation
branch yet.
condition: "{{ outputs.delivery_plan.vars.plan.prAction == 'CREATE' }}"
then:
- id: pr_create
type: io.kestra.plugin.github.pulls.Create
oauthToken: "{{ secret('GITHUB_TOKEN') }}"
repository: "{{ outputs.delivery_plan.vars.plan.repository }}"
sourceBranch: "{{ outputs.delivery_plan.vars.plan.branch }}"
targetBranch: "{{ outputs.delivery_plan.vars.plan.defaultBranch }}"
title: "Migrate {{ outputs.remediate.vars.remediation.changedCount }} Kestra
flow(s) in {{ inputs.flows_path }} to 2.0 (kestra-migrate {{
outputs.remediate.vars.remediation.migratorVersion }})"
body: |
{{ outputs.delivery_plan.vars.plan.marker }}
Closes #{{ outputs.issue_create.issueNumber ?? outputs.delivery_plan.vars.plan.issueNumber }}
Rewrites flows under `{{ inputs.flows_path }}` with `kestra-migrate {{ outputs.remediate.vars.remediation.migratorVersion }} -o`, applied to source commit `{{ inputs.commit }}`.
**Files changed**
{% for f in outputs.remediate.vars.remediation.changedFiles %}
- `{{ f }}`
{% endfor %}
**Validation performed**
- `kestra-migrate --check` on the rewritten scope: CLEAN
- changed files equal the files the scan marked as rewritten; no other file, mode or symlink changed
- the patch re-applies to the source commit and reproduces tree `{{ outputs.remediate.vars.remediation.resultTree }}`
- no Kestra execution of the migrated flows was performed
| | |
|---|---|
| Source commit | `{{ inputs.commit }}` |
| Patch SHA256 | `{{ outputs.remediate.vars.remediation.patchSha256 }}` |
| Remediation key | `{{ outputs.delivery_plan.vars.plan.remediationKey }}` |
| Approval | Kestra execution `{{ execution.id }}`; reviewer name entered at approval: {{ outputs.approve.onResume.reviewer }} (not an authenticated identity) |
- id: delivery_verify
type: io.kestra.plugin.scripts.shell.Commands
description: Reads GitHub back and checks exactly one branch, issue and pull
request for this key, the branch content against the validated
patch, and that no credential appears in what was delivered.
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: buildpack-deps:trixie-scm@sha256:ed6425e2963332c914dd78c0768b4b0f86680ac5f3d630f39f91ec8aa82384a1
env:
GITHUB_TOKEN: "{{ secret('GITHUB_TOKEN') }}"
PLAN_JSON: "{{ outputs.delivery_plan.vars.plan | toJson }}"
REMEDIATION_JSON: "{{ outputs.remediate.vars.remediation | toJson }}"
inputFiles:
verify.py: |
"""Reads GitHub back after delivery and emits the kestra-migrate-github-delivery/1.0 result. Read-only."""
import json, os, sys, time, urllib.error, urllib.request
plan = json.loads(os.environ['PLAN_JSON'])
rem = json.loads(os.environ['REMEDIATION_JSON'])
token = os.environ['GITHUB_TOKEN']
repo, branch, key, marker = plan['repository'], plan['branch'], plan['remediationKey'], plan['marker']
problems = []
def api(path):
req = urllib.request.Request('https://api.github.com' + path, headers={
'Authorization': 'Bearer ' + token, 'Accept': 'application/vnd.github+json',
'X-GitHub-Api-Version': '2022-11-28', 'User-Agent': 'kestra-migration-gate'})
for attempt in range(3): # read-only: transient network errors and 5xx are retried
try:
with urllib.request.urlopen(req, timeout=30) as r:
return json.loads(r.read())
except urllib.error.HTTPError as e:
if e.code < 500 or attempt == 2:
problems.append(f'GET {path.split("?")[0]} returned HTTP {e.code}')
return None
except (urllib.error.URLError, OSError) as e:
if attempt == 2:
problems.append(f'GET {path.split("?")[0]} failed: {type(e).__name__}')
return None
time.sleep(3 * (attempt + 1))
br = api(f'/repos/{repo}/branches/{branch}') or {}
head = (br.get('commit') or {}).get('sha')
if head != plan['branchHead']:
problems.append(f'branch head is {head}, expected {plan["branchHead"]}')
cmp = api(f'/repos/{repo}/compare/{rem["sourceCommit"]}...{head}') or {}
files = sorted(f['filename'] for f in cmp.get('files', []))
if cmp.get('ahead_by') != 1 or cmp.get('behind_by') != 0 or files != sorted(rem['changedFiles']):
problems.append(f'branch is {cmp.get("ahead_by")} ahead/{cmp.get("behind_by")} behind the source commit and changes {files}')
login = (api('/user') or {}).get('login')
items = api(f'/repos/{repo}/issues?state=all&creator={login}&per_page=100') or []
issues = [i for i in items if marker in (i.get('body') or '') and 'pull_request' not in i]
prs = api(f'/repos/{repo}/pulls?state=all&head={repo.split("/")[0]}:{branch}&per_page=100') or []
if len(issues) != 1:
problems.append(f'{len(issues)} issues carry the remediation key, expected 1')
if len(prs) != 1:
problems.append(f'{len(prs)} PRs come from {branch}, expected 1')
issue = issues[0] if len(issues) == 1 else {}
pr = prs[0] if len(prs) == 1 else {}
if pr:
body = pr.get('body') or ''
for needle, what in ((marker, 'remediation key'), (rem['sourceCommit'], 'source commit'), (rem['patchSha256'], 'patch SHA256')):
if needle not in body:
problems.append(f'PR body is missing the {what}')
if pr['base']['ref'] != plan['defaultBranch'] or pr['head']['ref'] != branch or pr['head']['sha'] != head:
problems.append('PR base, head branch or head commit is not the expected one')
if issue and f'#{issue["number"]}' not in body:
problems.append('PR body does not reference the issue')
texts = [issue.get('body') or '', issue.get('title') or '', pr.get('body') or '', pr.get('title') or '']
commit = api(f'/repos/{repo}/commits/{head}') if head else None
texts.append(((commit or {}).get('commit') or {}).get('message') or '')
if any(token in t for t in texts) or any(p in t for t in texts for p in ('github_pat_', 'ghp_', 'gho_', 'x-access-token')):
problems.append('a credential-like string appears in the delivered issue, PR or commit')
result = {
'contract': 'kestra-migrate-github-delivery/1.0', 'status': 'DELIVERED' if not problems else 'FAILED',
'reasonCode': '' if not problems else 'GITHUB_STATE_MISMATCH', 'reason': '; '.join(problems)[:600],
'remediationKey': key, 'sourceRepo': rem['sourceRepo'], 'sourceCommit': rem['sourceCommit'],
'branch': branch, 'branchHead': head, 'baseBranch': plan['defaultBranch'],
'issueNumber': issue.get('number', 0), 'issueUrl': issue.get('html_url', ''),
'prNumber': pr.get('number', 0), 'prUrl': pr.get('html_url', ''), 'prState': pr.get('state', ''),
'patchSha256': rem['patchSha256'], 'changedFiles': rem['changedFiles'],
'createdOrReused': {'branch': plan['branchAction'], 'issue': plan['issueAction'], 'pr': plan['prAction']},
}
print(f'Delivery {result["status"]}: branch {branch} @ {head}, issue #{result["issueNumber"]}, PR #{result["prNumber"]}')
print('::' + json.dumps({'outputs': {'delivery': result}}) + '::')
commands:
- |
set -eu
[ "$(sha256sum verify.py | cut -d' ' -f1)" = eb31755c18c74ada80f73438d2487ef830acf892717afe70e4067172ecfacc09 ] || { echo "verify.py integrity check failed."; exit 1; }
python3 verify.py
- id: delivery_result_guard
type: io.kestra.plugin.core.flow.If
condition: "{{ outputs.delivery_verify.vars.delivery.status != 'DELIVERED' }}"
then:
- id: delivery_failed
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Delivery verification failed: {{
outputs.delivery_verify.vars.delivery.reason }}"
- id: outcome_delivered
type: io.kestra.plugin.core.debug.Return
format: "{{ outputs.delivery_verify.vars.delivery | jq('{outcome:
\"REMEDIATION_DELIVERED\", scanStatus: \"AUTO_MIGRATABLE\"} +
.') | first | toJson }}"
defaults:
- id: unknown_status
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Unknown scan status '{{ outputs.scan.vars.status }}'."
outputs:
- id: readiness
type: STRING
description: Scan classification of the requested scope - CLEAN,
AUTO_MIGRATABLE, ADVISORY, BLOCKING or TOOL_ERROR.
value: "{{ outputs.scan.vars.status ?? 'NOT_SCANNED' }}"
- id: flow_counts
type: JSON
description: Number of scanned flows per classification.
value: "{{ outputs.scan.vars.counts is defined ? (outputs.scan.vars.counts |
jq('{total: (.CLEAN + .AUTO_MIGRATABLE + .ADVISORY + .BLOCKING)} + .') |
first | toJson) : '{}' }}"
- id: result
type: JSON
description: Final verdict. For a delivered remediation it includes the branch,
issue, pull request, patch SHA-256 and whether each was created or reused.
value: "{{ outputs.outcome_delivered.value ?? outputs.outcome_rejected.value ??
outputs.outcome_clean.value ?? outputs.outcome_advisory.value ??
outputs.outcome_blocking.value ?? '{}' }}"
triggers:
- id: readiness_webhook
type: io.kestra.plugin.core.trigger.Webhook
description: 'POST {"repo_url": "...", "commit": "...", "flows_path": "..."} to
start a readiness check, for example from CI after a push. The values go
through the same input validators as a manual run; an invalid request gets
HTTP 422 and no execution is created.'
key: "{{ secret('MIGRATION_GATE_WEBHOOK_KEY') }}"
inputs:
repo_url: "{{ trigger.body.repo_url }}"
commit: "{{ trigger.body.commit }}"
flows_path: "{{ trigger.body.flows_path ?? '.' }}"
Find out whether the Kestra flows in a GitHub repository are ready for Kestra 2.0 and, when the migration can be automated, let a human approve a safe, verified rewrite that is delivered as a pull request.
Kestra 2.0 removes or changes several 1.x constructs. kestra-migrate reports what has to change and can rewrite many flows automatically, but running a rewriting tool on a repository and pushing the result is a risky, manual operation. This blueprint turns it into a gated workflow: read-only scan → classification → human approval → remediation in a fresh clone → verified patch → idempotent GitHub delivery.
Scan (read-only). clone_a (io.kestra.plugin.git.Clone) checks out the exact commit. verify_a confirms the checkout is that commit, refuses repositories containing symbolic links, and keeps flows_path inside the clone. install_a downloads kestra-migrate 2.6.1 and verifies its SHA-256. scan runs kestra-migrate --check --summary and parses the output with a strict parser: any unexpected line, a missing summary, or totals that do not reconcile make the result TOOL_ERROR instead of a guess.
Classify. The policy Switch routes on the result:
| Classification | Meaning | What the gate does |
|---|---|---|
CLEAN |
every flow is already 2.0-compatible | reports NO_REMEDIATION_REQUIRED |
AUTO_MIGRATABLE |
kestra-migrate can rewrite every flow that needs it | asks for approval |
ADVISORY |
flows deploy on 2.0 but can break at run time | reports NOT_ELIGIBLE; a human must fix them |
BLOCKING |
Kestra 2.0 rejects at least one flow | reports NOT_ELIGIBLE; a human must fix them |
TOOL_ERROR |
the scan output cannot be trusted | fails the execution |
Approve. For AUTO_MIGRATABLE only, show_proposal logs the repository, commit, path, counts, warnings, the files that would be rewritten and the full proposed diff, then approve (io.kestra.plugin.core.flow.Pause) waits for a decision. decision defaults to REJECT; an unanswered request is cancelled after one day.
Remediate. On APPROVE, remediation_workdir clones the same commit again (the scanned checkout is never modified), re-runs the scan and requires the exact approved result, then runs kestra-migrate -o. remediate rejects the result if any file outside the expected set changed, if files were added, deleted, re-moded or turned into symbolic links, if anything outside flows_path changed, or if a post-migration scan is not CLEAN. The binary patch is re-applied to a pristine copy of the source commit and must reproduce the same tree.
Deliver. delivery_plan derives a remediation key (SHA-256 of repository, commit, path and migrator version), rebuilds the patch as a deterministic commit, and decides CREATE or REUSE for the branch kestra-migration/v2.6.1-<commit12>-<key12>, the issue and the pull request, which carry the hidden marker <!-- kestra-remediation-key: <key> -->. push_branch performs a create-only push. issue_create (io.kestra.plugin.github.issues.Create) and pr_create (io.kestra.plugin.github.pulls.Create) run only when nothing exists yet. A foreign branch or pull request with the same name or key stops the delivery instead of being overwritten.
Verify. delivery_verify reads GitHub back and requires exactly one branch, one issue and one pull request for the key, the branch to be one commit on top of the source commit with exactly the validated files, the pull request to target the default branch, and no credential in anything delivered.
https://github.com/<owner>/<repo> URLs and full 40-character commit SHAs are accepted; each clone is checked against the requested commit.flows_path is checked against the real path of the clone.APPROVE.github.com and api.github.com.plugin-git, plugin-github and plugin-scripts (included in the default Kestra image).GITHUB_TOKEN: GitHub token used, after approval only, to push the remediation branch, open the issue and the pull request, and read them back. Not used for CLEAN, ADVISORY, BLOCKING or rejected runs.MIGRATION_GATE_WEBHOOK_KEY: random key that forms the webhook URL. Required only if you use the webhook trigger.| Name | Type | Default | Description |
|---|---|---|---|
repo_url |
STRING | https://github.com/kestra-io/kestra2-flow-migration |
Public GitHub repository, https://github.com/<owner>/<repo>. No other host, scheme, port, credentials, query or fragment. |
commit |
STRING | 41e0fe57268711013841f633bd189cb04d83b198 |
Full 40-character lowercase commit SHA. Branch names and HEAD are rejected. |
flows_path |
STRING | input-flows/additional-test-cases/ion-read.yaml |
Folder or flow file inside the repository, or . for the whole repository. No absolute paths, .. or leading -. |
The defaults scan one example flow of the kestra-migrate repository at a pinned commit. It has a Kestra 2.0 advisory finding (an ION output that becomes binary), so a default run reports ADVISORY / NOT_ELIGIBLE and ends: it never pauses for approval, never uses GITHUB_TOKEN and never writes to GitHub. To remediate, run the flow with your own repository, commit and flows path.
GITHUB_TOKEN secret (and MIGRATION_GATE_WEBHOOK_KEY if you use the webhook).readiness and flow_counts outputs. If the result is AUTO_MIGRATABLE, review the proposal in the logs and resume the paused execution with decision: APPROVE and your name, or leave it on REJECT.result output.To check every push automatically, call the webhook from CI, for example from a GitHub Actions step:
curl -fsS -X POST -H 'Content-Type: application/json' \
-d "{\"repo_url\": \"https://github.com/${GITHUB_REPOSITORY}\", \"commit\": \"${GITHUB_SHA}\", \"flows_path\": \"flows\"}" \
"https://<your-kestra-host>/api/v1/<tenant>/executions/webhook/company.team/kestra-2-migration-readiness-gate/${MIGRATION_GATE_WEBHOOK_KEY}"
The webhook values go through the same input validators as a manual run: an invalid request is answered with HTTP 422 and no execution is created. flows_path defaults to . when omitted.
readiness: CLEAN, AUTO_MIGRATABLE, ADVISORY, BLOCKING or TOOL_ERROR.flow_counts: number of flows per classification, plus the total.result: NO_REMEDIATION_REQUIRED, NOT_ELIGIBLE, REJECTED or REMEDIATION_DELIVERED. A delivered result includes the branch and head commit, the issue and pull request numbers and URLs, the patch SHA-256, the changed files, and whether the branch, issue and pull request were created or reused.{{ outputs.scan.outputFiles['scan-report.txt'] }}: the full kestra-migrate report; {{ outputs.remediate.outputFiles['remediation.patch'] }}: the approved patch.