Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion plugins/kbagent/agents/keboola-expert.md
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,7 @@ been retired, so its absence is NOT a promise (see §1 Rule 6).
| Repartition / recluster a populated BigQuery table | `kbagent storage create-table --source-table-id <src> --time-partitioning-type DAY --time-partitioning-field created_at --clustering-field tenant_id` (BigQuery only) to copy the data into the new layout, then `swap-tables` to flip it in. `--source-table-id` derives the schema, so `--column` is forbidden. Range partitioning: all four `--range-partitioning-*` flags together. VERIFY the result with `storage table-detail --json` -> `.definition.timePartitioning` / `.clustering` (0.88.0+) -- `create-table` only echoes the layout you REQUESTED, so it proves nothing. See [storage-types-workflow.md](../skills/kbagent/references/storage-types-workflow.md) | -- | `CREATE TABLE ... AS SELECT` in a workspace (drops NOT NULL + primary key) |
| Re-seed a table without losing schema / PK / dependents | `kbagent storage truncate-table --project P --table-id in.c-foo.data [--dry-run] [--yes]` -- rows only, uniformly async-via-job on every branch; batch via repeated `--table-id` | -- | drop + recreate (loses descriptions, PK, sharing edges, and breaks every downstream reference); deleting rows via raw SQL in a workspace (bypasses the Storage audit trail) |
| Back up / restore a table around a risky change | `kbagent storage snapshot-create --table-id ...` then, to restore, `kbagent storage table-from-snapshot --snapshot-id ID --bucket-id B --name NEW` -- restore is always a NEW table (`--name` REQUIRED, no overwrite): verify it, then `swap-tables`. See [snapshot-workflow.md](../skills/kbagent/references/snapshot-workflow.md) | `storage snapshots` / `snapshot-detail` to find one | exporting to CSV as a "backup" (loses column types + PK); `create-table --snapshot-id` (not a thing) |
| Debug a failed job | `kbagent job detail --project P --job-id J --json` + `kbagent job run ... --log-tail-lines 200` | `kbagent workspace from-transformation` for SQL repro | "I think the issue is..." without reading logs |
| Debug a failed job | `kbagent job detail --project P --job-id J --log-tail-lines 200 --json` | `kbagent workspace from-transformation` for SQL repro | "I think the issue is..." without reading logs; a new job run only to see its logs (a writer can send its data again) |
| Ad-hoc SQL / row-count / type audit | `kbagent workspace create` + `workspace load` (since 0.91.0 auto-CLONEs eligible tables, else COPY; `--load-type` forces one and fails loudly if ineligible; COPY > 1 GiB needs `--force` outside a TTY; on a `--timeout` [default 300s] the job keeps running server-side and now exits 4/retryable, not 1 -- retry or poll `GET /v2/storage/jobs/{id}`, don't treat it as a hard failure) + `kbagent workspace query --sql "..."` -- results are inline and fast but **capped at `--limit`, default 500**: check `statements[].truncated` / `total_rows`, use `COUNT(*)` for counts, `--full` for the complete set | `kbagent workspace from-transformation` for existing-transform debugging; `workspace list --qs-compatible` for data-app reuse; read-only input-mapping (`KBC_<STACK>_<PROJECT>`) to query prod with no load at all | trusting a default `SELECT *` as the full result; querying Storage via raw Snowflake credentials outside the workspace abstraction |
| Export a FILTERED or INCREMENTAL slice of a table (no workspace) | `kbagent storage download-table --table-id ... --where-column status --where-value active [--where-operator eq\|neq] [--changed-since "-2 days"]` -- server-side filter on the credential-only export path | `kbagent workspace query` with a `WHERE` clause when you need real SQL | downloading the whole table then filtering locally |
| Run Keboola SQL / read-write Storage Files from INSIDE a Python process you control | `from keboola_agent_cli import Client` -- stateless `Client(url, token)`; `.query(workspace_id, sql)`, `.files.upload/.read_bytes/.list`; no subprocess, no `serve`, no config-dir. See [library-workflow.md](../skills/kbagent/references/library-workflow.md) | the CLI or `kbagent serve` REST when you are NOT already inside Python | shelling out to the `kbagent` binary from Python you control; using it for open-ended exploration (fixed set of typed ops) |
Expand Down
12 changes: 11 additions & 1 deletion plugins/kbagent/skills/kbagent/references/gotchas.md
Original file line number Diff line number Diff line change
Expand Up @@ -2511,7 +2511,7 @@ unknown -- do not try to parse a fallback message.
- `--timeout N` is a **local** deadline. When it elapses, kbagent issues `POST /jobs/{id}/kill` against the Queue API. Two outcomes:
- Kill succeeded -> exit **7** with `details.job` + `details.logTail`. The remote is definitely cancelled.
- Kill failed -> exit **4** with `details.logTail`, `retryable=True`. The remote **may still be running**; investigate before retrying.
- Inspecting events outside of `job run`: `kbagent job detail --project X --job-id N` does not fetch the log tail. To get the raw event stream, call the Storage Events API directly (`GET /v2/storage/events?runId=<runId>`) with the project token.
- Inspecting events outside of `job run`: `kbagent job detail --project X --job-id N --log-tail-lines N` (since v0.88.0). For the raw event stream, call the Storage Events API directly (`GET /v2/storage/events?runId=<job id>`) with the project token -- pass the plain job **id**, not the job's `runId`: a nested job (inside a flow, or a child row job) has a dotted `runId` (`<parent>.<child>.<id>`) that matches zero events.

## `--deny-writes` / `--deny-destructive` firewall (since 0.22.0)

Expand Down Expand Up @@ -4771,6 +4771,16 @@ The log tail is off by default so a plain `job detail` stays one API call.
that behaviour is unchanged, and its `--log-tail-lines` is capped at the same
maximum as `job detail`'s.

## `logTail` was empty for nested jobs

Fixed (since vNEXT). A job inside a flow, or a child row job of a row-based component, has a
Comment on lines +4774 to +4776

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔍 Update the job-debugging playbook

The debugging guidance still suggests a plain job detail followed by job run for logs. Plain detail omits logs; a tail-enabled detail now retrieves them from the existing nested job without rerunning it.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

dotted Queue `runId` (`<parent>.<child>.<job id>`). The Storage Events API
returns zero events for that dotted value, so `job detail --log-tail-lines`
and the `serve` job log stream came back with `logTail: []` for such jobs
(issue #787). kbagent now queries events by the job's own `id` (last `runId`
segment as a fallback). Top-level jobs, where `runId == id`, are unaffected.
`job run` starts only top-level jobs, so its failure tail was never affected.

## A scaffolded `keboola.flow` config can now be pushed from disk (since v0.89.0)

`config new --component-id keboola.flow` (no `--push`) used to write a
Expand Down
7 changes: 5 additions & 2 deletions src/keboola_agent_cli/client/queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -191,8 +191,11 @@ def fetch_job_events(self, run_id: str, limit: int | None = None) -> list[dict[s
a chronological "tail" should reverse the slice).

Args:
run_id: The job's ``runId`` (``job["runId"]``; falls back to
``job["id"]`` on legacy records where they match).
run_id: The value sent as the ``runId`` filter. Pass the job's own
``id`` (see ``services.job_service.resolve_events_run_id``):
a nested job's dotted Queue ``runId`` (``a.b.<id>``) matches
zero events, while the plain job id matches the job's events
(issue #787). For top-level jobs ``runId == id``.
limit: Optional server-side event cap. Storage API default is
about 100; pass an explicit value to cover long runs.

Expand Down
3 changes: 2 additions & 1 deletion src/keboola_agent_cli/server/routers/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
DEFAULT_LOG_TAIL_LINES,
DEFAULT_POLL_STRATEGY,
)
from ...services.job_service import resolve_events_run_id
from ..dependencies import ServiceRegistry, get_registry
from ..sse import json_event

Expand Down Expand Up @@ -164,7 +165,7 @@ async def gen() -> AsyncIterator[dict[str, str]]:
yield json_event({"status": current_status, "job": detail}, event="status")
last_status = current_status

run_id = str(detail.get("runId") or detail.get("id") or job_id)
run_id = resolve_events_run_id(detail) or str(job_id)
try:
client = registry.job._client_factory(proj.stack_url, proj.token)
try:
Expand Down
31 changes: 25 additions & 6 deletions src/keboola_agent_cli/services/job_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -132,14 +132,33 @@ def _job_sort_key(job: dict[str, Any], sort_by: str, sort_order: str) -> tuple[i
return (0, _DescStr(text) if desc else text, alias, id_tiebreak)


def resolve_events_run_id(job: dict[str, Any]) -> str:
"""Return the key to query Storage Events with for this job's log tail.

The job's own ``id`` wins. For a top-level Queue job ``runId == id``, so
this is unchanged there. For a NESTED job (a job inside a flow, or a
child row job of a row-based component) the Queue ``runId`` is dotted --
``<parent>.<child>.<job id>`` -- and ``GET /v2/storage/events?runId=``
with that dotted value returns ZERO events, while the plain job id
returns the job's events (issue #787). When ``id`` is absent we fall
back to the last segment of a dotted ``runId``, which is the job's own
id. Returns ``""`` when neither key is present.
"""
job_id = job.get("id")
if job_id is not None and str(job_id):
return str(job_id)
run_id = str(job.get("runId") or "")
return run_id.rsplit(".", 1)[-1]


def _safe_fetch_log_tail(client: Any, job: dict[str, Any], limit: int) -> list[dict[str, Any]]:
"""Fetch the last ``limit`` events for a job; never raises.

Resolves ``runId`` from the job dict (falls back to ``id`` on legacy
records where Queue v2 makes them equal). Storage Events API returns
newest -> oldest, which we keep as-is: a "tail" display wants the
most recent events first, and callers can reverse for chronology if
they prefer.
The events key comes from :func:`resolve_events_run_id` (job ``id``
first, never a dotted nested ``runId``, which returns zero events).
Storage Events API returns newest -> oldest, which we keep as-is: a
"tail" display wants the most recent events first, and callers can
reverse for chronology if they prefer.

Log-tail capture is a convenience surface; failing the whole command
because the events endpoint blipped would obscure the real underlying
Expand All @@ -148,7 +167,7 @@ def _safe_fetch_log_tail(client: Any, job: dict[str, Any], limit: int) -> list[d
"""
if limit <= 0:
return []
run_id = str(job.get("runId") or job.get("id") or "")
run_id = resolve_events_run_id(job)
if not run_id:
# Defensive: malformed job dict with neither runId nor id. Log so
# a real API regression doesn't get masked by the silent return.
Expand Down
50 changes: 46 additions & 4 deletions tests/test_services.py
Original file line number Diff line number Diff line change
Expand Up @@ -3780,8 +3780,9 @@ def test_run_job_warning_attaches_log_tail(self, tmp_config_dir: Path) -> None:
assert len(result["logTail"]) == 100
assert result["logTail"][0]["id"] == 249
assert result["logTail"][-1]["id"] == 150
# runId (not raw id) must have been the query key.
mock_client.fetch_job_events.assert_called_once_with("702-run", limit=100)
mock_client.fetch_job_events.assert_called_once_with(
"702", limit=100
) # job id wins over runId (#787)

def test_run_job_zero_tail_skips_fetch(self, tmp_config_dir: Path) -> None:
mock_client = MagicMock()
Expand Down Expand Up @@ -3880,8 +3881,9 @@ def test_run_job_timeout_issues_kill_and_raises_terminated(self, tmp_config_dir:
details = exc_info.value.details
assert details["job"]["status"] == "terminated"
assert details["logTail"] == [{"uuid": "u1", "message": "x"}]
# runId from the terminated job detail used as the lookup key.
mock_client.fetch_job_events.assert_called_once_with("705-run", limit=200)
mock_client.fetch_job_events.assert_called_once_with(
"705", limit=200
) # job id wins over runId (#787)

def test_run_job_timeout_kill_fails_falls_back(self, tmp_config_dir: Path) -> None:
"""If kill_job AND the follow-up GET fail, surface QUEUE_JOB_TIMEOUT (retryable)."""
Expand Down Expand Up @@ -4057,6 +4059,46 @@ def test_missing_created_sorts_last(self) -> None:
assert tail[1]["uuid"] == "no_created"


class TestSafeFetchLogTailNestedJobs:
"""Issue #787: nested jobs carry a dotted Queue runId that matches zero
Storage events; the log tail must be fetched by the job's own id."""

def _client(self) -> MagicMock:
client = MagicMock()
client.fetch_job_events.return_value = [{"uuid": "e1", "created": "2026-09-01"}]
return client

def test_nested_job_uses_id_not_dotted_run_id(self) -> None:
from keboola_agent_cli.services.job_service import _safe_fetch_log_tail

client = self._client()
job = {"id": "56453151", "runId": "56452146.56453149.56453151"}
tail = _safe_fetch_log_tail(client, job, limit=50)
client.fetch_job_events.assert_called_once_with("56453151", limit=50)
assert tail == [{"uuid": "e1", "created": "2026-09-01"}]

def test_only_dotted_run_id_uses_last_segment(self) -> None:
from keboola_agent_cli.services.job_service import _safe_fetch_log_tail

client = self._client()
_safe_fetch_log_tail(client, {"runId": "56452146.56453149.56453151"}, limit=5)
client.fetch_job_events.assert_called_once_with("56453151", limit=5)

def test_top_level_job_unchanged(self) -> None:
from keboola_agent_cli.services.job_service import _safe_fetch_log_tail

client = self._client()
_safe_fetch_log_tail(client, {"id": 123, "runId": "123"}, limit=5)
client.fetch_job_events.assert_called_once_with("123", limit=5)

def test_neither_id_nor_run_id_returns_empty_without_call(self) -> None:
from keboola_agent_cli.services.job_service import _safe_fetch_log_tail

client = self._client()
assert _safe_fetch_log_tail(client, {"status": "error"}, limit=5) == []
client.fetch_job_events.assert_not_called()


class TestJobServiceVariableValuesResolution:
"""Tests for `resolve_variable_values_id` + auto-resolution in `run_job`.

Expand Down
Loading