-
Notifications
You must be signed in to change notification settings - Fork 2
feat: API executor runs DataMixin prompts, reads python stages and reports token usage #157
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
22b222e
ad4752f
275ca45
67b795b
38a2a07
096a58a
0fb3d26
d2fe09c
b5583ea
630d55a
aff4e63
da6a2df
d3e31b8
4358a0b
dc9da05
da2aef1
0b4f12c
7dff652
9524d1c
0341004
3f9e5c7
40bf0a2
9691501
76e218c
d4ced9f
b6c01a6
de41f50
e4991dd
2109405
7d2ddb9
4a05001
d11996d
fc8d5c2
821b2d0
7ab51b6
b4cf53f
1141a99
d1c0e3e
13a42cf
f86ee96
657ba7b
bb949c8
5ad3ce1
fc4a960
ae0be15
16676b6
5db7f13
77e4bec
608e834
6bf1192
5ef1e96
c16c3c4
037d99b
d8ac38f
477a003
b1fe866
64e1f92
2cbe3d5
7a493be
a48f21f
5fee751
33f2f9e
a2be9e5
6e64b31
631c857
73fc8bc
ef976c8
2460ac5
fb29688
5217718
8257947
2bce1a1
d7af140
3ea5ac7
c33e6b2
4353593
01ee51c
349dfae
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -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`. | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
|
|
@@ -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. | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| ### 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
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. A python task's return value is
Suggested change
|
||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| ## Python task | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -20,7 +20,9 @@ | |
| AgentItem, | ||
| AgentMetadata, | ||
| AgentUsage, | ||
| APIGroupItem, | ||
| APIItem, | ||
| APIUsage, | ||
| CostEstimates, | ||
| DataRetrievalItem, | ||
| EchoItem, | ||
|
|
@@ -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 | ||
|
|
@@ -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 | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. v0.1.9 stored the upstream usage dict here ( |
||
| text: str | None = None | ||
| items: list[APIItem] = Field(default_factory=list) | ||
| items: list[APIItem | APIGroupItem] = Field(default_factory=list) | ||
|
|
||
|
|
||
| class SSHResult(StrictExecutorResult): | ||
|
|
||
There was a problem hiding this comment.
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.