Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
78 commits
Select commit Hold shift + click to select a range
22b222e
feat: add user-specified retries to the API executor
timzsu Sep 23, 2026
ad4752f
fix(worker): reset stale API executor cancellations and honor cancel …
timzsu Sep 27, 2026
275ca45
fix(worker): shorten stale cancellation comment to a single line
timzsu Sep 27, 2026
67b795b
fix(worker): reject cancellations addressed to a different active task
timzsu Sep 27, 2026
38a2a07
fix(worker): shorten stale cancellation comment to a single line
timzsu Sep 27, 2026
096a58a
fix(worker): keep pending cancellations by task id in the API executor
timzsu Sep 27, 2026
0fb3d26
feat: batch multiple HTTP requests in the API executor
timzsu Sep 18, 2026
d2fe09c
fix(sdk): mirror APIResult.items and APIItem in the SDK models
timzsu Sep 18, 2026
b5583ea
refactor: single parallel path for the API executor
timzsu Sep 18, 2026
630d55a
docs: fold batching into the two-stage api template
timzsu Sep 18, 2026
aff4e63
fix: let APIItem round-trip its own result output
timzsu Sep 18, 2026
da6a2df
fix: size the API HTTP pool from the capped concurrency
timzsu Sep 18, 2026
d3e31b8
fix: resolve API stage dependencies from the first row's text
timzsu Sep 18, 2026
4358a0b
docs: correct the APIResult docstring and guard SDK by-name construction
timzsu Sep 18, 2026
dc9da05
refactor: fold API dependency and cache-key notes into docstrings
timzsu Sep 18, 2026
da2aef1
fix: make the GPU test import-safe and clean up test style
timzsu Sep 18, 2026
0b4f12c
docs: fold the concurrency cap note into the API executor docstring
timzsu Sep 18, 2026
7dff652
test: cover row-specific fields, effective concurrency, and NVML defe…
timzsu Sep 18, 2026
9524d1c
fix: classify API HTTP errors before strict payload parsing
timzsu Sep 18, 2026
0341004
feat: add task-scoped cancellation to the API executor
timzsu Sep 18, 2026
3f9e5c7
test: correct the SDK APIItem docstring and narrow the NVML skip
timzsu Sep 18, 2026
40bf0a2
pin pre-start and in-flight cancellation guards separately
timzsu Sep 18, 2026
9691501
fix: recheck cancellation during in-flight requests and before run
timzsu Sep 19, 2026
76e218c
fix: make API executor cancel check-and-reset atomic under a lock
timzsu Sep 19, 2026
d4ced9f
fix: correct cancellation docs and rename collection-guard test
timzsu Sep 19, 2026
b6c01a6
test: merge API executor batch tests into test_api_executor.py
timzsu Sep 19, 2026
de41f50
refactor: revert gitignore and GPU test, condense added docstrings
timzsu Sep 19, 2026
e4991dd
fix: use flowmesh/v1 apiVersion in API executor tests
timzsu Sep 19, 2026
2109405
test: adapt API executor retry tests to the merged batch executor
timzsu Sep 24, 2026
7d2ddb9
fix(worker): substitute a structured prompt object into the API body
timzsu Sep 27, 2026
4a05001
fix(schema): serialize APIItem's json field under its wire alias
timzsu Sep 27, 2026
d11996d
test: assert completed requests before collection in cancellation test
timzsu Sep 27, 2026
fc8d5c2
test: bound the submission wait in the collection-guard cancellation …
timzsu Sep 27, 2026
821b2d0
feat: API executor runs DataMixin row-wise and aggregate prompts
timzsu Sep 24, 2026
7ab51b6
feat: echo task runs a function over whole upstream lists
timzsu Sep 24, 2026
b4cf53f
fix: echo function arguments and sandbox sources fail closed
timzsu Sep 24, 2026
1141a99
fix(worker): declare tabulate, which graph-template DataFrame renderi…
timzsu Sep 24, 2026
d1c0e3e
fix(sdk): mirror APIGroupItem and grouped APIResult.items
timzsu Sep 24, 2026
13a42cf
fix(worker): a dataframe over zero rows runs zero rows
timzsu Sep 24, 2026
f86ee96
fix(worker): expose ValueError to sandboxed functions
timzsu Sep 24, 2026
657ba7b
fix(worker): report status from first row across grouped results
timzsu Sep 24, 2026
bb949c8
fix(worker): decide dataframe grouping from upstream structure, not c…
timzsu Sep 25, 2026
5ad3ce1
fix(worker): group dataframe columns over list-mode lambda output
timzsu Sep 25, 2026
fc4a960
test(worker): assert APIResult dumps grouped rows under the json alias
timzsu Sep 25, 2026
ae0be15
test(sdk): assert grouped APIResult rows dump under the json alias
timzsu Sep 25, 2026
16676b6
feat(worker): log per-call statistics in the API executor
timzsu Sep 25, 2026
5db7f13
fix(worker): harden API executor prompt and telemetry handling
timzsu Sep 27, 2026
77e4bec
docs: document grouped API results and echo function data type
timzsu Sep 27, 2026
608e834
refactor(worker): pass summary snapshot as named locals
timzsu Sep 27, 2026
6bf1192
fix(worker): cancel a zero-row grouped API task after the pool drains
timzsu Sep 27, 2026
5ef1e96
fix(worker): log the real attempt count for an exhausted connection r…
timzsu Sep 27, 2026
c16c3c4
Merge main (PR 155) into zsu/api-executor-batch
timzsu Sep 27, 2026
037d99b
Merge zsu/api-executor-batch (main with PR 155) into zsu/api-executor…
timzsu Sep 27, 2026
d8ac38f
Address PR comments
timzsu Sep 28, 2026
477a003
Merge zsu/api-executor-batch (PR 141 review round 2) into zsu/api-exe…
timzsu Sep 28, 2026
b1fe866
Fix PR comments.
timzsu Sep 28, 2026
64e1f92
Merge origin/main into zsu/api-executor-batch
timzsu Sep 29, 2026
2cbe3d5
Merge zsu/api-executor-batch (PR 141 final) into zsu/api-executor-dat…
timzsu Sep 29, 2026
7a493be
Merge origin/main (PR 141 squash) into zsu/api-executor-dataframe
timzsu Sep 29, 2026
a48f21f
Merge remote-tracking branch 'origin/main' into zsu/pr157-rescope
timzsu Sep 30, 2026
5fee751
refactor: remove echo data.type function mode
timzsu Sep 30, 2026
33f2f9e
test: read a python stage's output in dataframe API tasks
timzsu Sep 30, 2026
a2be9e5
docs: document how a dataframe API task reads a python stage
timzsu Sep 30, 2026
6e64b31
fix: issue one aggregate prompt per group in graph_template
timzsu Sep 30, 2026
631c857
test: graph_template direct grouped column issues one request per group
timzsu Sep 30, 2026
73fc8bc
fix: record api calls that fail parsing or are cancelled
timzsu Sep 30, 2026
ef976c8
fix: reject echo data.type function with a pointer to python tasks
timzsu Sep 30, 2026
2460ac5
style: make APIGroupItem docstrings one paragraph
timzsu Sep 30, 2026
fb29688
docs: dataframe api results are always grouped into APIGroupItems
timzsu Sep 30, 2026
5217718
fix: keep per-row expansion for ungrouped graph-template columns
timzsu Sep 30, 2026
8257947
fix: address the review round on grouping, echo and API usage
timzsu Oct 1, 2026
2bce1a1
feat: report token usage per API task and per workflow
timzsu Oct 1, 2026
d7af140
fix: serve tasks contribute no usage to the workflow sum
timzsu Oct 1, 2026
3ea5ac7
docs: shorten the grouped results section
timzsu Oct 1, 2026
c33e6b2
docs: restore the embedded {{prompt}} substitution rule
timzsu Oct 1, 2026
4353593
fix: a condition-skipped task records no usage instead of nulling the…
timzsu Oct 1, 2026
01ee51c
test: guard APIUsage against server and SDK drift
timzsu Oct 1, 2026
349dfae
refactor: drop workflow-level usage; clients sum per-task usage
timzsu Oct 2, 2026
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
39 changes: 17 additions & 22 deletions docs/WORKFLOWS.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,10 +71,10 @@ contract.

## API task

`taskType: api` sends one HTTP request per `spec.data` row, in parallel, and
returns one `APIResult.items` entry per row, in order. `spec.data` is required
and accepts `list`, `dataset`, `graph_template`, and `dataframe`; row metadata
and `dataframe` table grouping are not applied.
`taskType: api` sends one HTTP request per `spec.data` prompt, in parallel, and
returns the responses in `APIResult.items`. `spec.data` is required; it supports
the same data types as the vLLM executor (`list`, `dataset`, `graph_template`,
`dataframe`).

By default it routes to the Nebula endpoint and authenticates with the worker's `NEBULA_API_TOKEN`.

Expand Down Expand Up @@ -113,24 +113,19 @@ exactly `{{prompt}}` takes the prompt as-is, so a message-list row fills
requests. Any failed row fails the task. Cancelling the task skips rows that
have not started and marks it cancelled once in-flight requests return.

```yaml
spec:
taskType: api
data:
type: list
items:
- Explain vector databases
- Explain attention
api:
method: POST
body:
model: gpt-4o
messages:
- role: user
content: "{{prompt}}"
response:
parse_json: true
```
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.
Comment on lines +116 to +120

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.


### Grouped results

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`.
Comment on lines +124 to +128

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.


## Python task

Expand Down
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,7 @@ runtime-analytics = [
"pandas>=2.3.3",
"psycopg[binary]>=3.2.12",
"sqlalchemy>=2.0.44",
"tabulate>=0.9.0",
]
runtime-worker-cpu = [
{ include-group = "runtime-worker-core" },
Expand Down
4 changes: 4 additions & 0 deletions sdk/src/flowmesh/models/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,10 @@
AgentResult,
AgentUsage,
AnyExecutorResult,
APIGroupItem,
APIItem,
APIResult,
APIUsage,
BaseExecutorResult,
CostEstimates,
DataProfilingResult,
Expand Down Expand Up @@ -106,8 +108,10 @@
)

__all__ = [
"APIGroupItem",
"APIItem",
"APIResult",
"APIUsage",
"ActiveWaitBreakdown",
"AgentBatchSummary",
"AgentItem",
Expand Down
4 changes: 4 additions & 0 deletions sdk/src/flowmesh/models/result/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,9 @@
AgentItem,
AgentMetadata,
AgentUsage,
APIGroupItem,
APIItem,
APIUsage,
CostEstimates,
DataRetrievalItem,
EchoItem,
Expand Down Expand Up @@ -91,8 +93,10 @@
_model.model_rebuild()

__all__ = [
"APIGroupItem",
"APIItem",
"APIResult",
"APIUsage",
"AgentBatchSummary",
"AgentItem",
"AgentMetadata",
Expand Down
6 changes: 4 additions & 2 deletions sdk/src/flowmesh/models/result/catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,9 @@
AgentItem,
AgentMetadata,
AgentUsage,
APIGroupItem,
APIItem,
APIUsage,
CostEstimates,
DataRetrievalItem,
EchoItem,
Expand Down Expand Up @@ -211,9 +213,9 @@ class APIResult(StrictExecutorResult):
truncated: bool = False
headers: dict[str, str] | None = None
response_json: Any = Field(default=None, alias="json")
usage: dict[str, Any] | None = None
usage: APIUsage | None = None
text: str | None = None
items: list[APIItem] = Field(default_factory=list)
items: list[APIItem | APIGroupItem] = Field(default_factory=list)


class SSHResult(StrictExecutorResult):
Expand Down
19 changes: 19 additions & 0 deletions sdk/src/flowmesh/models/result/payloads.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,17 @@ class GenerationUsage(StrictModel):
latency_sec: float


class APIUsage(StrictModel):
prompt_tokens: int
completion_tokens: int
reasoning_tokens: int
calls: int
failures: int
retries: int
truncated_calls: int
wall_sec: float


class EmbeddingUsage(StrictModel):
prompt_tokens: int
total_tokens: int
Expand Down Expand Up @@ -174,3 +185,11 @@ class APIItem(StrictModel):
usage: dict[str, Any] | None = None
text: str | None = None
prompt: str | None = None


class APIGroupItem(StrictModel):
"""One group's row responses in a batched API task over grouped data; ``rows``
holds the group's row responses in order."""

index: int
rows: list[APIItem]
4 changes: 4 additions & 0 deletions src/shared/schemas/result/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,9 @@
AgentItem,
AgentMetadata,
AgentUsage,
APIGroupItem,
APIItem,
APIUsage,
CostEstimates,
DataRetrievalItem,
EchoItem,
Expand Down Expand Up @@ -91,7 +93,9 @@

__all__ = [
"APIItem",
"APIGroupItem",
"APIResult",
"APIUsage",
"AgentBatchSummary",
"AgentItem",
"AgentMetadata",
Expand Down
12 changes: 7 additions & 5 deletions src/shared/schemas/result/catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,9 @@
AgentItem,
AgentMetadata,
AgentUsage,
APIGroupItem,
APIItem,
APIUsage,
CostEstimates,
DataRetrievalItem,
EchoItem,
Expand Down Expand Up @@ -245,9 +247,9 @@ class EchoResult(StrictExecutorResult):


class APIResult(StrictExecutorResult):
"""HTTP request output. ``response_json``/``usage``/``headers`` are the
upstream API's own payloads and stay open mappings. ``items`` carries one
entry per row."""
"""HTTP request output. ``response_json``/``headers`` are the upstream
API's own payloads and stay open mappings; ``usage`` is the task's summed
token/call accounting. ``items`` carries one entry per row."""

task_type: Literal[TaskType.API] = TaskType.API
executor: str
Expand All @@ -257,9 +259,9 @@ class APIResult(StrictExecutorResult):
truncated: bool = False
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]?

text: str | None = None
items: list[APIItem] = Field(default_factory=list)
items: list[APIItem | APIGroupItem] = Field(default_factory=list)


class SSHResult(StrictExecutorResult):
Expand Down
21 changes: 21 additions & 0 deletions src/shared/schemas/result/payloads.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,19 @@ class GenerationUsage(StrictModel):
latency_sec: float


class APIUsage(StrictModel):
"""Token/call accounting for an API task, summed over its requests."""

prompt_tokens: int
completion_tokens: int
reasoning_tokens: int
calls: int
failures: int
retries: int
truncated_calls: int
wall_sec: float


class EmbeddingUsage(StrictModel):
"""Token/latency accounting for embedding inference."""

Expand Down Expand Up @@ -230,3 +243,11 @@ class APIItem(StrictModel):
usage: dict[str, Any] | None = None
text: str | None = None
prompt: str | None = None


class APIGroupItem(StrictModel):
"""One group's row responses in a batched API task over grouped data; ``rows``
holds the group's row responses in order."""

index: int
rows: list[APIItem]
Loading
Loading