Skip to content

feat: API executor runs DataMixin prompts, reads python stages and reports token usage - #157

Open
timzsu wants to merge 78 commits into
mainfrom
zsu/api-executor-dataframe
Open

timzsu wants to merge 78 commits into
mainfrom
zsu/api-executor-dataframe

Conversation

@timzsu

@timzsu timzsu commented Sep 24, 2026 •

Copy link
Copy Markdown
Collaborator

Purpose

Lumilake's ingestion workflow runs as API-backed LLM ops mixed with deterministic list steps. The deterministic steps (explode a paper into passages, filter, split, pair into ragged groups, collapse groups) now run as FlowMesh python tasks. The batch API executor still served only list data, though. A row-wise op, an aggregate op and grouped rows couldn't run through it the way they run on the vLLM path, and nothing let an API task consume a python stage's groups. Nor did an API task report what it cost in model calls and tokens. This PR closes all three gaps.

Changes

  • API executor runs DataMixin prompts (api_executor.py, result schemas).
    • A dataframe spec issues one request per row. A graph_template spec issues one request over all rows, or one per group when the upstream is grouped.
    • A graph_template that reads a grouped column directly issues one prompt per group, carrying the whole group. The grouped column is repeated whole across its group's rows and never counts toward the row count, while an ungrouped column still expands to one value per row.
    • Grouped dataframe data returns one APIGroupItem {index, rows: [APIItem...]} per table, the way the vLLM executor regroups with _populate_table. List data stays one item per row.
  • Expressions read results through lists (graph_templates.py). _evaluate_expr maps attribute access over lists of result models, DataFrames and nested lists of groups, and it resolves a declared field by name or alias, never a method. So items.rows.json.choices[0].message.content (row-wise upstream) and items.json.choices[0].message.content (aggregate upstream) read API content downstream. A python stage is read the same way at value.items.output.
  • One grouping rule for the vLLM and API paths (graph_templates.py, mixins/data.py).
    • _evaluate_expr returns whether the value is grouped. It is grouped when the first attribute access over the items list yields one list per item, so each item whose output is a list of records is one group, and groups may be ragged. A list of DataFrames groups by item.
    • The dataframe and graph-template paths group on that flag instead of guessing from the value's shape, and the dataframe branch reuses _build_grouped_dataframes.
    • A cell whose per-row value is a list stays one cell, and index access into a per-item list (items.json.choices[0]) does not group.
  • Zero rows run zero requests. A grouped upstream with an empty group builds an empty table instead of a one-row table of blanks. A grouped spec with no rows at all returns an empty result instead of failing with "produced no rows". The task's status code comes from the first row found in any group, or 0 when there is none.
  • Grouped API results keep the json wire key. A plain dump of an APIResult holding APIGroupItems emits json for every row, which is where Lumilake reads items[].rows[].json.
  • The SDK mirrors APIGroupItem and APIUsage, and APIResult.items admits APIItem | APIGroupItem.
  • Token usage per API task (api_executor.py, result schemas). APIResult.usage is a typed APIUsage summed over the task's requests: prompt, completion and reasoning tokens, calls, failures, retries, truncated calls and wall time. A call whose response fails parsing, or that is cancelled, is still counted. Token counts are read only when they are integers, so an odd usage payload cannot fail a finished call.
  • Per-call statistics in the API executor. One debug line per finished call (row, attempts, status or exception, wall time, token counts, finish reason, backend), one debug summary per task (calls, failures, retries, latency p50/p95/max, summed tokens, per-backend counts), and a 60 s INFO heartbeat while calls are outstanding. No request or response content is logged. One lock guards the statistics, the heartbeat and the summary, and a call that exhausts its connection retries logs its real attempt count.
  • A cancelled grouped task with no rows reports cancelled. The task checks for a cancel once all its requests have returned. So a zero-row grouped task, or a cancel that lands during the last request, ends cancelled.
  • Declare tabulate (pyproject.toml, uv.lock, worker requirements). Graph-template rendering calls DataFrame.to_markdown, which needs tabulate, but nothing declared it. It was only in the lock through GPU packages, so a CPU worker failed any aggregate over a dataframe column with Missing optional dependency 'tabulate'. It goes in runtime-analytics, beside pandas, which the CPU worker image installs.
  • Docs. docs/WORKFLOWS.md documents grouped API results and how a dataframe API task reads a python stage (value.items.output.<field>, one group per item whose output is a list of records).

Design

No new task type. The API executor reuses DataMixin's prompt collection, and whole-list functions stay in main's python task. The result shape matches the vLLM executor's, so Lumilake reads a row-wise op the same way on the local path and the API path. Usage is reported per task, on the task's result. A client that wants a workflow's total sums its tasks' results, which also survives a Redis restart; the server keeps no per-workflow usage.

Test Plan

The repo's pre-commit gate (gitleaks, isort, black, ruff, codespell, mypy on the CI environment from uv sync --all-packages --group ci --frozen, requirements sync), then tests/server tests/worker tests/shared tests/sdk. The tests cover:

  • dataframe rows, a graph_template aggregate (a plain APIItem, read at items.json), ragged groups, and a graph_template reading a grouped column directly;
  • attribute mapping over lists, DataFrames and nested groups, and the json alias;
  • a python stage read at value.items.output: rows, ragged groups of three and two, an aggregate over those groups, and an empty group;
  • grouping decided from upstream structure: a list-valued cell and index access into a per-item list;
  • empty groups and a grouped spec with no rows, with their lineage, and the status code over grouped results;
  • the json wire key on a plain dump in the worker and the SDK, and SDK parity with the shared schema for APIGroupItem and APIUsage;
  • per-task usage: tokens and calls summed across rows, a retried 503 then 200 counted as one call with one retry, a single call's usage shape, and no usage when there are no rows;
  • a cancelled zero-row grouped task raising cancellation.

Then end to end on an isolated stack built from this branch, against https://lum.id/llm/v1/chat/completions with qwen3.8-27b: a python stage T0 emits two ragged groups of questions, a row-wise dataframe API task T1 reads them at value.items.output.q, a graph_template aggregate T2 reads T1 at items.rows.json.choices[0].message.content, a python task T3 collapses T1's groups into one item per answer, and an API stage T4 whose condition is never met.

Test Result

  • tests/server tests/worker tests/shared tests/sdk at 349dfae: 2431 passed, 16 skipped, 1 deselected (test_mp_executor_cleans_up_vllm, a GPU test whose vLLM engine does not start in this environment, with or without these changes).
  • Pre-commit: all hooks pass across the repo.
  • End to end at 349dfae, workflow wfl-e824a49e: DONE, 5 completed, 0 failed.
    • T0 emitted two groups of three and two question records.
    • T1 returned two APIGroupItems of three and two rows, every row HTTP 200 with its answer.
    • T2 returned two plain APIItems, both HTTP 200.
    • T3 returned five items, one per answer.
    • T4 was skipped: it never reached a worker and carries no usage.
    • T1 and T2 carry usage with calls 5 and 2, failures 0, retries 0, and non-zero tokens. Summed client-side: calls 7, prompt 186, completion 408, reasoning 370.
    • GET /workflows/{id} has no usage field.
    • T1 and T2 ran with timeout_sec: 300 (default 60) because the upstream model endpoint was slow during an earlier run.

@timzsu timzsu changed the title Zsu/api executor dataframe feat: API executor runs DataMixin prompts, and echo runs functions over whole lists Sep 24, 2026
@timzsu
timzsu added this pull request to stack #156 September 24, 2026 12:47
The API executor marks transient failures (5xx, 408, 429, connection
errors) as retryable but never retried them. Add a spec.api.retries field
(default 0, preserving current behavior) that controls how many times a
transient failure is re-issued before the task fails, with a fixed 1s
backoff between attempts. Non-retryable 4xx statuses and cancelled tasks
are never retried.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
@timzsu
timzsu force-pushed the zsu/api-executor-dataframe branch from f99ca7d to 9e357da Compare September 27, 2026 01:13
timzsu and others added 7 commits September 27, 2026 11:16
…during backoff

A cancellation left over from a previous task no longer leaks into the next
task on a reused warm executor: cancel() records the task it targets, and
run() clears a cancellation addressed to a different task while one aimed at
the task now starting still stands, all under a single lock so a racing
cancel is never lost.

The retry backoff now waits on the cancel event instead of sleeping, so a
cancelled task raises its cancellation as soon as it is signalled rather than
sitting out the full backoff.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
The interrupt monitor checks the runner's current task id and then calls
executor.cancel(task_id) without a lock spanning both, so a late
cancellation for a finished task can reach a warm executor that has since
started another task. The executor now records the active task id under
its cancel lock, sets the cancel event only when no run is in flight or
the cancellation matches the active task, and clears the active id when
the run returns or raises.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
A single recorded cancellation id let a late cancel for a prior task
overwrite a recorded cancel for the task about to start, so that task
ran despite its own cancellation. Track pending cancelled ids in a set
guarded by the cancel lock; run() consumes only its own id and clears
the rest as stale.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
When spec.data is present, the API executor issues one request per row,
substituting each row's prompt for the {{prompt}} placeholder in the
request body, and returns the responses row-aligned in APIResult.items.
The single-request path is unchanged. Reuses the DataMixin parsing infra
shared with the vLLM executor.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
The batch change added `items: list[APIItem]` to the server-side APIResult
without the SDK counterpart, so tests/sdk/test_schema_compat.py failed with
"APIResult missing server fields: ['items']".

The SDK mirrors every server result model and a compat test enforces that
they stay in step. This adds APIItem to the SDK payloads beside the existing
InferenceItem and Omni* items, carrying the identical `alias="json"` so the
wire-alias check passes too, adds `items` to the SDK APIResult, and registers
APIItem in _RESULT_MODEL_NAMES so the drift guard covers it from now on
rather than only the field that happened to break.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
@timzsu
timzsu force-pushed the zsu/api-executor-dataframe branch from 9e357da to b5aa734 Compare September 27, 2026 05:20
timzsu and others added 16 commits September 27, 2026 12:26
Require spec.data (as the vLLM executor does) and issue one request per
row in parallel, returning row-aligned APIResult.items. A single request
is a one-row spec.data. The request skeleton is built once and each worker
only substitutes its prompt into the prepared body. Concurrency is bounded
by spec.api.concurrency (default 8) and the client connection pool is sized
to it. Migrate api_two_stage.yaml and the n8n parser to carry spec.data.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Remove the separate api_batch.yaml template and make api_two_stage.yaml
demonstrate batching: multiple rows fan out through both stages, each row's
prompt fills the {{prompt}} worker-side slot, and stage two consumes stage
one. A single request is a one-row spec.data, so a dedicated batch template
no longer represents a distinct feature.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
The committed APIItem could not round-trip its own output. The executor
constructs the item by field name (response_json=...), which the json
alias plus extra="forbid" rejects. The worker serialises results with
model_dump_json() and no by_alias, so it emits the field name
response_json; the server then re-validates against the json alias and
rejects it. Every API-executor result would 422 on ingest.

populate_by_name=True fixes construction and ingest while keeping the
json alias accepted on input, so it is backward compatible in both
directions.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
The connection pool was hard-coded at 8 while the ThreadPoolExecutor used
the effective spec.api.concurrency, and the client cache key omitted
concurrency so a pool built for one value was reused for another. Size both
limits from the capped configured value and include concurrency in the cache
key.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
APIResult is batch-only: text lives per row, so a dependent stage's
placeholder must address items.0.text rather than a scalar text field that
is never populated. Update the n8n parser, the two-stage example, and add
an end-to-end dependent-stage resolution test.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
APIResult is batch-only, so the docstring no longer describes a scalar
single-request path. Add an SDK-side by-name construct/serialize/revalidate
test so the populate_by_name setting is guarded on both sides of the wire.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Replace the added inline comments with condensed docstring notes so the
non-obvious decisions (batch-only APIResult row addressing, the client cache
key and pool sizing) stay documented without inline comment blocks.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Move the module-scope nvmlInit() call in the GPU cleanup test behind a
skip guard so importing the module is side-effect free on GPU-less hosts,
including CI. Replace the inline import and inline type: ignore comments
with top-level imports and a cast, correct the pool-sizing comment to say
the pool is sized to the effective concurrency (capped), and assert the
issued requests as an unordered collection instead of assuming thread-pool
start order.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
…rral

Give the recording transport distinct status, usage, and headers per row
(with raise_for_status false and include_headers true) and assert them per
row so a mix-up in those fields is caught. Add a run-level test at an
uncapped concurrency (1 and 4) that isolates the effective value forwarded
to the client, and a test that observes concurrency: 1 serializing requests.
Defer nvmlInit() to a module fixture so importing the GPU test is
side-effect free, and add an integrated translated-n8n-to-stage-resolution
test covering the items.0.text production change end to end.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
A 5xx error body (e.g. a 503 {"error": ...}) has no OpenAI usage or
choices, so parsing it before the retryable classification raised a
non-retryable ExecutionError and a retryable 503 was permanently failed.
Classify HTTP errors first, and accept absent usage/text in the success
path so a no-usage 2xx or a 5xx under raise_for_status: false still
produces a row-aligned item.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
The executor inherited the no-op Executor.cancel(), so an interrupt during a
batch let queued rows continue and the batch could report success. Add a
per-instance cancel event, set it from cancel(), and check it before issuing
each request and before submitting queued futures, raising TaskCancelledError
so the runner aborts the batch.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
The SDK by-name round-trip test only validates and revalidates the SDK
model, so its docstring no longer claims a cross-package worker/server
boundary. Narrow the NVML fixture's skip to pynvml.NVMLError so a real
initialization bug is not swallowed as a missing GPU.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
A cancel landing between run()'s pre-run check and its event reset was
dropped. Hold a lock across the check-and-clear and in cancel() so no
interval exists where a cancel is accepted and then discarded. Document
the capped concurrency and cancellation behavior in WORKFLOWS.md.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
timzsu and others added 3 commits September 28, 2026 16:12
…cutor-dataframe

Adopt the typed ApiConfig / ApiResponseConfig from PR 141, reading the executor's config as typed attributes with no re-validation or clamping.

Keep PR 157's DataMixin row-wise and aggregate prompts, grouped dataframe results, per-call statistics, the zero-row grouped cancel and the retry attempt-count log, plus PR 141's early stop on a row failure.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Base automatically changed from zsu/api-executor-batch to main September 29, 2026 01:41
timzsu and others added 13 commits September 29, 2026 08:58
…aframe

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
The echo function mode ran whole-list caller code in-process through
safe_eval, duplicating the python task and exposing the same escapable
namespace. Remove the executor branch, argument resolution, the
expect_list switch and its safe_eval-only changes, their tests, and the
docs. Echo data.type list stays.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
A dataframe or graph_template spec reads a python stage's output through
the existing {node, path} argument, starting at the result as flowmesh
result fetch shows it (value.items.output). _evaluate_expr already
resolves a PythonResult and its grouping flag already follows the
upstream structure, so these tests cover the behavior without code
changes: rows from value.items.output, ragged groups of 3 and 2 into a
row-wise dataframe API task, a graph_template aggregate over those
groups, and an empty group.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
A dataframe column reads a python stage with node and a path that starts
at the result as flowmesh result fetch shows it, e.g. value.items.output.q.
When each item's output is a list of records, each item is one group and
group sizes may differ; a per-row list of scalars stays one cell value.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Carry the resolver's grouped flag through to the structural-message
renderer so a direct grouped column builds one aggregate prompt per
group instead of expanding every inner list into a row request.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude <noreply@anthropic.com>
A call that fails while parsing a completed response, or is cancelled
right after its response, was not recorded. Record it before
propagating, and count completed calls (not planned prompts) in the
summary.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude <noreply@anthropic.com>
A legacy type: function payload that also carries top-level items was
silently run as list mode. Reject it with an error directing the caller
to a python task instead.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude <noreply@anthropic.com>
The api executor groups every dataframe into one APIGroupItem per
table, matching the vLLM executor, so ungrouped data is a single group
rather than plain APIItems.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude <noreply@anthropic.com>
The previous fix always emitted one aggregate prompt per group, which
changed behaviour for every DataMixin caller over grouped data. Restore
main's one-message-per-row expansion and only keep a grouped column
whole per group, repeated across its rows.

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude <noreply@anthropic.com>

@kaiitunnz kaiitunnz left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Leave comments. PTAL.

Comment thread docs/WORKFLOWS.md Outdated
Comment thread src/worker/executors/utils/graph_templates.py
Comment thread src/worker/executors/utils/graph_templates.py
Comment thread src/worker/executors/utils/graph_templates.py Outdated
Comment thread src/shared/schemas/result/catalog.py Outdated
Comment thread docs/WORKFLOWS.md Outdated
Comment thread src/worker/executors/echo_executor.py Outdated
Comment thread src/worker/executors/api_executor.py Outdated
Comment thread docs/WORKFLOWS.md Outdated
Comment thread docs/WORKFLOWS.md Outdated
timzsu and others added 7 commits October 1, 2026 11:44
- docs: apply the suggested WORKFLOWS.md wording, drop the duplicate api example and the EXECUTORS.md echo section.
- echo: drop the function-type rejection branch.
- graph templates: a list of DataFrames groups by item, and one grouping rule (mapped items) covers the vLLM and API paths with no APIGroupItem special case; _model_attr reads model_extra.
- result catalog: drop the _route_group_items validator in shared and the SDK mirror.
- data mixin: the dataframe branch reuses _build_grouped_dataframes.
- API executor: APIResult.usage is a typed APIUsage (shared and SDK), and per-call and summary lines log at debug.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
- store each task's usage at result ingest (results.py, workflow registry, redis)
- GET /workflows/{id} sums stored usage; a merged vLLM parent records its own
  share so every call counts once; a completed model-calling task with no
  mappable usage makes the sum null (fail closed); no-model tasks contribute nothing
- delete usage keys on unregister; EventMonitor takes the workflow registry
- SDK Workflow.usage field

Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
… sum

The dispatcher marks a condition-skipped task succeeded without a usage
entry, so GET /workflows/{id} treated it as fail-closed and returned null
usage for the whole workflow. The skip path now records the no-usage marker
before the task is marked succeeded; the workflow registry is a required
Dispatcher argument.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
@timzsu
timzsu requested a review from kaiitunnz October 1, 2026 09:42
@timzsu timzsu changed the title feat: API executor runs DataMixin prompts, and echo runs functions over whole lists feat: API executor runs DataMixin prompts, reads python stages and reports token usage Oct 1, 2026
Remove the per-workflow usage store and sum: the usage keys in the Redis
client, the save/load/delete and UnknownUsage sentinel in the workflow
registry, the usage ingest and task-type mapping in the results router,
_sum_usage and the usage line in GET /workflows/{id}, and the registry
wiring threaded through main, monitoring, and the dispatcher. The SDK
Workflow.usage field and its test file go too.

Per-task APIResult.usage stays: clients sum each task's result themselves,
which also survives a Redis restart.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>

@kaiitunnz kaiitunnz left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Leave some comments. PTAL.

if all(isinstance(v, list) for v in value):
# Mapping over groups keeps grouped as-is; a raw list of lists
# groups only when its inner lists hold records.
if not grouped:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

This decides grouping from shape again after the first access, which the docstring rules out. items.json.choices.message.content comes back grouped (one group per item), while items.json.choices[0].message.content doesn't. With this change, the worker tests for graph templates, dataframe grouping and the API executor still pass.

Suggested change
if not grouped:
if not grouped and not mapped_items:

headers: dict[str, str] | None = None
response_json: Any = Field(default=None, alias="json")
usage: dict[str, Any] | None = None
usage: APIUsage | None = None

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

v0.1.9 stored the upstream usage dict here (prompt_tokens, completion_tokens, total_tokens), so those API results no longer validate. I checked with the SDK's AnyExecutorResult. Could the summed usage go in a new field, or should this PR be marked [BREAKING]?

dataframes: list[pd.DataFrame] = []
for group_idx in range(group_count):
max_len = 1
max_len = 0

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

With max_len starting at 0, an empty group fails when any column has a single repeated value (e.g. a constant data: list of one item). The length-1 column no longer matches 0 and raises "same number of rows per group". When max_len == 0, a one-value column should become [].

if not isinstance(arg, str):
materialized_args.append(_aggregate_structural_messages(columns, arg))
materialized_args.append(
_aggregate_structural_messages(columns, arg, set())

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Message arguments to a function step still expand grouped columns per row, while the step runs once per group. With groups [a0, a1] and [b0], the step gets a0 and a1, and b0 is dropped. Pass the grouped labels through here.

completion_tokens=completion_snapshot,
reasoning_tokens=reasoning_snapshot,
calls=total,
failures=failures_snapshot,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

failures is always 0. Every failed call raises and fails the task, and a failed task uploads no result, so no stored APIUsage ever counts a failure, and clients summing per-task usage never see failed calls. Drop the field, or report usage on failure too.

in_flight_lock = threading.Lock()
stats_lock = threading.Lock()

def _record_call(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

nit: run() now keeps nine counters as closure variables, copies each into a snapshot, and passes ten arguments to _log_summary, which makes it hard to follow. Could a small tracker class with a lock, record() and snapshot() -> APIUsage hold them? failed: bool here also shadows the failed event above.

if _is_expandable_group_value(group_value):
if key in grouped_labels:
# A grouped column is kept whole per group, repeated across its rows.
columns[key].extend([group_value] * row_count) # type: ignore

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

nit: CODE_STYLE.md rules out a bare # type: ignore; add the error code. The same applies to lines 320 and 322.

Comment thread docs/WORKFLOWS.md
Comment on lines +116 to +120
A body value that is exactly `{{prompt}}` is replaced by the row's prompt
object as-is (a message list stays a list of `{"role", "content"}` dicts). An
embedded `{{prompt}}` inside a longer string keeps string substitution: a
string prompt is inserted verbatim, and any other prompt value is rendered as
JSON.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

The first sentence repeats the paragraph above. Keep only the rule for an embedded placeholder. Consider merging the two paragraphs.

Comment thread docs/WORKFLOWS.md
Comment on lines +124 to +128
A `dataframe` spec returns one `APIGroupItem` per table, with that table's
responses in `rows`. A column is grouped when it reads one list of records per
upstream item, such as an upstream API task's `items.rows` or a python stage's
`items.output.<field>`; each list becomes one table. Downstream stages read the
responses through `items.rows`.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

A python task's return value is PythonResult.value, so the path is value.items.output.<field>. The per-group behavior of graph_template is also missing.

Suggested change
A `dataframe` spec returns one `APIGroupItem` per table, with that table's
responses in `rows`. A column is grouped when it reads one list of records per
upstream item, such as an upstream API task's `items.rows` or a python stage's
`items.output.<field>`; each list becomes one table. Downstream stages read the
responses through `items.rows`.
A `dataframe` spec returns one `APIGroupItem` per table, with that table's
responses in `rows`. A column is grouped when it reads one list of records per
upstream item, such as an upstream API task's `items.rows` or a python stage's
`value.items.output.<field>`; each list becomes one table. Downstream stages
read the responses through `items.rows`. A `graph_template` over a grouped
column sends one request per group and returns one `APIItem` for each.

Comment on lines +263 to +268
def test_api_result_items_union_matches() -> None:
"""APIResult.items must be the same union on both sides; the field-name
check alone misses a type drift (e.g. SDK-only list[APIItem])."""
srv_items = srv_results.APIResult.model_fields["items"].annotation
sdk_items = sdk_models.APIResult.model_fields["items"].annotation
assert _union_member_names(srv_items) == _union_member_names(sdk_items)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

This test is ad hoc. Consider extending assert_fields_match or add assert_field_types_match.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants