From 69355ad9bd4879c197dcf157161cd5ffe7239e09 Mon Sep 17 00:00:00 2001 From: soustruh Date: Thu, 24 Sep 2026 13:31:40 +0200 Subject: [PATCH] feat(sync): clone recreates the target's storage buckets by default sync clone copies configs, not storage, so a cloned config's input/output mappings point at buckets a fresh target does not have. A clone is a complete clone, so it now recreates the missing buckets by default. Pass --no-create-buckets to skip it. Clone reads the storage/buckets.json pull export and creates the missing buckets on the backend the export recorded. It maps each id through --bucket-map, so a created bucket matches the rewritten config refs. Pull now records each bucket's backend and the source of a linked bucket. Clone links a linked bucket to the same source, under the same id. The source project's sharing settings decide if the target can link it. An export from an older pull does not record linked buckets, so clone creates no bucket from it and records one error. Idempotent: clone skips an existing bucket and collects a per-bucket API failure, a refused link included, in bucket_errors instead of aborting. Only the buckets are created. Their tables and their data are not in the export. --- CLAUDE.md | 4 +- plugins/kbagent/agents/keboola-expert.md | 2 +- .../kbagent/references/commands-reference.md | 2 +- .../skills/kbagent/references/gotchas.md | 10 + .../kbagent/references/sync-workflow.md | 20 + src/keboola_agent_cli/commands/context.py | 6 +- src/keboola_agent_cli/commands/sync.py | 26 ++ src/keboola_agent_cli/result_models.py | 13 + src/keboola_agent_cli/services/_sync_clone.py | 19 +- .../services/_sync_storage.py | 154 +++++++- tests/test_result_models.py | 20 + tests/test_sync_clone.py | 353 +++++++++++++++++- tests/test_sync_storage_jobs.py | 23 ++ 13 files changed, 642 insertions(+), 10 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 0b8b09745..80bd1804f 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -964,8 +964,8 @@ kbagent sync push --project ALIAS [--all-projects] [--dry-run] [--force] [--allo # with SYNC_LEGACY_BOUNDARY telling you to pull first; genuine edits push normally. # 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] [--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. +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. 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/plugins/kbagent/agents/keboola-expert.md b/plugins/kbagent/agents/keboola-expert.md index 2ce012597..b9210fdf5 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; 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 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 | | 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 00ea497dd..b995b8326 100644 --- a/plugins/kbagent/skills/kbagent/references/commands-reference.md +++ b/plugins/kbagent/skills/kbagent/references/commands-reference.md @@ -361,7 +361,7 @@ Requires the project to be added with its **master ('owner') Storage API token** - `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 still deletes on push. 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] [--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`. +- `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 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`. - `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 f7052740e..b313619b8 100644 --- a/plugins/kbagent/skills/kbagent/references/gotchas.md +++ b/plugins/kbagent/skills/kbagent/references/gotchas.md @@ -939,6 +939,16 @@ A `keboola.data-apps` config's runtime type (`python-js` / `streamlit` / ...) li The DS `/apps` list also returns sandbox and workspace records. Each carries a parent component's id and a backend `type` such as `snowflake`. So kbagent builds the type map from `componentId == keboola.data-apps` records only. +## `sync clone` recreates the reference's storage buckets + +`sync clone` copies component configs, not storage. The pulled `storage/` tree is a read-only snapshot, so a cloned config's input/output mappings point at buckets a fresh target project does not have. Clone *(since vNEXT)* closes that gap for the buckets **by default**: it reads the `storage/buckets.json` pull export, maps each bucket id through `--bucket-map` (so a created bucket matches what the config refs were rewritten to), and creates the ones the target is missing. Pass `--no-create-buckets` to skip it -- a clone is a complete clone by default, so this is opt-out, never opt-in. + +Idempotent by design -- an existing bucket is skipped, and a per-bucket API failure is collected in the result's `bucket_errors` rather than aborting the clone. It runs even on an idempotent re-run (existing `--target-dir`), so a re-clone fills in any bucket the target is still missing. Bucket creation happens before the config push, and buckets are created at production level (no branch scoping). Each bucket is created on the backend the export recorded. If the target cannot list its buckets, clone records one `bucket_errors` entry, creates no bucket, and still pushes the configs. + +A linked (shared) bucket is **linked** in the target to the same source as in the reference, under the same id (or its `--bucket-map` id), with the stage taken from that id. It is not created as an empty bucket: nothing in the project writes into a linked bucket, so an empty one would stay empty. The source project's sharing settings decide whether the target may link it. A refused link lands in `bucket_errors` and the clone continues. Each link is listed in `linked_buckets` (`bucket_id`, `source_bucket_id`, `source_project_id`) and printed with its source project, so a link to an unexpected source is visible. A pull by an older version does not record the link, so clone cannot tell its linked buckets apart. Creating them empty would be permanent, because every later re-clone skips an existing bucket. So clone creates no bucket from such an export and records one `bucket_errors` entry. Re-pull the reference, or pass `--no-create-buckets`. + +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. + ## `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 cf86535ab..ecbe6aaee 100644 --- a/plugins/kbagent/skills/kbagent/references/sync-workflow.md +++ b/plugins/kbagent/skills/kbagent/references/sync-workflow.md @@ -550,6 +550,26 @@ nested mapping, list, or empty (`null`) value is rejected with `CONFIG_ERROR` (exit 5) naming the offending key and its actual type, instead of being silently stringified into a bogus ID. +**Storage buckets (since vNEXT):** clone copies configs, not storage — a cloned +config's input/output mappings point at buckets a fresh target does not have. By +default clone reads the `storage/buckets.json` pull export and creates the missing +buckets in the target. Pass `--no-create-buckets` to skip it — a clone is a +complete clone by default, so bucket creation is opt-out, not opt-in. It is +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 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: nothing in the project writes into a linked bucket, so 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 is collected in `bucket_errors`. Each link +is listed in `linked_buckets` with its source project and bucket id. 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 — their **tables and data are not** in +the export, so populate tables by bucket sharing, `storage upload-table`, or by +running the flows. Pull the source **with** its storage (`sync pull` without +`--no-storage`) so `buckets.json` exists in the reference tree. + **Why it just works on a fresh target:** the reference's config ids do not exist in the target project, so the push diff classifies every config as `added` and assigns new ULIDs. Because the push's `created_id_map` is keyed by the reference diff --git a/src/keboola_agent_cli/commands/context.py b/src/keboola_agent_cli/commands/context.py index f8250894e..1b83c96d0 100644 --- a/src/keboola_agent_cli/commands/context.py +++ b/src/keboola_agent_cli/commands/context.py @@ -1631,13 +1631,17 @@ tracked on another branch's tree are never planned as creates -- they ride along on the result envelope under `orphaned` instead (see sync diff). --dry-run agrees. - kbagent sync clone --source DIR --target ALIAS --target-dir DIR [--bucket-map FILE] [--variable-values FILE] [--instance-rename FILE] [--dry-run] [--branch ID] + 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] 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. 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 + storage/buckets.json BY DEFAULT (--no-create-buckets skips it; buckets only, + not tables/data). A linked (shared) bucket is linked to the same source as in + the reference (listed in linked_buckets; a refused link -> bucket_errors). Note: --dry-run still creates --target-dir on disk (copy + overrides + manifest) but does not push. diff --git a/src/keboola_agent_cli/commands/sync.py b/src/keboola_agent_cli/commands/sync.py index e60c32fc7..16543df08 100644 --- a/src/keboola_agent_cli/commands/sync.py +++ b/src/keboola_agent_cli/commands/sync.py @@ -1167,6 +1167,12 @@ def sync_clone( "--branch", help="Target dev branch id (defaults to the target project's production branch)", ), + create_buckets: bool = typer.Option( + True, + "--create-buckets/--no-create-buckets", + help="Create the reference tree's storage buckets in the target -- on by " + "default, --no-create-buckets skips it. Tables and their data are never copied.", + ), ) -> None: """Clone a reference project into a fresh target, parameterised by overrides. @@ -1187,6 +1193,7 @@ def sync_clone( overrides["variable_values"] = _load_override_file(variable_values) if instance_rename is not None: overrides["instance_rename"] = _load_override_file(instance_rename) + overrides["create_buckets"] = create_buckets result = service.clone_project( source=source, @@ -1215,6 +1222,23 @@ def sync_clone( _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", "") @@ -1237,11 +1261,13 @@ def _format_clone_result(formatter: Any, result: dict[str, Any]) -> None: 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')} " diff --git a/src/keboola_agent_cli/result_models.py b/src/keboola_agent_cli/result_models.py index f903249f9..54877e5c0 100644 --- a/src/keboola_agent_cli/result_models.py +++ b/src/keboola_agent_cli/result_models.py @@ -222,6 +222,19 @@ class CloneResult(_ApiResultModel): flow_task_remaps: int = Field( default=0, description="keboola.flow task configIds remapped reference->ULID." ) + buckets_created: int = Field( + default=0, description="Storage buckets created in the target (--create-buckets)." + ) + buckets_skipped: int = Field( + default=0, description="Buckets that already existed in the target and were skipped." + ) + bucket_errors: list[dict[str, Any]] = Field( + default_factory=list, description="Per-bucket create failures, if any." + ) + linked_buckets: list[dict[str, Any]] = Field( + default_factory=list, + description="Linked (shared) buckets linked in the target to the reference's source.", + ) push: SyncPushResult | None = Field( default=None, description="The underlying sync push result (None for dry_run)." ) diff --git a/src/keboola_agent_cli/services/_sync_clone.py b/src/keboola_agent_cli/services/_sync_clone.py index 5d544ee6e..c89ad28bc 100644 --- a/src/keboola_agent_cli/services/_sync_clone.py +++ b/src/keboola_agent_cli/services/_sync_clone.py @@ -23,6 +23,7 @@ repoint_manifest_project, ) from ..sync.manifest import load_manifest, save_manifest +from ._sync_storage import create_buckets_from_export from .base import find_default_branch_id if TYPE_CHECKING: @@ -77,8 +78,11 @@ def clone_project( the first clone). target_dir: Where to materialise the clone. overrides: Optional dict with ``bucket_map`` (old->new bucket id), - ``variable_values`` (var name->value), and ``instance_rename`` (old - path prefix->new path prefix) keys. + ``variable_values`` (var name->value), ``instance_rename`` (old + path prefix->new path prefix), and ``create_buckets`` (bool, + default True: also create the reference tree's storage buckets in + the target from ``storage/buckets.json`` -- tables and their data + are not copied) keys. dry_run: Apply overrides + report the diff without pushing. branch_override: Optional target branch id. @@ -98,6 +102,7 @@ def clone_project( bucket_map = overrides.get("bucket_map") or {} variable_values = overrides.get("variable_values") or {} instance_rename = overrides.get("instance_rename") or {} + create_buckets = bool(overrides.get("create_buckets", True)) source_dir = Path(source) target_path = Path(target_dir) @@ -195,6 +200,12 @@ def clone_project( "project, or remove the existing configs first." ) + bucket_result = None + if create_buckets: + bucket_client = service._client_factory(target_project.stack_url, target_project.token) + with bucket_client: + bucket_result = create_buckets_from_export(bucket_client, target_path, bucket_map) + push_result = service.push(target_alias, target_path, branch_override=branch_override) status = "no_changes" if push_result.get("status") == "no_changes" else "cloned" return { @@ -203,6 +214,10 @@ def clone_project( "target_dir": str(target_path), "created": push_result.get("created", 0), "flow_task_remaps": push_result.get("flow_task_remaps", 0), + "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", []), **override_counts, diff --git a/src/keboola_agent_cli/services/_sync_storage.py b/src/keboola_agent_cli/services/_sync_storage.py index ecb6b23d5..d9ede4ddf 100644 --- a/src/keboola_agent_cli/services/_sync_storage.py +++ b/src/keboola_agent_cli/services/_sync_storage.py @@ -1,10 +1,12 @@ -"""Pull-side storage-metadata, jobs, and sample writers (from sync_service.py). +"""Storage-metadata writers (``sync pull``) + bucket re-creation (``sync clone``). Free functions that materialise a project's storage metadata (buckets/tables), per-config job history, and table data samples to the filesystem during ``sync pull``. Only :func:`fetch_jobs_per_config` needs the ``SyncService`` (for its ``_resolve_max_workers`` helper); the rest are pure. ``_ensure_path_within`` (the storage-write path-traversal guard) moved here with its only callers. +:func:`create_buckets_from_export` reads the ``buckets.json`` written on pull +back, to recreate a clone's buckets in the target project. """ from __future__ import annotations @@ -15,6 +17,7 @@ import logging import threading from concurrent.futures import ThreadPoolExecutor, as_completed +from dataclasses import dataclass from pathlib import Path from typing import TYPE_CHECKING, Any @@ -26,11 +29,12 @@ STORAGE_DIR_NAME, STORAGE_SAMPLES_DIR_NAME, ) -from ..errors import ConfigError +from ..errors import ConfigError, KeboolaApiError from ..sync.manifest import ManifestConfiguration from ..sync.naming import sanitize_path_segment if TYPE_CHECKING: + from ..client import KeboolaClient from .sync_service import SyncService logger = logging.getLogger(__name__) @@ -60,6 +64,17 @@ def _ensure_path_within(base_dir: Path, target: Path, what: str) -> None: ) +def _linked_source(bucket: dict[str, Any]) -> dict[str, Any] | None: + """Return ``{"bucket_id", "project_id"}`` of a linked bucket's source, else ``None``.""" + source = bucket.get("sourceBucket") + if not source: + return None + return { + "bucket_id": source.get("id", ""), + "project_id": (source.get("project") or {}).get("id"), + } + + def write_storage_metadata( project_root: Path, buckets: list[dict[str, Any]], @@ -88,6 +103,8 @@ def write_storage_metadata( "name": b.get("name", ""), "stage": b.get("stage", ""), "description": b.get("description", ""), + "backend": b.get("backend", ""), + "source_bucket": _linked_source(b), "tables_count": b.get("tablesCount") or 0, "data_size_bytes": b.get("dataSizeBytes") or 0, "metadata": b.get("metadata", []), @@ -374,3 +391,136 @@ def mask_encrypted_columns(csv_data: str) -> str: writer.writerow(row) return output.getvalue() + + +_LEGACY_EXPORT_ERROR = ( + "storage/buckets.json comes from an older kbagent pull and does not record linked " + "buckets. No bucket was created. Pull the reference again, or use --no-create-buckets." +) + + +@dataclass +class BucketCreateResult: + """Outcome of re-creating a reference tree's buckets in a clone target. + + ``created`` / ``skipped`` hold the target bucket ids; ``errors`` collects a + ``{"bucket_id", "error"}`` record per bucket the API refused, so one bad + bucket never aborts the clone. ``linked`` holds a + ``{"bucket_id", "source_bucket_id", "source_project_id"}`` record per linked + (shared) bucket that was linked in the target to the reference's source. + """ + + created: list[str] + skipped: list[str] + errors: list[dict[str, str]] + linked: list[dict[str, Any]] + + +def _parse_bucket_id(bucket_id: str) -> tuple[str, str] | None: + """Split a ``.c-`` bucket id into ``(stage, name)``. + + Returns ``None`` for anything ``create-bucket`` cannot recreate: an id with + no ``in`` / ``out`` stage prefix (a system bucket) or an empty name. + """ + stage, _, remainder = bucket_id.partition(".") + if stage not in ("in", "out") or not remainder: + return None + name = remainder.removeprefix("c-") + return (stage, name) if name else None + + +def create_buckets_from_export( + storage_client: KeboolaClient, + project_root: Path, + bucket_map: dict[str, str], +) -> BucketCreateResult: + """Create a clone's reference buckets in the target project. + + ``sync clone`` copies configs, not storage: the pulled ``storage/`` tree is + a read-only snapshot, so a cloned config's input/output mappings point at + buckets that do not exist in a fresh target. This reads the + ``storage/buckets.json`` export, maps each bucket id through ``bucket_map`` + (so a created bucket matches what :func:`sync.clone.apply_bucket_map` + rewrote the configs to reference), and creates the ones that are missing. + + Idempotent: an existing bucket is skipped, and a per-bucket API failure is + collected in the result, never raised. Each bucket is created on the + backend the export recorded. A linked bucket is linked to the same source + as in the reference instead: nothing in the project writes into a linked + bucket, so 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 is collected like any other failure. An export from an older + pull (no ``source_bucket`` key) creates nothing and records one error + instead. Only the buckets are + created -- their tables and data are not in the export, so they are not + recreated. + """ + result = BucketCreateResult(created=[], skipped=[], errors=[], linked=[]) + buckets_file = project_root / STORAGE_DIR_NAME / STORAGE_BUCKETS_FILENAME + if not buckets_file.exists(): + return result + try: + records = json.loads(buckets_file.read_text(encoding="utf-8")) + except (OSError, ValueError): + logger.warning("Could not read %s; skipping bucket creation", buckets_file, exc_info=True) + return result + if not isinstance(records, list): + return result + # An older pull wrote no ``source_bucket`` key, so its linked buckets cannot + # be told apart. Creating them empty cannot be undone by a re-clone (the + # existing bucket is skipped), so create nothing. + if any(isinstance(r, dict) and "source_bucket" not in r for r in records): + result.errors.append({"bucket_id": "", "error": _LEGACY_EXPORT_ERROR}) + return result + + try: + existing = {str(b.get("id")) for b in storage_client.list_buckets()} + except KeboolaApiError as exc: + result.errors.append({"bucket_id": "", "error": f"Cannot list buckets: {exc.message}"}) + return result + for record in records: + if not isinstance(record, dict): + continue + source_id = str(record.get("id") or "") + target_id = bucket_map.get(source_id, source_id) + parsed = _parse_bucket_id(target_id) + if parsed is None: + continue + if target_id in existing: + result.skipped.append(target_id) + continue + stage, name = parsed + source = record.get("source_bucket") + link: dict[str, Any] | None = ( + { + "bucket_id": target_id, + "source_bucket_id": source.get("bucket_id", ""), + "source_project_id": source.get("project_id"), + } + if isinstance(source, dict) + else None + ) + try: + if link: + storage_client.link_bucket( + source_project_id=link["source_project_id"], + source_bucket_id=link["source_bucket_id"], + name=name, + stage=stage, + ) + else: + storage_client.create_bucket( + stage=stage, + name=name, + description=record.get("description") or "", + backend=record.get("backend") or None, + ) + except KeboolaApiError as exc: + result.errors.append({"bucket_id": target_id, "error": exc.message}) + continue + if link: + result.linked.append(link) + else: + result.created.append(target_id) + existing.add(target_id) + return result diff --git a/tests/test_result_models.py b/tests/test_result_models.py index 59c61b9bf..93b63a9ec 100644 --- a/tests/test_result_models.py +++ b/tests/test_result_models.py @@ -205,6 +205,26 @@ def test_dry_run_without_push(self) -> None: ) assert cr.status == "dry_run" and cr.push is None and cr.ok is True + def test_bucket_fields(self) -> None: + linked = {"bucket_id": "in.c-s", "source_bucket_id": "out.c-o", "source_project_id": 42} + cr = CloneResult.model_validate( + { + "status": "cloned", + "buckets_created": 2, + "buckets_skipped": 1, + "bucket_errors": [{"bucket_id": "in.c-bad", "error": "boom"}], + "linked_buckets": [linked], + } + ) + assert (cr.buckets_created, cr.buckets_skipped) == (2, 1) + assert cr.bucket_errors == [{"bucket_id": "in.c-bad", "error": "boom"}] + assert cr.linked_buckets == [linked] + + def test_bucket_fields_default_empty(self) -> None: + cr = CloneResult.model_validate({"status": "cloned"}) + assert (cr.buckets_created, cr.buckets_skipped) == (0, 0) + assert cr.bucket_errors == [] and cr.linked_buckets == [] + class TestBaseConfig: def test_all_models_allow_extra(self) -> None: diff --git a/tests/test_sync_clone.py b/tests/test_sync_clone.py index 90754e3f8..1f3308a5b 100644 --- a/tests/test_sync_clone.py +++ b/tests/test_sync_clone.py @@ -5,6 +5,7 @@ tree copy and manifest re-point -- no API client involved. """ +import json from pathlib import Path from typing import Any from unittest.mock import MagicMock @@ -13,9 +14,10 @@ import yaml from keboola_agent_cli.config_store import ConfigStore -from keboola_agent_cli.errors import ConfigError +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_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 ( _config_dir, @@ -986,6 +988,15 @@ def test_forwards_args_and_overrides(self, tmp_path: Path) -> None: "variable_overrides": 0, "renamed_instances": 0, "flow_task_remaps": 1, + "buckets_created": 1, + "buckets_skipped": 0, + "linked_buckets": [ + { + "bucket_id": "in.c-shared", + "source_bucket_id": "out.c-origin", + "source_project_id": 42, + } + ], "errors": [], } MockSync.return_value = svc @@ -1010,7 +1021,43 @@ def test_forwards_args_and_overrides(self, tmp_path: Path) -> None: call = svc.clone_project.call_args.kwargs assert call["target_alias"] == "target" assert call["overrides"]["bucket_map"] == {"in.c-ref": "in.c-prod"} + assert call["overrides"]["create_buckets"] is True assert "Cloned into target" in result.output + assert "Buckets: 1 created, 1 linked, 0 already present" in result.output + assert "Linked in.c-shared -> project 42 bucket out.c-origin" in result.output + + def test_no_create_buckets_forwards_false(self, tmp_path: Path) -> None: + from unittest.mock import patch + + from typer.testing import CliRunner + + from keboola_agent_cli.cli import app + + runner = CliRunner() + source = tmp_path / "golden" + _golden_source(source) + + with patch("keboola_agent_cli.cli.SyncService") as MockSync: + svc = MagicMock() + svc.clone_project.return_value = {"status": "cloned", "errors": []} + MockSync.return_value = svc + result = runner.invoke( + app, + [ + "sync", + "clone", + "--source", + str(source), + "--target", + "target", + "--target-dir", + str(tmp_path / "clone"), + "--no-create-buckets", + ], + ) + + assert result.exit_code == 0, result.output + assert svc.clone_project.call_args.kwargs["overrides"]["create_buckets"] is False def test_nested_override_value_errors(self, tmp_path: Path) -> None: """A non-scalar override value (fat-fingered colon) is rejected, not stringified.""" @@ -1078,3 +1125,307 @@ def test_missing_override_file_errors(self, tmp_path: Path) -> None: ], ) assert result.exit_code == 5 + + +def _write_buckets_export(root: Path, buckets: list[dict[str, Any]]) -> None: + """Write a current-format export: every record carries ``source_bucket`` (None if unset).""" + storage_dir = root / "storage" + storage_dir.mkdir(parents=True, exist_ok=True) + records = [{"source_bucket": None, **b} for b in buckets] + (storage_dir / "buckets.json").write_text(json.dumps(records), encoding="utf-8") + + +class TestParseBucketId: + def test_in_and_out_buckets_parse(self) -> None: + assert _parse_bucket_id("in.c-foo") == ("in", "foo") + assert _parse_bucket_id("out.c-bar") == ("out", "bar") + + def test_non_creatable_ids_return_none(self) -> None: + assert _parse_bucket_id("sys.foo") is None + assert _parse_bucket_id("in") is None + assert _parse_bucket_id("in.c-") is None + + +class TestCreateBucketsFromExport: + def test_creates_missing_and_skips_existing(self, tmp_path: Path) -> None: + _write_buckets_export( + tmp_path, + [ + { + "id": "in.c-new", + "name": "new", + "stage": "in", + "description": "d", + "backend": "snowflake", + }, + {"id": "in.c-existing", "name": "existing", "stage": "in", "description": ""}, + ], + ) + client = MagicMock() + client.list_buckets.return_value = [{"id": "in.c-existing"}] + + result = create_buckets_from_export(client, tmp_path, {}) + + assert result.created == ["in.c-new"] + assert result.skipped == ["in.c-existing"] + assert result.errors == [] + client.create_bucket.assert_called_once_with( + stage="in", name="new", description="d", backend="snowflake" + ) + + def test_bucket_map_is_applied_to_created_id(self, tmp_path: Path) -> None: + _write_buckets_export( + tmp_path, [{"id": "in.c-foo", "name": "foo", "stage": "in", "description": ""}] + ) + client = MagicMock() + client.list_buckets.return_value = [] + + result = create_buckets_from_export(client, tmp_path, {"in.c-foo": "out.c-bar"}) + + assert result.created == ["out.c-bar"] + client.create_bucket.assert_called_once_with( + stage="out", name="bar", description="", backend=None + ) + + def test_per_bucket_error_is_collected_not_raised(self, tmp_path: Path) -> None: + _write_buckets_export( + tmp_path, + [ + {"id": "in.c-ok", "name": "ok", "stage": "in", "description": ""}, + {"id": "in.c-bad", "name": "bad", "stage": "in", "description": ""}, + ], + ) + client = MagicMock() + client.list_buckets.return_value = [] + + def _create( + *, stage: str, name: str, description: str = "", backend: str | None = None + ) -> dict[str, Any]: + if name == "bad": + raise KeboolaApiError( + message="boom", status_code=400, error_code=ErrorCode.API_ERROR + ) + return {"id": f"{stage}.c-{name}"} + + client.create_bucket.side_effect = _create + + result = create_buckets_from_export(client, tmp_path, {}) + + assert result.created == ["in.c-ok"] + assert result.errors == [{"bucket_id": "in.c-bad", "error": "boom"}] + + def test_linked_bucket_is_linked_to_its_source(self, tmp_path: Path) -> None: + _write_buckets_export( + tmp_path, + [ + { + "id": "in.c-shared", + "name": "shared", + "stage": "in", + "description": "", + "source_bucket": {"bucket_id": "out.c-origin", "project_id": 42}, + } + ], + ) + client = MagicMock() + client.list_buckets.return_value = [] + + result = create_buckets_from_export(client, tmp_path, {}) + + client.create_bucket.assert_not_called() + client.link_bucket.assert_called_once_with( + source_project_id=42, source_bucket_id="out.c-origin", name="shared", stage="in" + ) + assert result.created == [] + assert result.linked == [ + { + "bucket_id": "in.c-shared", + "source_bucket_id": "out.c-origin", + "source_project_id": 42, + } + ] + + def test_refused_link_is_collected_not_raised(self, tmp_path: Path) -> None: + _write_buckets_export( + tmp_path, + [ + { + "id": "in.c-shared", + "name": "shared", + "stage": "in", + "description": "", + "source_bucket": {"bucket_id": "out.c-origin", "project_id": 42}, + } + ], + ) + client = MagicMock() + client.list_buckets.return_value = [] + client.link_bucket.side_effect = KeboolaApiError( + message="not shared", status_code=500, error_code=ErrorCode.STORAGE_JOB_FAILED + ) + + result = create_buckets_from_export(client, tmp_path, {}) + + client.create_bucket.assert_not_called() + assert result.linked == [] + assert result.errors == [{"bucket_id": "in.c-shared", "error": "not shared"}] + + def test_list_failure_is_collected_not_raised(self, tmp_path: Path) -> None: + _write_buckets_export( + tmp_path, [{"id": "in.c-foo", "name": "foo", "stage": "in", "description": ""}] + ) + client = MagicMock() + client.list_buckets.side_effect = KeboolaApiError( + message="denied", status_code=403, error_code=ErrorCode.API_ERROR + ) + + result = create_buckets_from_export(client, tmp_path, {}) + + client.create_bucket.assert_not_called() + assert result.errors == [{"bucket_id": "", "error": "Cannot list buckets: denied"}] + + def test_legacy_export_creates_nothing(self, tmp_path: Path) -> None: + # An older pull wrote no source_bucket key: a linked bucket would be + # created empty and every later re-clone would skip it. + storage_dir = tmp_path / "storage" + storage_dir.mkdir() + (storage_dir / "buckets.json").write_text( + json.dumps([{"id": "in.c-foo", "name": "foo", "stage": "in", "description": ""}]), + encoding="utf-8", + ) + client = MagicMock() + + result = create_buckets_from_export(client, tmp_path, {}) + + client.list_buckets.assert_not_called() + client.create_bucket.assert_not_called() + assert result.created == [] + assert len(result.errors) == 1 + assert "older kbagent pull" in result.errors[0]["error"] + + def test_no_export_returns_empty_without_listing(self, tmp_path: Path) -> None: + client = MagicMock() + + result = create_buckets_from_export(client, tmp_path, {}) + + assert (result.created, result.skipped, result.errors) == ([], [], []) + client.list_buckets.assert_not_called() + + +def _golden_source_with_buckets(root: Path) -> None: + """A golden reference tree plus a one-bucket ``storage/buckets.json`` export.""" + _golden_source(root) + storage_dir = root / "storage" + storage_dir.mkdir(parents=True, exist_ok=True) + (storage_dir / "buckets.json").write_text( + json.dumps( + [ + { + "id": "in.c-ref", + "name": "ref", + "stage": "in", + "description": "", + "source_bucket": None, + } + ] + ), + encoding="utf-8", + ) + + +class TestCloneCreatesBuckets: + """clone_project wires the storage-bucket create step (diff/push mocked).""" + + def _svc( + self, tmp_config_dir: Path, monkeypatch: pytest.MonkeyPatch + ) -> tuple[SyncService, MagicMock]: + client = MagicMock() + client.list_dev_branches.return_value = [{"id": 555, "isDefault": True}] + client.list_buckets.return_value = [] + svc = _service(tmp_config_dir, client) + monkeypatch.setattr( + svc, + "diff", + MagicMock( + return_value={ + "changes": [{"change_type": "added", "component_id": "keboola.ex-db"}] + } + ), + ) + monkeypatch.setattr( + svc, "push", MagicMock(return_value={"status": "pushed", "created": 1, "errors": []}) + ) + return svc, client + + def test_clone_creates_missing_buckets( + self, tmp_path: Path, tmp_config_dir: Path, monkeypatch: pytest.MonkeyPatch + ) -> None: + source = tmp_path / "golden" + _golden_source_with_buckets(source) + svc, client = self._svc(tmp_config_dir, monkeypatch) + + result = svc.clone_project( + source=source, + target_alias="target", + target_dir=tmp_path / "clone", + overrides={"create_buckets": True}, + ) + + client.create_bucket.assert_called_once_with( + stage="in", name="ref", description="", backend=None + ) + assert result["buckets_created"] == 1 + assert result["bucket_errors"] == [] + + def test_clone_creates_buckets_at_production_level_with_branch( + self, tmp_path: Path, tmp_config_dir: Path, monkeypatch: pytest.MonkeyPatch + ) -> None: + # Buckets are project-level: an explicit --branch scopes the push only. + source = tmp_path / "golden" + _golden_source_with_buckets(source) + svc, client = self._svc(tmp_config_dir, monkeypatch) + + svc.clone_project( + source=source, + target_alias="target", + target_dir=tmp_path / "clone", + overrides={"create_buckets": True}, + branch_override=52099, + ) + + client.create_bucket.assert_called_once() + assert "branch_id" not in client.create_bucket.call_args.kwargs + client.list_buckets.assert_called_once_with() + + def test_clone_opt_out_skips_bucket_creation( + self, tmp_path: Path, tmp_config_dir: Path, monkeypatch: pytest.MonkeyPatch + ) -> None: + source = tmp_path / "golden" + _golden_source_with_buckets(source) + svc, client = self._svc(tmp_config_dir, monkeypatch) + + result = svc.clone_project( + source=source, + target_alias="target", + target_dir=tmp_path / "clone", + overrides={"create_buckets": False}, + ) + + client.create_bucket.assert_not_called() + assert result["buckets_created"] == 0 + + def test_service_default_creates_buckets( + self, tmp_path: Path, tmp_config_dir: Path, monkeypatch: pytest.MonkeyPatch + ) -> None: + # The service default matches the CLI default, so a caller that passes + # no create_buckets key still gets a complete clone. + source = tmp_path / "golden" + _golden_source_with_buckets(source) + svc, client = self._svc(tmp_config_dir, monkeypatch) + + result = svc.clone_project( + source=source, target_alias="target", target_dir=tmp_path / "clone" + ) + + client.create_bucket.assert_called_once() + assert result["buckets_created"] == 1 diff --git a/tests/test_sync_storage_jobs.py b/tests/test_sync_storage_jobs.py index 849aacdae..e7383a702 100644 --- a/tests/test_sync_storage_jobs.py +++ b/tests/test_sync_storage_jobs.py @@ -85,6 +85,7 @@ "name": "c-data", "stage": "in", "description": "Input data bucket", + "backend": "snowflake", "tablesCount": 3, "dataSizeBytes": 1024000, "metadata": [{"key": "owner", "value": "team-a"}], @@ -483,9 +484,31 @@ def test_bucket_summary_format(self, tmp_config_dir: Path, tmp_path: Path) -> No assert b0["data_size_bytes"] == 1024000 assert b0["metadata"] == [{"key": "owner", "value": "team-a"}] + assert b0["backend"] == "snowflake" + assert b0["source_bucket"] is None + b1 = buckets[1] assert b1["id"] == "out.c-results" assert b1["tables_count"] == 1 + assert b1["backend"] == "" + + def test_linked_bucket_records_its_source(self, tmp_config_dir: Path, tmp_path: Path) -> None: + """A linked bucket's export keeps its source, so clone can skip it.""" + project_root = tmp_path / "project" + project_root.mkdir() + linked = { + "id": "in.c-shared", + "name": "c-shared", + "stage": "in", + "backend": "snowflake", + "sourceBucket": {"id": "out.c-origin", "project": {"id": 42, "name": "Origin"}}, + } + + write_storage_metadata(project_root, [linked], [], {}) + + buckets_file = project_root / STORAGE_DIR_NAME / STORAGE_BUCKETS_FILENAME + buckets = json.loads(buckets_file.read_text(encoding="utf-8")) + assert buckets[0]["source_bucket"] == {"bucket_id": "out.c-origin", "project_id": 42} def test_table_metadata_format(self, tmp_config_dir: Path, tmp_path: Path) -> None: """Per-table JSON files contain correct metadata fields."""