diff --git a/plugins/kbagent/agents/keboola-expert.md b/plugins/kbagent/agents/keboola-expert.md index 6a00d8f12..2c1befc96 100644 --- a/plugins/kbagent/agents/keboola-expert.md +++ b/plugins/kbagent/agents/keboola-expert.md @@ -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 --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__`) 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) | diff --git a/plugins/kbagent/skills/kbagent/references/gotchas.md b/plugins/kbagent/skills/kbagent/references/gotchas.md index cf6fb121e..a6105ef13 100644 --- a/plugins/kbagent/skills/kbagent/references/gotchas.md +++ b/plugins/kbagent/skills/kbagent/references/gotchas.md @@ -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=`) 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=`) 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` (`..`) that matches zero events. ## `--deny-writes` / `--deny-destructive` firewall (since 0.22.0) @@ -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 +dotted Queue `runId` (`..`). 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 diff --git a/src/keboola_agent_cli/client/queue.py b/src/keboola_agent_cli/client/queue.py index 1099c6d1e..e0b755a01 100644 --- a/src/keboola_agent_cli/client/queue.py +++ b/src/keboola_agent_cli/client/queue.py @@ -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.``) 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. diff --git a/src/keboola_agent_cli/server/routers/jobs.py b/src/keboola_agent_cli/server/routers/jobs.py index d3c324c98..446bfdec4 100644 --- a/src/keboola_agent_cli/server/routers/jobs.py +++ b/src/keboola_agent_cli/server/routers/jobs.py @@ -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 @@ -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: diff --git a/src/keboola_agent_cli/services/job_service.py b/src/keboola_agent_cli/services/job_service.py index e0aaf174e..aed0cf723 100644 --- a/src/keboola_agent_cli/services/job_service.py +++ b/src/keboola_agent_cli/services/job_service.py @@ -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 -- + ``..`` -- 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 @@ -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. diff --git a/tests/test_services.py b/tests/test_services.py index 896c8c86d..eddbe8bde 100644 --- a/tests/test_services.py +++ b/tests/test_services.py @@ -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() @@ -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).""" @@ -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`.