From 665baf6741241fe3258578e672e649fee881b3ac Mon Sep 17 00:00:00 2001 From: Petr Date: Tue, 24 Mar 2026 22:48:32 +0100 Subject: [PATCH 1/2] Add storage commands with shared bucket path resolution Resolves #35. New command namespace `kbagent storage`: - storage buckets: list buckets with sharing/linked metadata - storage bucket-detail: resolve Snowflake paths for shared buckets - storage tables: list tables by project/bucket Key feature: bucket-detail resolves the correct Snowflake database for linked/shared buckets (e.g. in.c-db linked from project 1507 -> "sapi_1507"."in.c-keboola-ex-db-mysql"."table_name"). This info is not available via MCP tools. --- plugins/kbagent/skills/kbagent/SKILL.md | 3 + src/keboola_agent_cli/cli.py | 5 + src/keboola_agent_cli/client.py | 49 +++++ src/keboola_agent_cli/commands/context.py | 26 +++ src/keboola_agent_cli/commands/storage.py | 208 ++++++++++++++++++ .../services/storage_service.py | 205 +++++++++++++++++ 6 files changed, 496 insertions(+) create mode 100644 src/keboola_agent_cli/commands/storage.py create mode 100644 src/keboola_agent_cli/services/storage_service.py diff --git a/plugins/kbagent/skills/kbagent/SKILL.md b/plugins/kbagent/skills/kbagent/SKILL.md index cc92119f..ec2702f3 100644 --- a/plugins/kbagent/skills/kbagent/SKILL.md +++ b/plugins/kbagent/skills/kbagent/SKILL.md @@ -69,6 +69,9 @@ This prints all commands, flags, workflows, and tips. Read it fully before proce | Load tables into a workspace | `kbagent workspace load --project PROJECT --workspace-id WORKSPACE-ID --tables TABLES` | | Execute SQL query in a workspace via Query Service | `kbagent workspace query --project PROJECT --workspace-id WORKSPACE-ID` | | Create a workspace from a transformation config | `kbagent workspace from-transformation --project PROJECT --component-id COMPONENT-ID --config-id CONFIG-ID` | +| List storage buckets with sharing/linked bucket information | `kbagent storage buckets` | +| Show detailed bucket info including Snowflake direct access paths | `kbagent storage bucket-detail --project PROJECT --bucket-id BUCKET-ID` | +| List storage tables from a project | `kbagent storage tables --project PROJECT` | | Initialize a sync working directory for a Keboola project | `kbagent sync init --project PROJECT` | | Download all configurations from a Keboola project to local files | `kbagent sync pull --project PROJECT` | | Show which local configurations have been modified, added, or deleted | `kbagent sync status` | diff --git a/src/keboola_agent_cli/cli.py b/src/keboola_agent_cli/cli.py index d67ee76b..670e5d10 100644 --- a/src/keboola_agent_cli/cli.py +++ b/src/keboola_agent_cli/cli.py @@ -17,6 +17,7 @@ from .commands.org import org_app from .commands.project import project_app from .commands.repl import repl_command +from .commands.storage import storage_app from .commands.sync import sync_app from .commands.tool import tool_app from .commands.version import version_command @@ -33,6 +34,7 @@ from .services.mcp_service import McpService from .services.org_service import OrgService from .services.project_service import ProjectService +from .services.storage_service import StorageService from .services.sync_service import SyncService from .services.version_service import VersionService from .services.workspace_service import WorkspaceService @@ -53,6 +55,7 @@ app.add_typer(explorer_app, name="explorer") app.add_typer(llm_app, name="llm") app.add_typer(workspace_app, name="workspace") +app.add_typer(storage_app, name="storage") app.add_typer(sync_app, name="sync") app.command("context")(context_command) app.command("doctor")(doctor_command) @@ -128,6 +131,7 @@ def main( org_service = OrgService(config_store=config_store) mcp_service = McpService(config_store=config_store) branch_service = BranchService(config_store=config_store) + storage_service = StorageService(config_store=config_store) sync_service = SyncService(config_store=config_store) workspace_service = WorkspaceService(config_store=config_store) kbc_service = KbcService(config_store=config_store) @@ -153,6 +157,7 @@ def main( ctx.obj["org_service"] = org_service ctx.obj["mcp_service"] = mcp_service ctx.obj["branch_service"] = branch_service + ctx.obj["storage_service"] = storage_service ctx.obj["sync_service"] = sync_service ctx.obj["workspace_service"] = workspace_service ctx.obj["kbc_service"] = kbc_service diff --git a/src/keboola_agent_cli/client.py b/src/keboola_agent_cli/client.py index 4fa817d0..6fce40b4 100644 --- a/src/keboola_agent_cli/client.py +++ b/src/keboola_agent_cli/client.py @@ -516,6 +516,55 @@ def list_buckets(self, include: str | None = None) -> list[dict[str, Any]]: response = self._request("GET", "/v2/storage/buckets", params=params) return response.json() + def get_bucket_detail( + self, + bucket_id: str, + branch_id: int | None = None, + ) -> dict[str, Any]: + """Get detailed information about a storage bucket. + + Returns full bucket metadata including sharing/linked info + (sourceBucket, sourceTable with project references). + + Args: + bucket_id: Bucket ID (e.g. 'in.c-db'). + branch_id: If set, target a specific dev branch. + + Returns: + Bucket detail dict from the API. + """ + prefix = f"/v2/storage/branch/{branch_id}" if branch_id else "/v2/storage" + safe_id = quote(bucket_id, safe="") + response = self._request("GET", f"{prefix}/buckets/{safe_id}") + return response.json() + + def list_tables( + self, + bucket_id: str | None = None, + branch_id: int | None = None, + include: str | None = None, + ) -> list[dict[str, Any]]: + """List storage tables, optionally filtered by bucket. + + Args: + bucket_id: If set, list tables only from this bucket. + branch_id: If set, target a specific dev branch. + include: Optional include parameter (e.g. 'columns'). + + Returns: + List of table dicts from the API. + """ + prefix = f"/v2/storage/branch/{branch_id}" if branch_id else "/v2/storage" + params: dict[str, str] = {} + if include: + params["include"] = include + if bucket_id: + safe_id = quote(bucket_id, safe="") + response = self._request("GET", f"{prefix}/buckets/{safe_id}/tables", params=params) + else: + response = self._request("GET", f"{prefix}/tables", params=params) + return response.json() + def list_jobs( self, component_id: str | None = None, diff --git a/src/keboola_agent_cli/commands/context.py b/src/keboola_agent_cli/commands/context.py index 24031c00..95b2fb26 100644 --- a/src/keboola_agent_cli/commands/context.py +++ b/src/keboola_agent_cli/commands/context.py @@ -142,6 +142,32 @@ Example: kbagent --json job detail --project prod --job-id 148512262 +### Storage (Buckets and Tables) + + kbagent storage buckets [--project NAME] + List storage buckets with sharing/linked bucket information. + Shows which buckets are linked from other projects, including the + source project ID and name. This info is NOT available via MCP tools. + Examples: + kbagent --json storage buckets + kbagent --json storage buckets --project prod + + kbagent storage bucket-detail --project NAME --bucket-id BUCKET_ID + Show detailed bucket info including Snowflake direct access paths. + For linked/shared buckets, resolves the correct Snowflake database + and schema from the source project. Each table includes a ready-to-use + fully-qualified Snowflake path with proper quoting. + CRITICAL for direct Snowflake access: linked buckets live in a different + database than the current project (e.g. sapi_1507 instead of sapi_226). + Example: + kbagent --json storage bucket-detail --project slevomat --bucket-id in.c-db + + kbagent storage tables --project NAME [--bucket-id BUCKET_ID] + List storage tables, optionally filtered by bucket. + Example: + kbagent --json storage tables --project prod + kbagent --json storage tables --project prod --bucket-id in.c-main + ### Data Lineage kbagent lineage [--project NAME] diff --git a/src/keboola_agent_cli/commands/storage.py b/src/keboola_agent_cli/commands/storage.py new file mode 100644 index 00000000..404e69fb --- /dev/null +++ b/src/keboola_agent_cli/commands/storage.py @@ -0,0 +1,208 @@ +"""Storage commands - buckets, tables, and direct access path resolution. + +Provides direct Storage API access including sharing/linked bucket metadata +that is not available via MCP tools. +""" + +import typer + +from ..errors import ConfigError, KeboolaApiError +from ._helpers import emit_project_warnings, get_formatter, get_service, map_error_to_exit_code + +storage_app = typer.Typer(help="Browse storage buckets and tables") + + +@storage_app.command("buckets") +def storage_buckets( + ctx: typer.Context, + project: list[str] | None = typer.Option( + None, + "--project", + help="Project alias (can be repeated for multiple projects)", + ), +) -> None: + """List storage buckets with sharing/linked bucket information. + + Shows which buckets are linked from other projects, including the + source project ID and name. This information is not available via + MCP tools. + """ + formatter = get_formatter(ctx) + service = get_service(ctx, "storage_service") + + try: + result = service.list_buckets(aliases=project) + except ConfigError as exc: + formatter.error(message=exc.message, error_code="CONFIG_ERROR") + raise typer.Exit(code=5) from None + + if formatter.json_mode: + formatter.output(result) + else: + from rich.table import Table + + buckets = result["buckets"] + if not buckets: + formatter.console.print("[dim]No buckets found.[/dim]") + return + + # Group by project + by_project: dict[str, list[dict]] = {} + for b in buckets: + alias = b["project_alias"] + by_project.setdefault(alias, []).append(b) + + for alias, proj_buckets in by_project.items(): + table = Table(title=f"Buckets - {alias}") + table.add_column("Bucket ID", style="bold cyan") + table.add_column("Stage", style="dim") + table.add_column("Rows", justify="right") + table.add_column("Linked From", style="yellow") + + for b in proj_buckets: + linked = "" + if b["is_linked"]: + linked = f"{b['source_project_name']} (#{b['source_project_id']})" + table.add_row( + b["id"], + b["stage"], + str(b["rows_count"]), + linked, + ) + + formatter.console.print(table) + formatter.console.print() + + emit_project_warnings(formatter, result) + + +@storage_app.command("bucket-detail") +def storage_bucket_detail( + ctx: typer.Context, + project: str = typer.Option( + ..., + "--project", + help="Project alias", + ), + bucket_id: str = typer.Option( + ..., + "--bucket-id", + help="Bucket ID (e.g. in.c-db)", + ), +) -> None: + """Show detailed bucket info including Snowflake direct access paths. + + For linked/shared buckets, resolves the correct Snowflake database + and schema from the source project. Each table includes a ready-to-use + fully-qualified Snowflake path with proper quoting. + """ + formatter = get_formatter(ctx) + service = get_service(ctx, "storage_service") + + try: + result = service.get_bucket_detail(alias=project, bucket_id=bucket_id) + except ConfigError as exc: + formatter.error(message=exc.message, error_code="CONFIG_ERROR") + raise typer.Exit(code=5) from None + except KeboolaApiError as exc: + formatter.error(message=exc.message, error_code=exc.error_code) + raise typer.Exit(code=map_error_to_exit_code(exc)) from None + + if formatter.json_mode: + formatter.output(result) + else: + formatter.console.print(f"[bold]Bucket:[/bold] {result['bucket_id']}") + formatter.console.print(f" Display name: {result['display_name']}") + formatter.console.print(f" Backend: {result['backend']}") + + if result["is_linked"]: + formatter.console.print( + f" [yellow]Linked from:[/yellow] " + f"{result['source_project_name']} (#{result['source_project_id']})" + ) + formatter.console.print(f" Source bucket: {result['source_bucket_id']}") + + formatter.console.print(f" Snowflake DB: {result['snowflake_database']}") + formatter.console.print(f" Snowflake schema: {result['snowflake_schema']}") + formatter.console.print(f" Tables: {result['table_count']}") + + if result["tables"]: + formatter.console.print() + from rich.table import Table + + table = Table(title="Tables with Snowflake paths") + table.add_column("Table", style="bold") + table.add_column("Snowflake Path", style="green") + table.add_column("Alias", style="dim") + + for t in result["tables"][:50]: # limit display + table.add_row( + t["name"], + t["snowflake_path"], + "yes" if t["is_alias"] else "", + ) + + formatter.console.print(table) + + if len(result["tables"]) > 50: + formatter.console.print( + f" ... and {len(result['tables']) - 50} more (use --json for full list)" + ) + + +@storage_app.command("tables") +def storage_tables( + ctx: typer.Context, + project: str = typer.Option( + ..., + "--project", + help="Project alias", + ), + bucket_id: str | None = typer.Option( + None, + "--bucket-id", + help="Filter tables by bucket ID", + ), +) -> None: + """List storage tables from a project.""" + formatter = get_formatter(ctx) + service = get_service(ctx, "storage_service") + + try: + result = service.list_tables(alias=project, bucket_id=bucket_id) + except ConfigError as exc: + formatter.error(message=exc.message, error_code="CONFIG_ERROR") + raise typer.Exit(code=5) from None + except KeboolaApiError as exc: + formatter.error(message=exc.message, error_code=exc.error_code) + raise typer.Exit(code=map_error_to_exit_code(exc)) from None + + if formatter.json_mode: + formatter.output(result) + else: + from rich.table import Table + + tables = result["tables"] + if not tables: + formatter.console.print("[dim]No tables found.[/dim]") + return + + table = Table(title=f"Tables - {result['project_alias']}") + table.add_column("Table ID", style="bold cyan") + table.add_column("Rows", justify="right") + table.add_column("Size", justify="right", style="dim") + table.add_column("Last Import", style="dim") + + for t in tables: + size_mb = t["data_size_bytes"] / (1024 * 1024) if t["data_size_bytes"] else 0 + last_import = t.get("last_import_date", "") + if last_import and "T" in last_import: + last_import = last_import.split("T")[0] + table.add_row( + t["id"], + str(t["rows_count"]), + f"{size_mb:.1f} MB", + last_import, + ) + + formatter.console.print(table) diff --git a/src/keboola_agent_cli/services/storage_service.py b/src/keboola_agent_cli/services/storage_service.py new file mode 100644 index 00000000..cf4234b2 --- /dev/null +++ b/src/keboola_agent_cli/services/storage_service.py @@ -0,0 +1,205 @@ +"""Storage service - business logic for bucket and table operations. + +Provides direct access to Storage API data including sharing/linked bucket +metadata that MCP tools strip from responses. +""" + +import logging +from typing import Any + +from ..models import ProjectConfig +from .base import BaseService + +logger = logging.getLogger(__name__) + + +class StorageService(BaseService): + """Business logic for storage bucket and table operations. + + Supports multi-project parallel queries for listing operations. + """ + + def list_buckets( + self, + aliases: list[str] | None = None, + ) -> dict[str, Any]: + """List storage buckets from one or more projects. + + Includes sharing/linked bucket metadata (sourceBucket, sourceProject) + that is not available via MCP tools. + + Returns: + Dict with 'buckets' list and 'errors' list. + """ + projects = self.resolve_projects(aliases) + successes, errors = self._run_parallel(projects, self._fetch_buckets) + + buckets: list[dict[str, Any]] = [] + for result in successes: + alias = result[0] + for bucket in result[1]: + entry: dict[str, Any] = { + "project_alias": alias, + "id": bucket.get("id", ""), + "display_name": bucket.get("displayName", bucket.get("name", "")), + "stage": bucket.get("stage", ""), + "backend": bucket.get("backend", ""), + "rows_count": bucket.get("rowsCount", 0), + "data_size_bytes": bucket.get("dataSizeBytes", 0), + "description": bucket.get("description", ""), + "is_linked": False, + "source_project_id": None, + "source_project_name": "", + "source_bucket_id": "", + } + + # Enrich with sharing info + source = bucket.get("sourceBucket") + if source: + entry["is_linked"] = True + src_project = source.get("project", {}) + entry["source_project_id"] = src_project.get("id") + entry["source_project_name"] = src_project.get("name", "") + entry["source_bucket_id"] = source.get("id", "") + + buckets.append(entry) + + return {"buckets": buckets, "errors": errors} + + def get_bucket_detail( + self, + alias: str, + bucket_id: str, + ) -> dict[str, Any]: + """Get detailed bucket info including tables and sharing metadata. + + For linked buckets, includes the Snowflake direct access path. + + Returns: + Dict with bucket detail, tables, and resolved Snowflake paths. + """ + projects = self.resolve_projects([alias]) + project = projects[alias] + + client = self._client_factory(project.stack_url, project.token) + try: + token_info = client.verify_token() + bucket = client.get_bucket_detail(bucket_id) + finally: + client.close() + + project_id = token_info.project_id + source = bucket.get("sourceBucket") + + result: dict[str, Any] = { + "project_alias": alias, + "project_id": project_id, + "bucket_id": bucket.get("id", ""), + "display_name": bucket.get("displayName", ""), + "stage": bucket.get("stage", ""), + "description": bucket.get("description", ""), + "backend": bucket.get("backend", ""), + "is_linked": source is not None, + } + + # Resolve Snowflake paths + if source: + src_project = source.get("project", {}) + src_project_id = src_project.get("id") + src_bucket_id = source.get("id", "") + result["source_project_id"] = src_project_id + result["source_project_name"] = src_project.get("name", "") + result["source_bucket_id"] = src_bucket_id + result["snowflake_database"] = f"sapi_{src_project_id}" + result["snowflake_schema"] = src_bucket_id + else: + result["source_project_id"] = None + result["source_project_name"] = "" + result["source_bucket_id"] = "" + result["snowflake_database"] = f"sapi_{project_id}" + result["snowflake_schema"] = bucket.get("id", "") + + # Build table list with Snowflake paths + tables: list[dict[str, Any]] = [] + for table in bucket.get("tables", []): + table_name = table.get("name", "") + sf_db = result["snowflake_database"] + sf_schema = result["snowflake_schema"] + tables.append( + { + "id": table.get("id", ""), + "name": table_name, + "display_name": table.get("displayName", table_name), + "is_alias": table.get("isAlias", False), + "snowflake_path": f'"{sf_db}"."{sf_schema}"."{table_name}"', + } + ) + + result["tables"] = tables + result["table_count"] = len(tables) + + return result + + def list_tables( + self, + alias: str, + bucket_id: str | None = None, + ) -> dict[str, Any]: + """List tables from a project, optionally filtered by bucket. + + Returns: + Dict with 'tables' list. + """ + projects = self.resolve_projects([alias]) + project = projects[alias] + + client = self._client_factory(project.stack_url, project.token) + try: + raw_tables = client.list_tables(bucket_id=bucket_id) + finally: + client.close() + + tables = [ + { + "project_alias": alias, + "id": t.get("id", ""), + "name": t.get("name", ""), + "display_name": t.get("displayName", t.get("name", "")), + "bucket_id": t.get("bucket", {}).get("id", "") + if isinstance(t.get("bucket"), dict) + else "", + "rows_count": t.get("rowsCount", 0), + "data_size_bytes": t.get("dataSizeBytes", 0), + "is_alias": t.get("isAlias", False), + "last_import_date": t.get("lastImportDate", ""), + } + for t in raw_tables + ] + + return {"tables": tables, "project_alias": alias} + + # ------------------------------------------------------------------ + # Parallel workers + # ------------------------------------------------------------------ + + def _fetch_buckets( + self, alias: str, project: ProjectConfig + ) -> tuple[str, list[dict[str, Any]], bool]: + """Fetch buckets for a single project (worker for _run_parallel).""" + from ..errors import KeboolaApiError + + client = self._client_factory(project.stack_url, project.token) + try: + buckets = client.list_buckets(include="linkedBuckets") + return (alias, buckets, True) + except KeboolaApiError as exc: + return ( + alias, + { + "project_alias": alias, + "error_code": exc.error_code, + "message": exc.message, + }, + ) + finally: + client.close() From 472a3e3b719111c972134de1218bbb3f0cf92492 Mon Sep 17 00:00:00 2001 From: Petr Date: Tue, 24 Mar 2026 22:54:18 +0100 Subject: [PATCH 2/2] Add Snowflake quoting best practice to workspace tips --- src/keboola_agent_cli/commands/context.py | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/src/keboola_agent_cli/commands/context.py b/src/keboola_agent_cli/commands/context.py index 95b2fb26..27af31c1 100644 --- a/src/keboola_agent_cli/commands/context.py +++ b/src/keboola_agent_cli/commands/context.py @@ -584,6 +584,19 @@ kbagent --json workspace load --project prod --workspace-id WS_ID --tables in.c-bucket.my-table kbagent --json workspace query --project prod --workspace-id WS_ID --sql "SELECT * FROM \"my-table\" LIMIT 10" + IMPORTANT -- Snowflake quoting rules for workspace queries: + Snowflake converts unquoted identifiers to UPPERCASE. If a database, + schema, or table name contains lowercase letters, dots, or hyphens, + you MUST double-quote it. This applies to ALL identifiers: + WRONG: SELECT * FROM sap_9.my_schema.my_table + (Snowflake reads this as SAP_9.MY_SCHEMA.MY_TABLE -- not found!) + RIGHT: SELECT * FROM "sap_9"."my_schema"."my_table" + Best practice: ALWAYS double-quote database, schema, and table names + in workspace queries, even if they look like they don't need it. + Keboola workspace database/schema names are often lowercase. + For shared/linked buckets, use 'kbagent storage bucket-detail' to get + the correct fully-qualified Snowflake path (source project DB differs). + 16. Setting up projects -- two approaches: a) Single project (you have a Storage API token):